From aa02826dbbddf10542bf236b95498e200a1a4114 Mon Sep 17 00:00:00 2001 From: dmgnr Date: Sat, 5 Sep 2026 12:09:16 +0000 Subject: [PATCH] feat: add CI deployment and live progress --- README.md | 2 +- command/ci.ts | 80 ++++ command/main.ts | 4 +- command/up.ts | 357 +++++++++++++++- command/users.ts | 142 +++++++ lib/k8s.ts | 5 +- package.json | 2 +- server/app.ts | 573 ++++++++++++++++++++------ server/auth.ts | 103 ++++- server/authorization.ts | 8 + server/index.ts | 2 +- server/kubernetes-state.ts | 5 +- server/kubernetes-store.ts | 143 ++++++- server/management.ts | 122 +++++- server/operation-store.ts | 108 ++++- server/registry.ts | 46 ++- shared/api.ts | 31 ++ tests/command/administration.test.ts | 51 +++ tests/command/ci.test.ts | 19 + tests/command/up-api.test.ts | 159 +++++++ tests/server/api-keys.test.ts | 463 +++++++++++++++++++++ tests/server/app.test.ts | 190 +++++++++ tests/server/kubernetes-state.test.ts | 47 ++- tests/server/kubernetes-store.test.ts | 40 +- tests/server/management.test.ts | 73 ++++ tests/server/operation-store.test.ts | 53 +++ tests/server/registry.test.ts | 63 +++ 27 files changed, 2711 insertions(+), 180 deletions(-) create mode 100644 command/ci.ts create mode 100644 tests/command/ci.test.ts create mode 100644 tests/server/api-keys.test.ts diff --git a/README.md b/README.md index 316f32f..db3a031 100644 --- a/README.md +++ b/README.md @@ -93,7 +93,7 @@ The server grants capabilities through three roles: - `viewer` — read-only cluster access (`kubernetes:read`) - `operator` — `viewer` plus `kubernetes:write` and `kubernetes:exec` - `admin` — all capabilities, including user administration - (`users:read`, `users:write`, `sessions:revoke`) + (`users:read`, `users:write`, `sessions:revoke`, `platform:adopt`) Administer users with the `users` command tree: diff --git a/command/ci.ts b/command/ci.ts new file mode 100644 index 0000000..9094f0d --- /dev/null +++ b/command/ci.ts @@ -0,0 +1,80 @@ +import { defineCommand } from "citty"; +import { + apiRequest, + type ApiRequestInit, + type ApiRequestOptions, +} from "../lib/api"; +import { ctx } from "../lib/context"; +import { resolveTrustIdentity } from "../lib/trust"; +import { runUp } from "./up"; + +type CiRequest = ( + path: string, + init?: ApiRequestInit, + options?: ApiRequestOptions, +) => Promise; + +function requireApiKey(value: string | undefined): string { + const token = value?.trim(); + if (!token) + throw new Error("KUBER_API_KEY or --api-key is required for kuber ci"); + return token; +} + +export function apiKeyRequest(token: string): CiRequest { + return (path, init = {}, options = {}) => { + const headers = new Headers(init.headers); + headers.set("authorization", `Bearer ${token}`); + return apiRequest( + path, + { ...init, headers }, + { ...options, authenticated: false }, + ); + }; +} + +export async function runCi( + build: boolean, + grantTrust: boolean, + token: string, + request: CiRequest = apiKeyRequest(token), +): Promise { + const { project, cwd } = ctx(); + const trust = await resolveTrustIdentity(project, cwd); + if (grantTrust) { + await request(`/workspaces/${encodeURIComponent(project)}/trust`, { + method: "POST", + json: { fingerprint: trust.fingerprint }, + }); + } + await runUp(build, request, { trust: grantTrust ? trust : undefined }); +} + +export const ci = defineCommand({ + meta: { + name: "ci", + description: "Deploy using an API key without a session", + }, + args: { + apiKey: { + type: "string", + description: "API key (defaults to KUBER_API_KEY)", + }, + trust: { + type: "boolean", + description: "Grant current directory trust before deploying", + }, + build: { + type: "boolean", + default: true, + negativeDescription: "Resolve existing images without building", + }, + }, + async run({ args }) { + await runCi( + args.build, + Boolean(args.trust), + requireApiKey(args.apiKey ?? process.env.KUBER_API_KEY), + ); + }, +}); diff --git a/command/main.ts b/command/main.ts index 330fb6c..7651131 100644 --- a/command/main.ts +++ b/command/main.ts @@ -1,6 +1,7 @@ import tab from "@bomb.sh/tab/citty"; import { defineCommand } from "citty"; import { audit } from "./audit"; +import { ci } from "./ci"; import { login, logout, whoami } from "./auth"; import { db } from "./db"; import { down } from "./down"; @@ -21,7 +22,7 @@ import { trust } from "./trust"; export const main = defineCommand({ meta: { name: "kuber", - version: "2.2.0", + version: "2.3.0", description: "Docker Compose -> K8s translation layer", }, args: { @@ -32,6 +33,7 @@ export const main = defineCommand({ }, subCommands: { audit, + ci, db, down, export: exportCommand, diff --git a/command/up.ts b/command/up.ts index 13f62f1..ef8dae4 100644 --- a/command/up.ts +++ b/command/up.ts @@ -1,5 +1,5 @@ import { defineCommand } from "citty"; -import { Listr } from "listr2"; +import { Listr, type ListrTaskWrapper } from "listr2"; import { randomUUID } from "node:crypto"; import type { ComposeSpecification } from "../schema/docker.d"; import type { KuberResource } from "../types"; @@ -24,6 +24,7 @@ import { requireLocalTrust, resolveTrustIdentity, trustHeaders, + type TrustIdentity, } from "../lib/trust"; type Workspace = { @@ -54,9 +55,44 @@ type OperationResponse = { operation?: { status?: OperationStatus }; }; +export type ResourceOperationEvent = { + sequence: number; + data: { + resource: { + apiVersion: string; + kind: string; + name: string; + namespace?: string; + }; + phase: "apply" | "wait" | "delete"; + state: "started" | "succeeded" | "failed" | "aborted"; + }; +}; + export type OperationResumeOptions = { now?: () => number; sleep?: (milliseconds: number) => Promise; + onEvent?: (event: ResourceOperationEvent) => void; +}; + +type OperationEventsResponse = { + items?: ResourceOperationEvent[]; + retainedFirstSequence?: unknown; + cursorGap?: unknown; +}; + +type ResourceOperationTarget = { + apiVersion: string; + kind: string; + name: string; + namespace?: string; +}; + +export type LiveResourceOperationOptions = { + onTaskStarted?: ( + target: ResourceOperationTarget, + task: ListrTaskWrapper, + ) => void; }; type UpContext = { @@ -78,6 +114,14 @@ const RECOVERABLE_API_ERROR_CODES = new Set([ "HTTP_504", ]); +class OperationProgressCursorGapError extends Error { + constructor() { + super( + "Operation progress history was truncated; refusing to report an incomplete resource stream", + ); + } +} + export function workspaceAdoptionRoute(project: string): string { return `/workspaces/${encodeURIComponent(project)}/adopt`; } @@ -99,6 +143,44 @@ export function getDeploymentNames(resources: KubernetesResource[]): string[] { .filter((name): name is string => Boolean(name)); } +export function resourceOperationEventTitle( + event: ResourceOperationEvent, +): string { + const { resource, phase, state } = event.data; + return `${phase === "apply" ? "Applied" : phase === "wait" ? "Waited for" : "Deleted"} ${resource.kind}/${resource.name}: ${state}`; +} + +function isResourceOperationEvent( + event: unknown, +): event is ResourceOperationEvent { + if (!event || typeof event !== "object") return false; + const { sequence, data } = event as { + sequence?: unknown; + data?: unknown; + }; + if (!Number.isSafeInteger(sequence) || !data || typeof data !== "object") + return false; + const { resource, phase, state } = data as { + resource?: unknown; + phase?: unknown; + state?: unknown; + }; + if (!resource || typeof resource !== "object") return false; + const identity = resource as Record; + return ( + typeof identity.apiVersion === "string" && + typeof identity.kind === "string" && + typeof identity.name === "string" && + (identity.namespace === undefined || + typeof identity.namespace === "string") && + (phase === "apply" || phase === "wait" || phase === "delete") && + (state === "started" || + state === "succeeded" || + state === "failed" || + state === "aborted") + ); +} + function workspaceInput( compose: ComposeSpecification, snapshot: WorkspaceSnapshot, @@ -235,6 +317,17 @@ function isInterruptedOperation(status: OperationStatus): boolean { ); } +function hasOperationEventCursorGap( + events: OperationEventsResponse, + after: number, +): boolean { + return ( + events.cursorGap === true || + (Number.isSafeInteger(events.retainedFirstSequence) && + (events.retainedFirstSequence as number) > after + 1) + ); +} + function operationResumeDeadline( rolloutTimeoutMs: number, now: number, @@ -255,10 +348,12 @@ async function resumeManagedOperation( const deadline = operationResumeDeadline(rolloutTimeoutMs, now()); const headers = new Headers(init.headers); headers.set("idempotency-key", randomUUID()); + headers.set("prefer", "respond-async"); let operationInit = { ...init, headers }; let operationId: string | undefined; let backoffMs = OPERATION_RESUME_INITIAL_BACKOFF_MS; let restartRetryPending = false; + let eventCursor = 0; for (;;) { if (restartRetryPending && now() >= deadline) @@ -285,10 +380,37 @@ async function resumeManagedOperation( status = operationStatusFromResponse(response); } + if (operationId) { + try { + const events = await managementRequest( + project, + request, + `/operations/${encodeURIComponent(operationId)}/events?after=${eventCursor}`, + {}, + ); + if (hasOperationEventCursorGap(events, eventCursor)) + throw new OperationProgressCursorGapError(); + for (const event of events.items ?? []) { + if ( + !Number.isSafeInteger(event.sequence) || + event.sequence <= eventCursor + ) + continue; + eventCursor = event.sequence; + if (!isResourceOperationEvent(event)) continue; + options.onEvent?.(event); + } + } catch (error) { + if (error instanceof OperationProgressCursorGapError) throw error; + // Event polling is additive; never delay resumable operation status polling. + } + } + if (!operationId || !status || status.state === "succeeded") return; if (isInterruptedOperation(status)) { if (now() >= deadline) throw operationFailure(operationId, status); operationId = undefined; + eventCursor = 0; restartRetryPending = true; headers.set("idempotency-key", randomUUID()); operationInit = { ...init, headers }; @@ -301,6 +423,7 @@ async function resumeManagedOperation( ) { if (now() >= deadline) throw adoptionHint(project, error); operationId = undefined; + eventCursor = 0; restartRetryPending = true; headers.set("idempotency-key", randomUUID()); operationInit = { ...init, headers }; @@ -319,23 +442,31 @@ async function resumeManagedOperation( } } -export async function reconcileResources( +async function planResources( project: string, resources: KubernetesResource[], - rolloutTimeoutMs: number, - postApply: - | ((resources: KubernetesResource[]) => void | Promise) - | undefined, request: ApiRequester = apiRequest, - resumeOptions: OperationResumeOptions = {}, ): Promise { const workspacePath = `/workspaces/${encodeURIComponent(project)}`; - const plan = await managementRequest( + return managementRequest( project, request, `${workspacePath}/resources/plan`, { method: "POST", json: { resources } }, ); +} + +async function applyResourcePlan( + project: string, + plan: ResourcePlan, + rolloutTimeoutMs: number, + postApply: + | ((resources: KubernetesResource[]) => void | Promise) + | undefined, + request: ApiRequester, + resumeOptions: OperationResumeOptions, +): Promise { + const workspacePath = `/workspaces/${encodeURIComponent(project)}`; await resumeManagedOperation( project, request, @@ -376,17 +507,117 @@ export async function reconcileResources( resumeOptions, ); } +} + +export async function reconcileResources( + project: string, + resources: KubernetesResource[], + rolloutTimeoutMs: number, + postApply: + | ((resources: KubernetesResource[]) => void | Promise) + | undefined, + request: ApiRequester = apiRequest, + resumeOptions: OperationResumeOptions = {}, +): Promise { + const plan = await planResources(project, resources, request); + await applyResourcePlan( + project, + plan, + rolloutTimeoutMs, + postApply, + request, + resumeOptions, + ); return plan; } +function resourceOperationTarget(resource: { + apiVersion?: string; + kind?: string; + name?: string; + namespace?: string; + metadata?: { name?: string; namespace?: string }; +}): ResourceOperationTarget | undefined { + const name = resource.name ?? resource.metadata?.name; + if (!resource.apiVersion || !resource.kind || !name) return; + return { + apiVersion: resource.apiVersion, + kind: resource.kind, + name, + namespace: resource.namespace ?? resource.metadata?.namespace, + }; +} + +function resourceOperationKey(resource: ResourceOperationTarget): string { + return `${resource.apiVersion}\0${resource.kind}\0${resource.namespace ?? ""}\0${resource.name}`; +} + +function resourceOperationTaskTitle( + phase: ResourceOperationEvent["data"]["phase"], + resource: ResourceOperationTarget, +): string { + const action = + phase === "apply" ? "Apply" : phase === "wait" ? "Wait for" : "Delete"; + return `${action} ${resource.kind}/${resource.name}`; +} + +export async function runLiveResourceOperation( + task: ListrTaskWrapper, + phase: ResourceOperationEvent["data"]["phase"], + targets: ResourceOperationTarget[], + operation: ( + onEvent: (event: ResourceOperationEvent) => void, + ) => Promise, + options: LiveResourceOperationOptions = {}, +): Promise { + if (targets.length === 0) return operation(() => {}); + const active = new Map< + string, + ListrTaskWrapper + >(); + const complete: Array<() => void> = []; + let started = 0; + let markStarted!: () => void; + const ready = new Promise((resolve) => { + markStarted = resolve; + }); + const children = task.newListr( + targets.map((target) => ({ + title: resourceOperationTaskTitle(phase, target), + task: (_target, child) => { + active.set(resourceOperationKey(target), child); + options.onTaskStarted?.(target, child); + started += 1; + if (started === targets.length) markStarted(); + return new Promise((resolve) => complete.push(resolve)); + }, + })), + ); + const childRun = children.run(); + await ready; + try { + await operation((event) => { + const resource = event.data.resource; + const child = active.get(resourceOperationKey(resource)); + if (!child) return; + child.output = event.data.state; + task.output = resourceOperationEventTitle(event); + }); + } finally { + for (const resolve of complete) resolve(); + await childRun; + } +} + export async function runUp( build: boolean, request: ApiRequester = apiRequest, + options: { trust?: TrustIdentity } = {}, ) { const { project, compose, cwd, config, hookContext: getHookContext } = ctx(); - const trusted = await requireLocalTrust( - await resolveTrustIdentity(project, cwd), - ); + const trusted = + options.trust ?? + (await requireLocalTrust(await resolveTrustIdentity(project, cwd))); const baseRequest = request; request = async ( path: string, @@ -539,20 +770,104 @@ export async function runUp( { title: "Reconcile resources", task: async (taskCtx, task) => { - taskCtx.plan = await reconcileResources( + taskCtx.plan = await planResources( project, taskCtx.resources!, - config.rolloutTimeoutMs, - config.postApply - ? async (resources) => - config.postApply?.( - resources as KuberResource[], - await getHookContext(), - ) - : undefined, request, ); - task.output = `${taskCtx.plan.desired.length} applied, ${taskCtx.plan.stale.length} stale deleted`; + const plan = taskCtx.plan; + const desired = plan.desired + .map(resourceOperationTarget) + .filter((resource): resource is ResourceOperationTarget => + Boolean(resource), + ); + const deployments = plan.desired + .filter((resource) => resource.kind === "Deployment") + .map(resourceOperationTarget) + .filter((resource): resource is ResourceOperationTarget => + Boolean(resource), + ); + const stale = plan.stale + .map(resourceOperationTarget) + .filter((resource): resource is ResourceOperationTarget => + Boolean(resource), + ); + const resourcePath = `${workspacePath}/resources`; + return task.newListr([ + { + title: "Apply resources", + task: async (_ctx, phaseTask) => + runLiveResourceOperation( + phaseTask, + "apply", + desired, + (onEvent) => + resumeManagedOperation( + project, + request, + `${resourcePath}/apply`, + { method: "POST", json: { resources: plan.desired } }, + config.rolloutTimeoutMs, + { onEvent }, + ), + ), + }, + { + title: "Run post-apply hook", + skip: !config.postApply, + task: async () => + config.postApply?.( + plan.desired as KuberResource[], + await getHookContext(), + ), + }, + { + title: "Wait for deployments", + skip: deployments.length === 0, + task: async (_ctx, phaseTask) => + runLiveResourceOperation( + phaseTask, + "wait", + deployments, + (onEvent) => + resumeManagedOperation( + project, + request, + `${resourcePath}/wait`, + { + method: "POST", + json: { + deployments: deployments.map( + (deployment) => deployment.name, + ), + timeoutMs: config.rolloutTimeoutMs, + }, + }, + config.rolloutTimeoutMs, + { onEvent }, + ), + ), + }, + { + title: "Delete stale resources", + skip: stale.length === 0, + task: async (_ctx, phaseTask) => + runLiveResourceOperation( + phaseTask, + "delete", + stale, + (onEvent) => + resumeManagedOperation( + project, + request, + `${resourcePath}/delete`, + { method: "POST", json: { resources: plan.stale } }, + config.rolloutTimeoutMs, + { onEvent }, + ), + ), + }, + ]); }, }, ], diff --git a/command/users.ts b/command/users.ts index b17de2e..c4d987b 100644 --- a/command/users.ts +++ b/command/users.ts @@ -5,6 +5,10 @@ import { apiRequest, type ApiRequestInit } from "../lib/api"; import { toTable } from "../lib/format"; import type { CreateUserRequest, + ApiKey, + CreateApiKeyRequest, + CreateApiKeyResponse, + ListApiKeysResponse, ListUsersResponse, UpdateUserRequest, User, @@ -40,6 +44,27 @@ function parseRoles(value: unknown): UserRole[] { return [...new Set(roles)]; } +function parseCapabilities( + value: unknown, +): CreateApiKeyRequest["capabilities"] { + const capabilities = String(value ?? "") + .split(",") + .map((capability) => capability.trim()) + .filter(Boolean); + const allowed = new Set([ + "kubernetes:read", + "kubernetes:write", + "kubernetes:exec", + "users:read", + "users:write", + "sessions:revoke", + "platform:adopt", + ]); + if (!capabilities.length || capabilities.some((item) => !allowed.has(item))) + throw new Error("Provide at least one valid capability"); + return [...new Set(capabilities)] as CreateApiKeyRequest["capabilities"]; +} + function renderUsers(users: User[]): string { if (users.length === 0) return "No users"; return toTable( @@ -172,6 +197,56 @@ export async function revokeUserSessions( return `Revoked ${result.revoked} session${result.revoked === 1 ? "" : "s"} for ${result.username}`; } +function renderApiKeys(keys: ApiKey[]): string { + if (!keys.length) return "No API keys"; + return toTable( + keys.map((key) => ({ + id: key.id, + capabilities: key.capabilities.join(","), + workspace: key.workspace ?? "", + expires: key.expiresAt, + disabled: key.disabled ? "yes" : "no", + })), + ); +} + +export async function listApiKeys( + username: string, + request: UsersApiRequest = apiRequest, +): Promise { + const response = await request( + `/users/${encodeURIComponent(username)}/keys`, + ); + return renderApiKeys(response.items); +} + +export async function createApiKey( + username: string, + body: CreateApiKeyRequest, + request: UsersApiRequest = apiRequest, +): Promise { + return request( + `/users/${encodeURIComponent(username)}/keys`, + { method: "POST", json: body }, + ); +} + +export async function revokeApiKey( + username: string, + id: string, + confirmed: boolean, + request: UsersApiRequest = apiRequest, +): Promise { + if (!confirmed) return "API key revocation cancelled"; + await request( + `/users/${encodeURIComponent(username)}/keys/${encodeURIComponent(id)}`, + { + method: "DELETE", + }, + ); + return `Revoked API key ${id} for ${username}`; +} + const list = defineCommand({ meta: { name: "ls", description: "List users" }, async run() { @@ -239,6 +314,72 @@ const revoke = defineCommand({ }, }); +const keys = defineCommand({ + meta: { name: "keys", description: "Manage user API keys" }, + subCommands: { + ls: defineCommand({ + meta: { name: "ls", description: "List a user's API keys" }, + async run({ args }) { + console.log(await listApiKeys(requireUsername(args._[0]))); + }, + }), + create: defineCommand({ + meta: { name: "create", description: "Create an API key" }, + args: { + capabilities: { + type: "string", + required: true, + description: "Comma-separated capabilities", + }, + workspace: { + type: "string", + description: "Restrict the key to a workspace", + }, + expiresDays: { + type: "string", + default: "90", + description: "Expiry in days (1-365)", + }, + }, + async run({ args }) { + const days = Number(args.expiresDays); + if (!Number.isSafeInteger(days) || days < 1 || days > 365) + throw new Error("--expires-days must be an integer from 1 to 365"); + const key = await createApiKey(requireUsername(args._[0]), { + capabilities: parseCapabilities(args.capabilities), + ...(args.workspace && { workspace: args.workspace }), + expiresAt: new Date( + Date.now() + days * 24 * 60 * 60 * 1000, + ).toISOString(), + }); + console.log( + "Store this API key securely now. It will not be shown again:", + ); + console.log(key.token); + console.log(renderApiKeys([key])); + }, + }), + revoke: defineCommand({ + meta: { name: "revoke", description: "Revoke an API key" }, + args: { + yes: { + type: "boolean", + alias: "y", + description: "Revoke without confirmation", + }, + }, + async run({ args }) { + const username = requireUsername(args._[0]); + const id = requireUsername(args._[1]); + const confirmed = + Boolean(args.yes) || + (await confirmDeletion(`API key '${id}' for ${username}`)); + console.log(await revokeApiKey(username, id, confirmed)); + }, + }), + }, +}); + export const users = defineCommand({ meta: { name: "users", description: "Administer users" }, subCommands: { @@ -246,6 +387,7 @@ export const users = defineCommand({ delete: remove, disable: disabledCommand("disable", true), enable: disabledCommand("enable", false), + keys, ls: list, revoke, update, diff --git a/lib/k8s.ts b/lib/k8s.ts index 50b67b0..86a3908 100644 --- a/lib/k8s.ts +++ b/lib/k8s.ts @@ -11,7 +11,10 @@ import { createKubernetesHttpLibrary } from "./k8s-http"; const kc = new KubeConfig(); kc.loadFromDefault(); -const bunHttpLibrary = createKubernetesHttpLibrary(); +const bunHttpLibrary = createKubernetesHttpLibrary({ + maxConcurrent: 4, + minIntervalMs: 0, +}); const makeApiClient = kc.makeApiClient.bind(kc); kc.makeApiClient = ((apiClientType) => { diff --git a/package.json b/package.json index 18aae4b..99f3df4 100644 --- a/package.json +++ b/package.json @@ -1,6 +1,6 @@ { "name": "@dmgnr/kuber", - "version": "2.2.0", + "version": "2.3.0", "description": "Docker Compose to Kubernetes translation layer", "bin": { "kuber": "dist/index.js" diff --git a/server/app.ts b/server/app.ts index 6894363..8a72db8 100644 --- a/server/app.ts +++ b/server/app.ts @@ -9,10 +9,12 @@ import { tokenHashesEqual, type AuthStore, type KuberUser, + type ApiKeyRecord, type SessionRecord, } from "./auth"; import { hasCapability, + isCapability, isRole, type Capability, type Role, @@ -36,7 +38,11 @@ import type { ExecServerWireFrame, ExecStartFrame, } from "../lib/exec-api"; -import type { ManagementService, ResourceIdentity } from "./management"; +import type { + ManagementService, + OperationProgressEmitter, + ResourceIdentity, +} from "./management"; import { OperationConflictError, OperationNotFoundError, @@ -67,6 +73,8 @@ 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; @@ -102,7 +110,11 @@ export type AppOptions = { }; type LoginFailures = { count: number; resetAt: number }; -type Identity = { user: KuberUser; session: SessionRecord }; +type Identity = { + user: KuberUser; + session?: SessionRecord; + apiKey?: ApiKeyRecord; +}; export type WorkspaceAdoptionResult = { workspaceId: string; @@ -126,7 +138,8 @@ export async function cleanupExpiredSessions( store: AuthStore, now = Date.now(), ): Promise { - return store.deleteExpiredSessions(now); + const sessions = await store.deleteExpiredSessions(now); + return sessions + (await store.deleteExpiredApiKeys(now)); } class HttpError extends Error { @@ -173,18 +186,29 @@ export async function authenticateRequest( if (!token) return; const tokenHash = hashToken(token); const session = await options.store.getSession(tokenHash); - if ( - !session || - !tokenHashesEqual(tokenHash, session.tokenHash) || - Date.parse(session.expiresAt) <= now() - ) { - if (session) await options.store.deleteSession(tokenHash); - return; + 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 user = await options.store.getUser(session.username); - if (!user || user.disabled || user.authVersion !== session.authVersion) + const apiKey = await options.store.getApiKey(tokenHash); + if ( + !apiKey || + !tokenHashesEqual(tokenHash, apiKey.tokenHash) || + Date.parse(apiKey.expiresAt) <= now() + ) return; - return { user, session }; + const user = await options.store.getUser(apiKey.username); + if (!user || user.disabled) return; + return { user, apiKey }; } function pathPart(value: string): string { @@ -482,7 +506,12 @@ export function createApp( request: Request, capability: Capability, ): Promise { - if (hasCapability(identity.user.roles, capability)) return; + 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, @@ -492,6 +521,60 @@ export function createApp( ); } + 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( @@ -594,7 +677,10 @@ export function createApp( action: operation.spec.action, }, status: { - ...operation.status, + ...(() => { + const { events: _events, ...status } = operation.status; + return status; + })(), ...(operation.status.error && { error: sanitizeOperationError( operation.status.error, @@ -673,7 +759,10 @@ export function createApp( workspace: Workspace, action: string, input: unknown, - execute: (signal: AbortSignal) => Promise, + execute: ( + signal: AbortSignal, + emit: OperationProgressEmitter, + ) => Promise, ): Promise { if (!options.operationStore) throw new HttpError( @@ -772,110 +861,134 @@ export function createApp( operation.metadata.name, ); }; - 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); - await requireLeaseOwnership(); - const completed = await options.operationStore.transition( - operation.metadata.name, - "succeeded", - { result }, - ); + const executeClaimed = async (): Promise => { try { - await audit( - identity, - request, - action, - "success", - undefined, - workspace.metadata.name, + 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) { + console.error( + "Failed to append successful operation audit event", + error, + ); + } + 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) { + console.error( + "Failed to append failed operation audit event", + auditError, + ); + } + 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) { + console.error("Failed to release workspace operation lease", error); + } + } + }; + if (request.headers.get("prefer") === "respond-async") { + void executeClaimed().catch((error) => { console.error( - "Failed to append successful operation audit event", + `Background operation '${operation.metadata.name}' failed`, error, ); - } - 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, + return response( { - code: leaseLost ? "WORKSPACE_LEASE_LOST" : "OPERATION_FAILED", - message: leaseLost - ? "Workspace operation lease ownership was lost" - : message, + operationId: operation.metadata.name, + operation: publicOperation(operation), }, + 202, + { location: `${API_PREFIX}/operations/${operation.metadata.name}` }, ); - 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) { - console.error( - "Failed to append failed operation audit event", - auditError, - ); - } - 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) { - console.error("Failed to release workspace operation lease", error); - } } + return executeClaimed(); } async function handleLogin(request: Request): Promise { @@ -956,7 +1069,8 @@ export function createApp( roles: identity.user.roles, }); if (request.method === "POST" && path === `${API_PREFIX}/logout`) { - await options.store.deleteSession(identity.session.tokenHash); + if (identity.session) + await options.store.deleteSession(identity.session.tokenHash); return new Response(null, { status: 204 }); } @@ -965,6 +1079,7 @@ export function createApp( ).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({ @@ -1041,6 +1156,7 @@ export function createApp( "BUILD_INVALID", "project is required", ); + await requireWorkspaceScope(identity, request, body.project); return response( await requireBuilds().submitBuild(body as unknown as BuildRequest), 202, @@ -1060,6 +1176,7 @@ export function createApp( "BUILD_INVALID", "project is required", ); + await requireWorkspaceScope(identity, request, body.project); return response( await requireBuilds().negotiateSnapshot(body.workspace as Sha256Digest), ); @@ -1067,13 +1184,6 @@ export function createApp( if (path === `${API_PREFIX}/images/resolve` && request.method === "POST") { await requireCapability(identity, request, "kubernetes:write"); - if (!options.resolveImage) - throw new HttpError( - 503, - "Service unavailable", - "BUILDS_UNAVAILABLE", - "Image resolution is not configured", - ); const body = await readJson(request); if (typeof body.project !== "string" || typeof body.service !== "string") throw new HttpError( @@ -1082,6 +1192,14 @@ export function createApp( "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)); } @@ -1093,6 +1211,11 @@ export function createApp( 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") @@ -1128,6 +1251,7 @@ export function createApp( "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)); @@ -1203,6 +1327,118 @@ export function createApp( 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]) { @@ -1293,11 +1529,27 @@ export function createApp( if (path === `${API_PREFIX}/workspaces`) { if (request.method === "GET") { await requireCapability(identity, request, "kubernetes:read"); - return response({ items: await requireWorkspaceStore().list() }); + 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, ); @@ -1320,13 +1572,18 @@ export function createApp( path === `${API_PREFIX}/platform/kuber-system/adopt` && request.method === "POST" ) { - if (!identity.user.roles.includes("admin")) + 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", - "The kuber-system platform adoption route requires the admin role", + "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( @@ -1353,6 +1610,7 @@ export function createApp( 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"); @@ -1580,7 +1838,7 @@ export function createApp( workspace, subpath!, body, - async (signal) => { + async (signal, emit) => { const management = requireManagement(); if (subpath === "resources/apply") { if (!Array.isArray(body.resources)) @@ -1588,7 +1846,7 @@ export function createApp( return management.applyResources( workspaceIdentity(workspace), body.resources as KubernetesObject[], - { signal }, + { signal, emit }, ); } if (subpath === "resources/wait") { @@ -1598,7 +1856,7 @@ export function createApp( workspaceIdentity(workspace), body.deployments as string[], typeof body.timeoutMs === "number" ? body.timeoutMs : undefined, - { signal }, + { signal, emit }, ); return { ready: true }; } @@ -1607,7 +1865,7 @@ export function createApp( await management.deleteResources( workspaceIdentity(workspace), body.resources as unknown as ResourceIdentity[], - { signal }, + { signal, emit }, ); return { deleted: body.resources.length }; }, @@ -1724,6 +1982,14 @@ export function createApp( 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, @@ -1734,11 +2000,47 @@ export function createApp( return response({ items: ( await options.operationStore.list( - url.searchParams.get("workspaceId") ?? undefined, + 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"); @@ -1750,6 +2052,11 @@ export function createApp( "OPERATION_NOT_FOUND", "Operation not found", ); + await requireWorkspaceScope( + identity, + request, + operation.spec.workspaceId, + ); return response(publicOperation(operation)); } @@ -1762,9 +2069,17 @@ export function createApp( "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( - url.searchParams.get("workspaceId") ?? undefined, + identity.apiKey?.workspace ?? requestedWorkspace, ), }); } @@ -1923,7 +2238,13 @@ export async function authorizeExecConnection( "A valid kuber login is required", { "www-authenticate": "Bearer" }, ); - if (!hasCapability(identity.user.roles, "kubernetes:exec")) { + 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({ diff --git a/server/auth.ts b/server/auth.ts index 6291302..969b3f9 100644 --- a/server/auth.ts +++ b/server/auth.ts @@ -1,5 +1,10 @@ import { createHash, randomBytes, timingSafeEqual } from "node:crypto"; -import { isRole, type Role } from "./authorization"; +import { + isCapability, + isRole, + type Capability, + type Role, +} from "./authorization"; export type KuberUser = { username: string; @@ -25,6 +30,18 @@ export type SessionInput = expiresAt: string; }; +export type ApiKeyRecord = { + id: string; + tokenHash: string; + username: string; + capabilities: Capability[]; + workspace?: string; + expiresAt: string; + disabled?: boolean; +}; + +export type NewApiKey = ApiKeyRecord; + export type NewKuberUser = Omit & { authVersion?: number; }; @@ -49,6 +66,11 @@ export interface AuthStore { revokeUserSessions(username: string): Promise; listExpiredSessions(now?: number): Promise; deleteExpiredSessions(now?: number): Promise; + getApiKey(tokenHash: string): Promise; + createApiKey(key: NewApiKey): Promise; + listApiKeys(username: string): Promise; + revokeApiKey(username: string, id: string): Promise; + deleteExpiredApiKeys(now?: number): Promise; } export function normalizeUser(user: NewKuberUser | KuberUser): KuberUser { @@ -85,6 +107,42 @@ export function normalizeSession(session: SessionRecord): SessionRecord { return { ...session }; } +export function normalizeApiKey(key: NewApiKey): ApiKeyRecord { + if (!/^[a-zA-Z0-9_-]{16,128}$/.test(key.id)) + throw new Error("API key ID is invalid"); + if (!/^[a-f0-9]{64}$/.test(key.tokenHash)) + throw new Error("API key token hash must be a SHA-256 hex digest"); + if (!key.username || key.username !== key.username.trim()) + throw new Error("API key username is required"); + if ( + !Array.isArray(key.capabilities) || + key.capabilities.length === 0 || + new Set(key.capabilities).size !== key.capabilities.length || + !key.capabilities.every(isCapability) + ) { + throw new Error("API key requires unique valid capabilities"); + } + if ( + key.workspace !== undefined && + (!/^[a-z0-9](?:[-a-z0-9]*[a-z0-9])?$/.test(key.workspace) || + key.workspace.length > 63) + ) { + throw new Error("API key workspace scope is invalid"); + } + const expiresAt = new Date(key.expiresAt); + if ( + !Number.isFinite(expiresAt.getTime()) || + expiresAt.toISOString() !== key.expiresAt + ) { + throw new Error("API key expiration must be an ISO timestamp"); + } + return { + ...key, + capabilities: [...key.capabilities], + disabled: Boolean(key.disabled), + }; +} + export function hashToken(token: string): string { return createHash("sha256").update(token).digest("hex"); } @@ -107,6 +165,7 @@ export function tokenHashesEqual(left: string, right: string): boolean { export class MemoryAuthStore implements AuthStore { readonly users = new Map(); readonly sessions = new Map(); + readonly apiKeys = new Map(); async getUser(username: string): Promise { return this.users.get(username); @@ -148,6 +207,8 @@ export class MemoryAuthStore implements AuthStore { async deleteUser(username: string): Promise { await this.revokeUserSessions(username); + for (const [id, key] of this.apiKeys) + if (key.username === username) this.apiKeys.delete(id); return this.users.delete(username); } @@ -201,4 +262,44 @@ export class MemoryAuthStore implements AuthStore { for (const session of expired) this.sessions.delete(session.tokenHash); return expired.length; } + + async getApiKey(tokenHash: string): Promise { + const key = [...this.apiKeys.values()].find( + (item) => item.tokenHash === tokenHash, + ); + if (!key || key.disabled || Date.parse(key.expiresAt) <= Date.now()) return; + const user = this.users.get(key.username); + if (!user || user.disabled) return; + return key; + } + + async createApiKey(key: NewApiKey): Promise { + const normalized = normalizeApiKey(key); + const user = this.users.get(normalized.username); + if (!user || user.disabled) throw new Error("API key user is not active"); + if (this.apiKeys.has(normalized.id)) + throw new Error("API key already exists"); + this.apiKeys.set(normalized.id, normalized); + } + + async listApiKeys(username: string): Promise { + return [...this.apiKeys.values()] + .filter((key) => key.username === username) + .sort((left, right) => left.id.localeCompare(right.id)); + } + + async revokeApiKey(username: string, id: string): Promise { + const key = this.apiKeys.get(id); + if (!key || key.username !== username) return false; + this.apiKeys.delete(id); + return true; + } + + async deleteExpiredApiKeys(now = Date.now()): Promise { + const expired = [...this.apiKeys.values()].filter( + (key) => Date.parse(key.expiresAt) <= now, + ); + for (const key of expired) this.apiKeys.delete(key.id); + return expired.length; + } } diff --git a/server/authorization.ts b/server/authorization.ts index 96370e7..9a7ce19 100644 --- a/server/authorization.ts +++ b/server/authorization.ts @@ -8,6 +8,7 @@ export const CAPABILITIES = [ "users:read", "users:write", "sessions:revoke", + "platform:adopt", ] as const; export type Capability = (typeof CAPABILITIES)[number]; @@ -23,6 +24,13 @@ export function isRole(value: unknown): value is Role { ); } +export function isCapability(value: unknown): value is Capability { + return ( + typeof value === "string" && + (CAPABILITIES as readonly string[]).includes(value) + ); +} + export function capabilitiesForRoles( roles: readonly string[], ): ReadonlySet { diff --git a/server/index.ts b/server/index.ts index bd0e855..79a3e69 100644 --- a/server/index.ts +++ b/server/index.ts @@ -172,7 +172,7 @@ const sessionCleanupTimer = setInterval( try { await cleanupExpiredSessions(store); } catch (error) { - console.error("Expired session cleanup failed", error); + console.error("Expired authentication cleanup failed", error); } }, Number.isFinite(sessionCleanupIntervalMs) && sessionCleanupIntervalMs > 0 diff --git a/server/kubernetes-state.ts b/server/kubernetes-state.ts index 6d33915..b197986 100644 --- a/server/kubernetes-state.ts +++ b/server/kubernetes-state.ts @@ -124,7 +124,10 @@ export function createKubernetesConfig(): KubeConfig { if (process.env.KUBERNETES_SERVICE_HOST) config.loadFromCluster(); else config.loadFromDefault(); const makeApiClient = config.makeApiClient.bind(config); - const httpLibrary = createKubernetesHttpLibrary(); + const httpLibrary = createKubernetesHttpLibrary({ + maxConcurrent: 4, + minIntervalMs: 0, + }); config.makeApiClient = ((apiClientType) => { const client = makeApiClient(apiClientType) as unknown as { api?: { configuration?: { httpApi?: typeof httpLibrary } }; diff --git a/server/kubernetes-store.ts b/server/kubernetes-store.ts index e6f38a0..44ca185 100644 --- a/server/kubernetes-store.ts +++ b/server/kubernetes-store.ts @@ -9,6 +9,8 @@ import { createKubernetesHttpLibrary } from "../lib/k8s-http"; import { normalizeSession, normalizeUser, + normalizeApiKey, + type ApiKeyRecord, type AuthStore, type KuberUser, type NewKuberUser, @@ -16,7 +18,7 @@ import { type SessionRecord, type UserUpdate, } from "./auth"; -import { isRole } from "./authorization"; +import { isCapability, isRole, type Capability } from "./authorization"; const FIELD_MANAGER = "kuber-server"; export const KUBER_SYSTEM_NAMESPACE = "kuber-system"; @@ -81,7 +83,10 @@ function parseRoles(value: unknown): KuberUser["roles"] | undefined { } } -function isSecret(secret: SecretObject, type: "user" | "session"): boolean { +function isSecret( + secret: SecretObject, + type: "user" | "session" | "api-key", +): boolean { return ( secret.apiVersion === "v1" && secret.kind === "Secret" && @@ -94,6 +99,24 @@ function isSecret(secret: SecretObject, type: "user" | "session"): boolean { ); } +function parseCapabilities(value: unknown): Capability[] | undefined { + const decoded = decode(value); + if (!decoded) return; + try { + const capabilities: unknown = JSON.parse(decoded); + if ( + !Array.isArray(capabilities) || + capabilities.length === 0 || + new Set(capabilities).size !== capabilities.length || + !capabilities.every(isCapability) + ) + return; + return capabilities; + } catch { + return; + } +} + function parseUser( secret: SecretObject, expectedUsername?: string, @@ -160,6 +183,51 @@ function parseSession(secret: SecretObject): SessionRecord | undefined { } } +function parseApiKey(secret: SecretObject): ApiKeyRecord | undefined { + if (!isSecret(secret, "api-key")) return; + const keys = [ + "id", + "tokenHash", + "username", + "capabilities", + "workspace", + "expiresAt", + "disabled", + ]; + if (!secret.data || !hasOnlyKeys(secret.data, keys)) return; + const id = decode(secret.data.id); + const tokenHash = decode(secret.data.tokenHash); + const username = decode(secret.data.username); + const capabilities = parseCapabilities(secret.data.capabilities); + const workspace = decode(secret.data.workspace); + const expiresAt = decode(secret.data.expiresAt); + const disabled = decode(secret.data.disabled); + if ( + !id || + !tokenHash || + !username || + !capabilities || + !expiresAt || + (workspace !== "" && workspace === undefined) || + (disabled !== "true" && disabled !== "false") || + secret.metadata?.name !== objectName("api-key", tokenHash) + ) + return; + try { + return normalizeApiKey({ + id, + tokenHash, + username, + capabilities, + ...(workspace && { workspace }), + expiresAt, + disabled: disabled === "true", + }); + } catch { + return; + } +} + function isNotFound(error: unknown): boolean { return Boolean( error && typeof error === "object" && "code" in error && error.code === 404, @@ -172,7 +240,10 @@ function createObjectApi(): KubernetesObjectApi { else config.loadFromDefault(); const makeApiClient = config.makeApiClient.bind(config); - const httpLibrary = createKubernetesHttpLibrary(); + const httpLibrary = createKubernetesHttpLibrary({ + maxConcurrent: 4, + minIntervalMs: 0, + }); config.makeApiClient = ((apiClientType) => { const client = makeApiClient(apiClientType) as unknown as { api?: { configuration?: { httpApi?: typeof httpLibrary } }; @@ -205,7 +276,7 @@ export class KubernetesAuthStore implements AuthStore { private async applySecret( name: string, - type: "user" | "session", + type: "user" | "session" | "api-key", stringData: Record, ): Promise { await this.objects.patch( @@ -228,7 +299,9 @@ export class KubernetesAuthStore implements AuthStore { ); } - private async listSecrets(type: "user" | "session"): Promise { + private async listSecrets( + type: "user" | "session" | "api-key", + ): Promise { const result = await this.objects.list( "v1", "Secret", @@ -309,6 +382,8 @@ export class KubernetesAuthStore implements AuthStore { async deleteUser(username: string): Promise { await this.revokeUserSessions(username); + for (const key of await this.listApiKeys(username)) + await this.deleteSecret(objectName("api-key", key.tokenHash)); return this.deleteSecret(objectName("user", username)); } @@ -378,4 +453,62 @@ export class KubernetesAuthStore implements AuthStore { } return expired.length; } + + async getApiKey(tokenHash: string): Promise { + if (!/^[a-f0-9]{64}$/.test(tokenHash)) return; + const key = (await this.listSecrets("api-key")) + .map(parseApiKey) + .find((item): item is ApiKeyRecord => item?.tokenHash === tokenHash); + if (!key || key.disabled || Date.parse(key.expiresAt) <= Date.now()) return; + const user = await this.getUser(key.username); + if (!user || user.disabled) return; + return key; + } + + async createApiKey(key: ApiKeyRecord): Promise { + const normalized = normalizeApiKey(key); + const user = await this.getUser(normalized.username); + if (!user || user.disabled) throw new Error("API key user is not active"); + await this.applySecret( + objectName("api-key", normalized.tokenHash), + "api-key", + { + id: normalized.id, + tokenHash: normalized.tokenHash, + username: normalized.username, + capabilities: JSON.stringify(normalized.capabilities), + workspace: normalized.workspace ?? "", + expiresAt: normalized.expiresAt, + disabled: String(Boolean(normalized.disabled)), + }, + ); + } + + async listApiKeys(username: string): Promise { + return (await this.listSecrets("api-key")) + .map(parseApiKey) + .filter((key): key is ApiKeyRecord => key?.username === username) + .sort((left, right) => left.id.localeCompare(right.id)); + } + + async revokeApiKey(username: string, id: string): Promise { + const key = (await this.listApiKeys(username)).find( + (item) => item.id === id, + ); + return key + ? this.deleteSecret(objectName("api-key", key.tokenHash)) + : false; + } + + async deleteExpiredApiKeys(now = Date.now()): Promise { + const expired = (await this.listSecrets("api-key")) + .map(parseApiKey) + .filter( + (key): key is ApiKeyRecord => + key !== undefined && Date.parse(key.expiresAt) <= now, + ); + for (const key of expired) + await this.deleteSecret(objectName("api-key", key.tokenHash)); + return expired.length; + } } diff --git a/server/management.ts b/server/management.ts index 8535a9b..43df58e 100644 --- a/server/management.ts +++ b/server/management.ts @@ -98,7 +98,20 @@ export type CredentialMetadata = { secretName?: string; }; -export type OperationExecution = { signal?: AbortSignal }; +export type OperationProgressEmitter = (event: { + resource: { + apiVersion: string; + kind: string; + name: string; + namespace?: string; + }; + phase: "apply" | "wait" | "delete"; + state: "started" | "succeeded" | "failed" | "aborted"; +}) => Promise; +export type OperationExecution = { + signal?: AbortSignal; + emit?: OperationProgressEmitter; +}; export type ManagementDependencies = { readNamespace(project: string): Promise; @@ -222,6 +235,20 @@ function resourceName(resource: KubernetesObject): string { return name; } +function resourceProgressIdentity( + resource: Pick< + ResourceIdentity, + "apiVersion" | "kind" | "name" | "namespace" + >, +) { + return { + apiVersion: resource.apiVersion, + kind: resource.kind, + name: resource.name, + ...(resource.namespace && { namespace: resource.namespace }), + }; +} + function assertResourceOwnership( workspace: Workspace, resource: KubernetesObject, @@ -426,8 +453,22 @@ export function createManagementService(dependencies: ManagementDependencies) { } } for (const resource of resources) { - throwIfExecutionAborted(execution); - await dependencies.deleteResource(resource, execution); + const event = { + resource: resourceProgressIdentity(resource), + phase: "delete" as const, + }; + try { + throwIfExecutionAborted(execution); + await execution?.emit?.({ ...event, state: "started" }); + await dependencies.deleteResource(resource, execution); + await execution?.emit?.({ ...event, state: "succeeded" }); + } catch (error) { + await execution?.emit?.({ + ...event, + state: execution?.signal?.aborted ? "aborted" : "failed", + }); + throw error; + } } } @@ -522,7 +563,12 @@ export function createManagementService(dependencies: ManagementDependencies) { } for (const name of selected) { throwIfExecutionAborted(execution); - await dependencies.scaleDeployment(workspace.project, name, 0, execution); + await dependencies.scaleDeployment( + workspace.project, + name, + 0, + execution, + ); } return selected; }, @@ -535,7 +581,11 @@ export function createManagementService(dependencies: ManagementDependencies) { const selected = await targets(workspace, names); for (const name of selected) { throwIfExecutionAborted(execution); - await dependencies.restartDeployment(workspace.project, name, execution); + await dependencies.restartDeployment( + workspace.project, + name, + execution, + ); } return selected; }, @@ -554,7 +604,11 @@ export function createManagementService(dependencies: ManagementDependencies) { ); for (const candidate of candidates) { throwIfExecutionAborted(execution); - await dependencies.rollbackDeployment(workspace.project, candidate, execution); + await dependencies.rollbackDeployment( + workspace.project, + candidate, + execution, + ); } for (const candidate of candidates) { throwIfExecutionAborted(execution); @@ -710,8 +764,27 @@ export function createManagementService(dependencies: ManagementDependencies) { for (const resource of sortResources( resources.map((item) => labelDesired(workspace, item)), )) { - throwIfExecutionAborted(execution); - applied.push(await dependencies.applyResource(resource, execution)); + const event = { + resource: resourceProgressIdentity({ + apiVersion: resource.apiVersion!, + kind: resource.kind!, + name: resourceName(resource), + namespace: resource.metadata?.namespace, + }), + phase: "apply" as const, + }; + try { + throwIfExecutionAborted(execution); + await execution?.emit?.({ ...event, state: "started" }); + applied.push(await dependencies.applyResource(resource, execution)); + await execution?.emit?.({ ...event, state: "succeeded" }); + } catch (error) { + await execution?.emit?.({ + ...event, + state: execution?.signal?.aborted ? "aborted" : "failed", + }); + throw error; + } } return applied; }, @@ -724,13 +797,32 @@ export function createManagementService(dependencies: ManagementDependencies) { ): Promise { const selected = await targets(workspace, deploymentTargets); for (const name of selected) { - throwIfExecutionAborted(execution); - await dependencies.waitForDeployment( - workspace.project, - name, - timeoutMs, - execution, - ); + const event = { + resource: { + apiVersion: "apps/v1", + kind: "Deployment", + name, + namespace: workspace.project, + }, + phase: "wait" as const, + }; + try { + throwIfExecutionAborted(execution); + await execution?.emit?.({ ...event, state: "started" }); + await dependencies.waitForDeployment( + workspace.project, + name, + timeoutMs, + execution, + ); + await execution?.emit?.({ ...event, state: "succeeded" }); + } catch (error) { + await execution?.emit?.({ + ...event, + state: execution?.signal?.aborted ? "aborted" : "failed", + }); + throw error; + } } }, diff --git a/server/operation-store.ts b/server/operation-store.ts index 488e9a7..066b50a 100644 --- a/server/operation-store.ts +++ b/server/operation-store.ts @@ -1,9 +1,25 @@ import { createHash, randomUUID } from "node:crypto"; import { KUBER_API_VERSION, type ObjectMeta } from "./workspace-store"; -import type { OperationError, OperationState } from "../shared/api"; +import type { + OperationError, + OperationEvent, + OperationState, +} from "../shared/api"; import { redactString, REDACTED } from "./redact"; export const MAX_IDEMPOTENCY_KEY_BYTES = 256; +export const MAX_OPERATION_EVENTS = 256; + +export type OperationProgress = { + resource: { + apiVersion: string; + kind: string; + name: string; + namespace?: string; + }; + phase: "apply" | "wait" | "delete"; + state: "started" | "succeeded" | "failed" | "aborted"; +}; export type { OperationError, OperationState } from "../shared/api"; @@ -23,6 +39,7 @@ export interface Operation { finishedAt?: string; result?: unknown; error?: OperationError; + events?: OperationEvent[]; }; } @@ -75,6 +92,14 @@ export interface OperationStore { state: OperationState, options?: { result?: unknown; error?: OperationError }, ): Promise; + emit(id: string, progress: OperationProgress): Promise; + events(id: string, after?: number): Promise; +} + +export interface OperationEvents { + items: OperationEvent[]; + retainedFirstSequence: number; + cursorGap: boolean; } export interface WorkspaceLease { @@ -330,10 +355,91 @@ export class PersistentOperationStore implements OperationStore { ...(options.error && { error: sanitizeOperationError(options.error, current.spec.action), }), + ...(current.status.events && { events: clone(current.status.events) }), }; await this.persistence.replace(operation, current.metadata.resourceVersion); return clone(operation); } + + async emit(id: string, progress: OperationProgress): Promise { + const data = sanitizeOperationProgress(progress); + for (let attempt = 0; attempt < 8; attempt++) { + const current = await this.persistence.get(id); + if (!current) + throw new OperationNotFoundError(`Operation '${id}' not found`); + const events = current.status.events ?? []; + const event: OperationEvent = { + operationId: id, + sequence: (events.at(-1)?.sequence ?? 0) + 1, + timestamp: this.now().toISOString(), + type: "progress", + data, + }; + const operation = clone(current); + operation.metadata.resourceVersion = String( + Number(current.metadata.resourceVersion) + 1, + ); + operation.status = { + ...operation.status, + events: [...events, event].slice(-MAX_OPERATION_EVENTS), + }; + try { + await this.persistence.replace( + operation, + current.metadata.resourceVersion, + ); + return clone(event); + } catch (error) { + if (!(error instanceof OperationConflictError) || attempt === 7) + throw error; + } + } + throw new OperationConflictError("Operation event could not be persisted"); + } + + async events(id: string, after = 0): Promise { + if (!Number.isSafeInteger(after) || after < 0) + throw new OperationValidationError( + "Operation event cursor must be a non-negative integer", + ); + const operation = await this.persistence.get(id); + if (!operation) + throw new OperationNotFoundError(`Operation '${id}' not found`); + const events = operation.status.events ?? []; + const retainedFirstSequence = events[0]?.sequence ?? 1; + return { + items: clone(events.filter((event) => event.sequence > after)), + retainedFirstSequence, + cursorGap: after < retainedFirstSequence - 1, + }; + } +} + +function sanitizeOperationProgress( + progress: OperationProgress, +): OperationProgress { + const resource = progress.resource; + if ( + !resource || + !resource.apiVersion || + !resource.kind || + !resource.name || + !["apply", "wait", "delete"].includes(progress.phase) || + !["started", "succeeded", "failed", "aborted"].includes(progress.state) + ) + throw new OperationValidationError("Invalid operation progress event"); + return { + resource: { + apiVersion: resource.apiVersion.slice(0, 128), + kind: resource.kind.slice(0, 128), + name: resource.name.slice(0, 253), + ...(resource.namespace && { + namespace: resource.namespace.slice(0, 253), + }), + }, + phase: progress.phase, + state: progress.state, + }; } export class MemoryOperationPersistence implements OperationPersistence { diff --git a/server/registry.ts b/server/registry.ts index 62e6955..f8aa8d3 100644 --- a/server/registry.ts +++ b/server/registry.ts @@ -21,8 +21,18 @@ export type RegistryResolveOptions = { credentials?: RegistryCredentials; insecure?: boolean; origin?: string; + cacheTtlMs?: number; + cacheMaxEntries?: number; + clock?: () => number; }; +const DEFAULT_DIGEST_CACHE_TTL_MS = 30_000; +const DEFAULT_DIGEST_CACHE_MAX_ENTRIES = 256; +const digestCache = new Map< + string, + { digest: Sha256Digest; expiresAt: number } +>(); + export type ParsedImageReference = { registry: string; repository: string; @@ -106,6 +116,13 @@ export async function resolveRegistryDigest( ): Promise { const parsed = parseImageReference(image); if (parsed.digest) return parsed.digest; + const now = options.clock ?? Date.now; + const cacheTtlMs = Number.isFinite(options.cacheTtlMs) + ? Math.max(0, options.cacheTtlMs!) + : DEFAULT_DIGEST_CACHE_TTL_MS; + const cacheMaxEntries = Number.isFinite(options.cacheMaxEntries) + ? Math.max(0, Math.floor(options.cacheMaxEntries!)) + : DEFAULT_DIGEST_CACHE_MAX_ENTRIES; const fetcher: RegistryFetch = options.fetch ?? globalThis.fetch; const scheme = options.insecure ? "http" : "https"; const repository = parsed.repository @@ -116,6 +133,19 @@ export async function resolveRegistryDigest( /\/+$/, "", ); + const credentialsKey = options.credentials + ? createHash("sha256") + .update( + `${options.credentials.username}\0${options.credentials.password}`, + ) + .digest("hex") + : "anonymous"; + const cacheKey = `${origin}\0${parsed.repository}\0${parsed.reference}\0${credentialsKey}`; + const cached = digestCache.get(cacheKey); + if (cached) { + if (cached.expiresAt > now()) return cached.digest; + digestCache.delete(cacheKey); + } const manifestUrl = `${origin}/v2/${repository}/manifests/${encodeURIComponent(parsed.reference)}`; const headers = new Headers({ accept: ACCEPT }); if (options.credentials) { @@ -167,9 +197,21 @@ export async function resolveRegistryDigest( .get("docker-content-digest") ?.trim() ?.toLowerCase(); + let digest: Sha256Digest; if (advertised !== undefined) { assertSha256Digest(advertised); - return advertised; + digest = advertised; + } else { + digest = `sha256:${createHash("sha256").update(body).digest("hex")}`; } - return `sha256:${createHash("sha256").update(body).digest("hex")}`; + if (cacheTtlMs && cacheMaxEntries) { + digestCache.delete(cacheKey); + while (digestCache.size >= cacheMaxEntries) + digestCache.delete(digestCache.keys().next().value!); + digestCache.set(cacheKey, { + digest, + expiresAt: now() + cacheTtlMs, + }); + } + return digest; } diff --git a/shared/api.ts b/shared/api.ts index 15ee09f..12e86c4 100644 --- a/shared/api.ts +++ b/shared/api.ts @@ -202,6 +202,14 @@ export type OperationEvent = { type: "status" | "progress" | "log" | "result" | "error"; data: JsonValue; }; +export type ListOperationEventsResponse = { + items: OperationEvent[]; + nextCursor?: number; + /** The oldest event sequence still retained for this operation. */ + retainedFirstSequence: number; + /** True when `after` predates the retained event history. */ + cursorGap: boolean; +}; export type UserRole = "admin" | "operator" | "viewer" | (string & {}); export type User = { @@ -223,6 +231,29 @@ export type UpdateUserRequest = { }; export type UserResponse = User; export type ListUsersResponse = Page; +export type Capability = + | "kubernetes:read" + | "kubernetes:write" + | "kubernetes:exec" + | "users:read" + | "users:write" + | "sessions:revoke" + | "platform:adopt"; +export type ApiKey = { + id: string; + username: string; + capabilities: Capability[]; + workspace?: string; + expiresAt: string; + disabled: boolean; +}; +export type CreateApiKeyRequest = { + capabilities: Capability[]; + workspace?: string; + expiresAt?: string; +}; +export type CreateApiKeyResponse = ApiKey & { token: string }; +export type ListApiKeysResponse = Page; export type AuditEvent = { id: string; diff --git a/tests/command/administration.test.ts b/tests/command/administration.test.ts index 6c5e182..b023b45 100644 --- a/tests/command/administration.test.ts +++ b/tests/command/administration.test.ts @@ -4,8 +4,11 @@ import { main } from "../../command/main"; import { getOperation, listOperations } from "../../command/operations"; import { addUser, + createApiKey, deleteUser, + listApiKeys, listUsers, + revokeApiKey, revokeUserSessions, setUserDisabled, updateUser, @@ -107,6 +110,54 @@ describe("user administration commands", () => { }, ]); }); + + test("manages API keys below the user route without leaking token in lists", async () => { + const calls: Call[] = []; + const key = { + id: "key-identifier-123", + username: "alice", + capabilities: ["kubernetes:write"] as const, + expiresAt: "2026-12-01T00:00:00.000Z", + disabled: false, + token: "shown-once-token", + }; + expect( + await listApiKeys( + "alice/example", + requestReturning({ items: [{ ...key, token: undefined }] }, calls), + ), + ).not.toContain(key.token); + await createApiKey( + "alice", + { capabilities: ["kubernetes:write"] }, + requestReturning(key, calls), + ); + expect( + await revokeApiKey( + "alice", + key.id, + false, + requestReturning(undefined, calls), + ), + ).toContain("cancelled"); + await revokeApiKey( + "alice", + key.id, + true, + requestReturning(undefined, calls), + ); + expect(calls).toEqual([ + { path: "/users/alice%2Fexample/keys", init: undefined }, + { + path: "/users/alice/keys", + init: { method: "POST", json: { capabilities: ["kubernetes:write"] } }, + }, + { + path: "/users/alice/keys/key-identifier-123", + init: { method: "DELETE" }, + }, + ]); + }); }); const operation = { diff --git a/tests/command/ci.test.ts b/tests/command/ci.test.ts new file mode 100644 index 0000000..5a49668 --- /dev/null +++ b/tests/command/ci.test.ts @@ -0,0 +1,19 @@ +import { expect, spyOn, test } from "bun:test"; +import { apiKeyRequest } from "../../command/ci"; + +test("CI requester supplies an API key without session authentication", async () => { + const fetch = spyOn(globalThis, "fetch").mockResolvedValue( + new Response(JSON.stringify({ ok: true }), { status: 200 }), + ); + try { + await apiKeyRequest("ci-secret")("/workspaces/shop/resources/plan", { + headers: { "x-kuber-trust-project": "shop" }, + }); + const [, init] = fetch.mock.calls[0]!; + const headers = new Headers(init?.headers); + expect(headers.get("authorization")).toBe("Bearer ci-secret"); + expect(headers.get("x-kuber-trust-project")).toBe("shop"); + } finally { + fetch.mockRestore(); + } +}); diff --git a/tests/command/up-api.test.ts b/tests/command/up-api.test.ts index 1ea6d80..cceb955 100644 --- a/tests/command/up-api.test.ts +++ b/tests/command/up-api.test.ts @@ -1,4 +1,5 @@ import { describe, expect, test } from "bun:test"; +import { Listr } from "listr2"; import { mkdtemp, rm, writeFile } from "node:fs/promises"; import { tmpdir } from "node:os"; import { join } from "node:path"; @@ -8,6 +9,7 @@ import type { ApiRequester } from "../../lib/build"; import { ensureWorkspace, reconcileResources, + runLiveResourceOperation, runUp, workspaceAdoptionRoute, } from "../../command/up"; @@ -226,8 +228,10 @@ describe("up API pipeline", () => { "/workspaces/shop/resources/plan", "/workspaces/shop/resources/apply", "/workspaces/shop/resources/apply", + "/operations/operation-apply/events?after=0", "/operations/operation-apply", "/operations/operation-apply", + "/operations/operation-apply/events?after=0", ]); expect(applyIdempotencyKeys[0]).toBeTruthy(); expect(applyIdempotencyKeys[1]).toBe(applyIdempotencyKeys[0]); @@ -272,6 +276,161 @@ describe("up API pipeline", () => { expect(applyRequests[2]?.json).toEqual(applyRequests[0]?.json); }); + test("polls persisted resource progress without duplicating events after reconnect", async () => { + const progress: string[] = []; + let poll = 0; + const request: ApiRequester = async (path: string) => { + if (path.endsWith("/plan")) return { desired: [], stale: [] } as T; + if (path.endsWith("/apply")) + return { + operationId: "operation-apply", + operation: { status: { state: "running" } }, + } as T; + if (path.includes("/events")) { + poll++; + return { + items: + poll === 1 + ? [ + { + sequence: 1, + data: null, + }, + { + sequence: 2, + data: { + resource: { + apiVersion: "v1", + kind: "Service", + name: "web", + }, + phase: "apply", + state: "started", + }, + }, + ] + : [ + { + sequence: 2, + data: { + resource: { + apiVersion: "v1", + kind: "Service", + name: "web", + }, + phase: "apply", + state: "started", + }, + }, + { + sequence: 3, + data: { + resource: { + apiVersion: "v1", + kind: "Service", + name: "web", + }, + phase: "apply", + state: "succeeded", + }, + }, + ], + } as T; + } + if (path === "/operations/operation-apply") { + return { status: { state: poll > 1 ? "succeeded" : "running" } } as T; + } + throw new Error(`Unexpected request: ${path}`); + }; + + await reconcileResources("shop", [], 1, undefined, request, { + sleep: async () => {}, + onEvent: (event) => + progress.push(`${event.sequence}:${event.data.state}`), + }); + expect(progress).toEqual(["2:started", "3:succeeded"]); + }); + + test("fails visibly when retained operation progress has a cursor gap", async () => { + let applyAttempts = 0; + let operationPolls = 0; + const request: ApiRequester = async (path: string) => { + if (path.endsWith("/plan")) return { desired: [], stale: [] } as T; + if (path.endsWith("/apply")) { + applyAttempts += 1; + return { + operationId: "operation-apply", + operation: { status: { state: "running" } }, + } as T; + } + if (path.includes("/events")) + return { + retainedFirstSequence: 3, + cursorGap: true, + items: [], + } as T; + if (path === "/operations/operation-apply") { + operationPolls += 1; + return { status: { state: "succeeded" } } as T; + } + throw new Error(`Unexpected request: ${path}`); + }; + + await expect( + reconcileResources("shop", [], 1, undefined, request, { + sleep: async () => {}, + }), + ).rejects.toThrow("progress history was truncated"); + expect(applyAttempts).toBe(1); + expect(operationPolls).toBe(0); + }); + + test("starts live resource subtasks before the operation and updates them from progress", async () => { + let finish!: () => void; + let child: { title: string; output: string } | undefined; + const completed = new Promise((resolve) => (finish = resolve)); + const listr = new Listr([ + { + title: "Apply resources", + task: async (_ctx, task) => + runLiveResourceOperation( + task, + "apply", + [ + { + apiVersion: "v1", + kind: "Service", + name: "web", + }, + ], + async (onEvent) => { + expect(child?.title).toBe("Apply Service/web"); + onEvent({ + sequence: 1, + data: { + resource: { + apiVersion: "v1", + kind: "Service", + name: "web", + }, + phase: "apply", + state: "started", + }, + }); + expect(child?.output).toBe("started"); + await completed; + }, + { onTaskStarted: (_target, activeTask) => (child = activeTask) }, + ), + }, + ]); + + const run = listr.run(); + await Bun.sleep(0); + finish(); + await run; + }); + test("resubmits with a fresh key after server restart interruption", async () => { const applyRequests: Array<{ key: string | null; json: unknown }> = []; let operationPolls = 0; diff --git a/tests/server/api-keys.test.ts b/tests/server/api-keys.test.ts new file mode 100644 index 0000000..0c55004 --- /dev/null +++ b/tests/server/api-keys.test.ts @@ -0,0 +1,463 @@ +import { describe, expect, test } from "bun:test"; +import { cleanupExpiredSessions, createApp } from "../../server/app"; +import { hashToken, MemoryAuthStore } from "../../server/auth"; +import { MemoryAuditStore } from "../../server/audit-store"; +import { MemoryOperationStore } from "../../server/operation-store"; + +const now = Date.parse("2026-09-05T00:00:00.000Z"); + +function request(path: string, init: RequestInit = {}, token = "admin-token") { + const headers = new Headers(init.headers); + headers.set("authorization", `Bearer ${token}`); + return new Request(`https://kuber.astrxl.dev${path}`, { ...init, headers }); +} + +async function setup() { + const store = new MemoryAuthStore(); + const auditStore = new MemoryAuditStore(() => new Date(now)); + const operationStore = new MemoryOperationStore(() => new Date(now)); + await store.putUser({ + username: "admin", + passwordHash: "hash", + roles: ["admin"], + }); + await store.putUser({ + username: "ci", + passwordHash: "hash", + roles: ["operator"], + }); + await store.putSession({ + tokenHash: hashToken("admin-token"), + username: "admin", + authVersion: 1, + expiresAt: "2026-10-05T00:00:00.000Z", + }); + return { + store, + auditStore, + operationStore, + app: createApp({ store, auditStore, operationStore, now: () => now }), + }; +} + +describe("API keys", () => { + test("creates once, lists without a token or hash, and revokes", async () => { + const { app, auditStore } = await setup(); + const created = await app( + request("/api/v2/users/ci/keys", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ + capabilities: ["kubernetes:read"], + workspace: "shop", + }), + }), + ); + expect(created.status).toBe(201); + const key = (await created.json()) as { id: string; token: string }; + expect(key.token).toHaveLength(43); + + const listed = await app(request("/api/v2/users/ci/keys")); + const body = JSON.stringify(await listed.json()); + expect(body).not.toContain(key.token); + expect(body).not.toContain(hashToken(key.token)); + expect(JSON.stringify(await auditStore.list())).not.toContain(key.token); + expect(JSON.stringify(await auditStore.list())).not.toContain( + hashToken(key.token), + ); + + expect( + ( + await app( + request(`/api/v2/users/ci/keys/${key.id}`, { method: "DELETE" }), + ) + ).status, + ).toBe(204); + expect((await app(request("/api/v2/me", {}, key.token))).status).toBe(401); + }); + + test("uses key capabilities rather than owner roles and invalidates disabled owners", async () => { + const { app, store } = await setup(); + await store.createApiKey({ + id: "key_capability_test", + tokenHash: hashToken("ci-key"), + username: "ci", + capabilities: ["kubernetes:read"], + expiresAt: "2026-10-05T00:00:00.000Z", + }); + expect((await app(request("/api/v2/users", {}, "ci-key"))).status).toBe( + 403, + ); + expect((await app(request("/api/v2/me", {}, "ci-key"))).status).toBe(200); + await store.updateUser("ci", { disabled: true }); + expect((await app(request("/api/v2/me", {}, "ci-key"))).status).toBe(401); + }); + + test("limits API key children to the parent's user, capabilities, and workspace", async () => { + const { app, store, auditStore } = await setup(); + await store.createApiKey({ + id: "key_delegation_parent", + tokenHash: hashToken("delegation-parent"), + username: "ci", + capabilities: ["users:write", "kubernetes:read"], + workspace: "shop", + expiresAt: "2026-10-05T00:00:00.000Z", + }); + const create = (body: unknown) => + app( + request( + "/api/v2/users/ci/keys", + { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify(body), + }, + "delegation-parent", + ), + ); + + const capabilities = await create({ + capabilities: ["kubernetes:write"], + workspace: "shop", + }); + expect(capabilities.status).toBe(403); + const capabilityError = (await capabilities.json()) as { code: string }; + expect(capabilityError.code).toBe("API_KEY_DELEGATION_FORBIDDEN"); + + const workspace = await create({ capabilities: ["kubernetes:read"] }); + expect(workspace.status).toBe(403); + + const differentWorkspace = await create({ + capabilities: ["kubernetes:read"], + workspace: "other", + }); + expect(differentWorkspace.status).toBe(403); + + const subset = await create({ + capabilities: ["kubernetes:read"], + workspace: "shop", + }); + expect(subset.status).toBe(201); + expect( + ((await subset.json()) as { capabilities: string[] }).capabilities, + ).toEqual(["kubernetes:read"]); + + const otherUser = await app( + request( + "/api/v2/users/admin/keys", + { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ + capabilities: ["kubernetes:read"], + workspace: "shop", + }), + }, + "delegation-parent", + ), + ); + expect(otherUser.status).toBe(403); + + const denied = await auditStore.list(); + expect( + denied.filter( + (event) => + event.spec.action === "api_key.create" && + event.spec.outcome === "denied", + ), + ).toHaveLength(4); + expect(JSON.stringify(denied)).not.toContain("delegation-parent"); + + await store.createApiKey({ + id: "key_unscoped_delegation", + tokenHash: hashToken("unscoped-delegation"), + username: "ci", + capabilities: ["users:write", "kubernetes:read"], + expiresAt: "2026-10-05T00:00:00.000Z", + }); + const unscopedSubset = await app( + request( + "/api/v2/users/ci/keys", + { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ capabilities: ["kubernetes:read"] }), + }, + "unscoped-delegation", + ), + ); + expect(unscopedSubset.status).toBe(201); + expect(await unscopedSubset.json()).not.toHaveProperty("workspace"); + }); + + test("rejects expired keys and expiry longer than 365 days", async () => { + const { app, store } = await setup(); + await store.createApiKey({ + id: "key_expiry_test_1", + tokenHash: hashToken("expired-key"), + username: "ci", + capabilities: ["kubernetes:read"], + expiresAt: "2026-09-04T00:00:00.000Z", + }); + expect((await app(request("/api/v2/me", {}, "expired-key"))).status).toBe( + 401, + ); + expect( + ( + await app( + request("/api/v2/users/ci/keys", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ + capabilities: ["kubernetes:read"], + expiresAt: "2027-09-06T00:00:00.000Z", + }), + }), + ) + ).status, + ).toBe(400); + }); + + test("uses the default expiry, cleans up expired keys, and does not log keys out", async () => { + const { app, store } = await setup(); + const created = await app( + request("/api/v2/users/ci/keys", { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ capabilities: ["kubernetes:read"] }), + }), + ); + const key = (await created.json()) as { token: string; expiresAt: string }; + expect(key.expiresAt).toBe("2026-12-04T00:00:00.000Z"); + expect( + (await app(request("/api/v2/logout", { method: "POST" }, key.token))) + .status, + ).toBe(204); + expect((await app(request("/api/v2/me", {}, key.token))).status).toBe(200); + expect(await cleanupExpiredSessions(store, Date.parse(key.expiresAt))).toBe( + 2, + ); + expect((await app(request("/api/v2/me", {}, key.token))).status).toBe(401); + }); + + test("denies a workspace-scoped key outside its workspace", async () => { + const { app, store } = await setup(); + await store.createApiKey({ + id: "key_scope_test_1", + tokenHash: hashToken("scoped-key"), + username: "ci", + capabilities: ["kubernetes:read"], + workspace: "shop", + expiresAt: "2026-10-05T00:00:00.000Z", + }); + expect( + (await app(request("/api/v2/workspaces/other", {}, "scoped-key"))).status, + ).toBe(403); + }); + + test("limits workspace-scoped keys to their own audit records", async () => { + const { app, store, auditStore } = await setup(); + await store.createApiKey({ + id: "key_audit_scope_1", + tokenHash: hashToken("scoped-audit-key"), + username: "ci", + capabilities: ["users:read"], + workspace: "shop", + expiresAt: "2026-10-05T00:00:00.000Z", + }); + await auditStore.append({ + actor: { username: "admin" }, + action: "workspace.shop", + workspaceId: "shop", + outcome: "success", + }); + await auditStore.append({ + actor: { username: "admin" }, + action: "workspace.other", + workspaceId: "other", + outcome: "success", + }); + await auditStore.append({ + actor: { username: "admin" }, + action: "platform.global", + outcome: "success", + }); + + const unfiltered = await app( + request("/api/v2/audit", {}, "scoped-audit-key"), + ); + expect(unfiltered.status).toBe(200); + const audit = (await unfiltered.json()) as { + items: { spec: { action: string } }[]; + }; + expect(audit.items.map((event) => event.spec.action)).toEqual([ + "workspace.shop", + ]); + expect( + ( + await app( + request("/api/v2/audit?workspaceId=other", {}, "scoped-audit-key"), + ) + ).status, + ).toBe(403); + }); + + test("automatically scopes unfiltered operation lists for workspace keys", async () => { + const { app, store, operationStore } = await setup(); + await operationStore.create({ + workspaceId: "shop", + action: "resources.apply", + idempotencyKey: "shop-operation", + }); + await operationStore.create({ + workspaceId: "other", + action: "resources.apply", + idempotencyKey: "other-operation", + }); + await store.createApiKey({ + id: "key_operation_scope_1", + tokenHash: hashToken("scoped-operation-key"), + username: "ci", + capabilities: ["kubernetes:read"], + workspace: "shop", + expiresAt: "2026-10-05T00:00:00.000Z", + }); + + const unfiltered = await app( + request("/api/v2/operations", {}, "scoped-operation-key"), + ); + expect(unfiltered.status).toBe(200); + expect( + ( + (await unfiltered.json()) as { + items: { spec: { workspaceId: string } }[]; + } + ).items.map((operation) => operation.spec.workspaceId), + ).toEqual(["shop"]); + expect( + ( + await app( + request( + "/api/v2/operations?workspaceId=shop", + {}, + "scoped-operation-key", + ), + ) + ).status, + ).toBe(200); + expect( + ( + await app( + request( + "/api/v2/operations?workspaceId=other", + {}, + "scoped-operation-key", + ), + ) + ).status, + ).toBe(403); + }); + + test("does not inherit an admin owner's platform adoption privilege", async () => { + const { store } = await setup(); + const adopted: string[] = []; + const app = createApp({ + store, + adoption: { + adopt: async () => ({ + workspaceId: "", + workspaceUid: "", + resourcesAdopted: 0, + }), + adoptPlatform: async (workspaceUid) => { + adopted.push(workspaceUid); + return { + workspaceId: "kuber-system", + workspaceUid, + resourcesAdopted: 1, + }; + }, + }, + now: () => now, + }); + await store.createApiKey({ + id: "key_platform_owner", + tokenHash: hashToken("admin-owner-key"), + username: "admin", + capabilities: ["kubernetes:write"], + workspace: "kuber-system", + expiresAt: "2026-10-05T00:00:00.000Z", + }); + await store.createApiKey({ + id: "key_platform_scope", + tokenHash: hashToken("wrong-scope-key"), + username: "admin", + capabilities: ["platform:adopt"], + workspace: "shop", + expiresAt: "2026-10-05T00:00:00.000Z", + }); + await store.createApiKey({ + id: "key_platform_allowed", + tokenHash: hashToken("platform-key"), + username: "admin", + capabilities: ["platform:adopt"], + workspace: "kuber-system", + expiresAt: "2026-10-05T00:00:00.000Z", + }); + const adopt = (token: string) => + app( + request( + "/api/v2/platform/kuber-system/adopt", + { + method: "POST", + body: JSON.stringify({ workspaceUid: "platform" }), + }, + token, + ), + ); + + expect((await adopt("admin-owner-key")).status).toBe(403); + expect((await adopt("wrong-scope-key")).status).toBe(403); + expect((await adopt("platform-key")).status).toBe(200); + expect(adopted).toEqual(["platform"]); + }); + + test("allows a workspace-scoped key to use matching project build routes", async () => { + const { app, store } = await setup(); + await store.createApiKey({ + id: "key_scope_build_1", + tokenHash: hashToken("scoped-build-key"), + username: "ci", + capabilities: ["kubernetes:write"], + workspace: "shop", + expiresAt: "2026-10-05T00:00:00.000Z", + }); + const matching = await app( + request( + "/api/v2/images/resolve", + { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ project: "shop", service: "web" }), + }, + "scoped-build-key", + ), + ); + expect(matching.status).toBe(503); + expect( + ( + await app( + request( + "/api/v2/images/resolve", + { + method: "POST", + headers: { "content-type": "application/json" }, + body: JSON.stringify({ project: "other", service: "web" }), + }, + "scoped-build-key", + ), + ) + ).status, + ).toBe(403); + }); +}); diff --git a/tests/server/app.test.ts b/tests/server/app.test.ts index 47dccff..302dbb0 100644 --- a/tests/server/app.test.ts +++ b/tests/server/app.test.ts @@ -391,6 +391,94 @@ describe("kuber v2 HTTP routes", () => { expect((await apply(app)).status).toBe(403); }); + test("starts preferred resource operations before returning so progress can be polled", async () => { + const workspaceStore = new MemoryWorkspaceStore({ + uid: () => "workspace-uid", + }); + await workspaceStore.create({ + id: "demo", + source: { uri: "oci://example/demo", digest: "sha256:abc" }, + }); + const trustStore = new MemoryTrustStore(); + const fingerprint = "b".repeat(64); + let start!: () => void; + let finish!: () => void; + let complete!: () => void; + const started = new Promise((resolve) => (start = resolve)); + const unblock = new Promise((resolve) => (finish = resolve)); + const completed = new Promise((resolve) => (complete = resolve)); + const management = { + applyResources: async ( + _workspace: unknown, + _resources: unknown, + execution: { + emit?: (event: { + resource: { apiVersion: string; kind: string; name: string }; + phase: "apply"; + state: "started" | "succeeded"; + }) => Promise; + }, + ) => { + await execution.emit?.({ + resource: { apiVersion: "v1", kind: "Service", name: "web" }, + phase: "apply", + state: "started", + }); + start(); + await unblock; + await execution.emit?.({ + resource: { apiVersion: "v1", kind: "Service", name: "web" }, + phase: "apply", + state: "succeeded", + }); + complete(); + return []; + }, + } as unknown as ManagementService; + const app = createApp({ + store: await authenticatedStore("operator"), + workspaceStore, + operationStore: new MemoryOperationStore(), + trustStore, + management, + }); + await app( + request( + "/api/v2/workspaces/demo/trust", + { method: "POST", body: JSON.stringify({ fingerprint }) }, + "token", + ), + ); + + const submitted = await app( + request( + "/api/v2/workspaces/demo/resources/apply", + { + method: "POST", + headers: { + "idempotency-key": "progress-once", + prefer: "respond-async", + "x-kuber-trust-project": "demo", + "x-kuber-trust-fingerprint": fingerprint, + }, + body: JSON.stringify({ resources: [] }), + }, + "token", + ), + ); + expect(submitted.status).toBe(202); + const { operationId } = (await submitted.json()) as { operationId: string }; + await started; + const events = await app( + request(`/api/v2/operations/${operationId}/events?after=0`, {}, "token"), + ); + expect(await events.json()).toMatchObject({ + items: [{ sequence: 1, data: { state: "started" } }], + }); + finish(); + await completed; + }); + test("uses exact origins, request IDs, and problem+json errors", async () => { const app = createApp({ store: new MemoryAuthStore(), @@ -548,6 +636,108 @@ describe("kuber v2 HTTP routes", () => { expect(await operationStore.list("demo")).toHaveLength(1); }); + test("lists persisted operation events with capability and workspace scope checks", async () => { + const store = await authenticatedStore("operator"); + const workspaceStore = new MemoryWorkspaceStore({ + uid: () => "workspace-uid", + }); + await workspaceStore.create({ + id: "demo", + source: { uri: "oci://example/demo", digest: "sha256:abc" }, + }); + const operationStore = new MemoryOperationStore( + undefined, + () => "event-id", + ); + const operation = await operationStore.create({ + workspaceId: "demo", + action: "resources.apply", + idempotencyKey: "event-key", + }); + await operationStore.emit(operation.metadata.name, { + resource: { apiVersion: "v1", kind: "Service", name: "web" }, + phase: "apply", + state: "succeeded", + }); + await store.createApiKey({ + id: "demo-reader-key-01", + tokenHash: hashToken("demo-reader-token"), + username: "operator", + capabilities: ["kubernetes:read"], + workspace: "demo", + expiresAt: "2030-01-01T00:00:00.000Z", + }); + await store.createApiKey({ + id: "other-reader-key-1", + tokenHash: hashToken("other-reader-token"), + username: "operator", + capabilities: ["kubernetes:read"], + workspace: "other", + expiresAt: "2030-01-01T00:00:00.000Z", + }); + const app = createApp({ store, workspaceStore, operationStore }); + const path = `/api/v2/operations/${operation.metadata.name}/events?after=0`; + const allowed = await app(request(path, {}, "demo-reader-token")); + expect(allowed.status).toBe(200); + expect(await allowed.json()).toMatchObject({ + items: [ + { + sequence: 1, + data: { resource: { kind: "Service", name: "web" }, phase: "apply" }, + }, + ], + nextCursor: 1, + retainedFirstSequence: 1, + cursorGap: false, + }); + expect((await app(request(path, {}, "other-reader-token"))).status).toBe( + 403, + ); + expect( + (await app(request(`${path}x`, {}, "demo-reader-token"))).status, + ).toBe(400); + }); + + test("reports operation event cursor gaps after retained history is truncated", async () => { + const operationStore = new MemoryOperationStore( + undefined, + () => "retained", + ); + const operation = await operationStore.create({ + workspaceId: "demo", + action: "resources.apply", + idempotencyKey: "retained-events", + }); + for (let sequence = 0; sequence < 257; sequence += 1) { + await operationStore.emit(operation.metadata.name, { + resource: { apiVersion: "v1", kind: "ConfigMap", name: "config" }, + phase: "apply", + state: "succeeded", + }); + } + const app = createApp({ + store: await authenticatedStore("operator"), + operationStore, + }); + + const response = await app( + request( + `/api/v2/operations/${operation.metadata.name}/events?after=0`, + {}, + "token", + ), + ); + expect(response.status).toBe(200); + const events = (await response.json()) as { + retainedFirstSequence: number; + cursorGap: boolean; + items: { sequence: number }[]; + }; + expect(events.retainedFirstSequence).toBe(2); + expect(events.cursorGap).toBe(true); + expect(events.items[0]?.sequence).toBe(2); + }); + test("rejects JSON bodies over the configured limit", async () => { const app = createApp({ store: await authenticatedStore("admin"), diff --git a/tests/server/kubernetes-state.test.ts b/tests/server/kubernetes-state.test.ts index 02df54d..e168f67 100644 --- a/tests/server/kubernetes-state.test.ts +++ b/tests/server/kubernetes-state.test.ts @@ -14,7 +14,10 @@ import { KubernetesWorkspacePersistence, type LeaseObjects, } from "../../server/kubernetes-state"; -import type { Operation } from "../../server/operation-store"; +import { + PersistentOperationStore, + type Operation, +} from "../../server/operation-store"; import type { Workspace, WorkspaceRevision, @@ -183,6 +186,42 @@ describe("Kubernetes state persistence", () => { ); }); + test("keeps operation progress in the Kubernetes-backed operation record", async () => { + const fake = new FakeObjects(); + const store = new PersistentOperationStore( + new KubernetesOperationPersistence( + fake as unknown as KubernetesObjectApi, + ), + () => new Date(timestamp), + () => "event-id", + ); + const created = await store.create({ + workspaceId: "demo", + action: "resources.apply", + idempotencyKey: "events", + }); + await store.emit(created.metadata.name, { + resource: { apiVersion: "v1", kind: "ConfigMap", name: "settings" }, + phase: "apply", + state: "succeeded", + }); + const reloaded = new PersistentOperationStore( + new KubernetesOperationPersistence( + fake as unknown as KubernetesObjectApi, + ), + ); + expect(await reloaded.events(created.metadata.name)).toMatchObject({ + items: [ + expect.objectContaining({ + sequence: 1, + data: expect.objectContaining({ phase: "apply" }), + }), + ], + retainedFirstSequence: 1, + cursorGap: false, + }); + }); + test("recovers replacement from a matching precreated revision", async () => { const fake = new FakeObjects(); const persistence = new KubernetesWorkspacePersistence( @@ -365,7 +404,11 @@ class FakeLeaseStore implements LeaseObjects { return structuredClone(stored); } - async delete(name: string, namespace: string, expectedResourceVersion?: string) { + async delete( + name: string, + namespace: string, + expectedResourceVersion?: string, + ) { this.beforeDelete?.(); this.beforeDelete = undefined; const key = this.key(name, namespace); diff --git a/tests/server/kubernetes-store.test.ts b/tests/server/kubernetes-store.test.ts index bc92359..ef0174e 100644 --- a/tests/server/kubernetes-store.test.ts +++ b/tests/server/kubernetes-store.test.ts @@ -49,8 +49,10 @@ class FakeObjects { readonly secrets = new Map(); readonly patches: StoredSecret[] = []; readonly deleted: string[] = []; + reads = 0; async read(value: KubernetesObject): Promise { + this.reads += 1; const found = this.secrets.get(value.metadata?.name ?? ""); if (!found) throw { code: 404 }; return found; @@ -97,7 +99,10 @@ class FakeObjects { } } -function setup(): { fake: FakeObjects; store: KubernetesAuthStore } { +function setup(): { + fake: FakeObjects; + store: KubernetesAuthStore; +} { const fake = new FakeObjects(); return { fake, @@ -106,6 +111,39 @@ function setup(): { fake: FakeObjects; store: KubernetesAuthStore } { } describe("KubernetesAuthStore", () => { + test("validates the session and user from the store on every request", async () => { + const { fake, store } = setup(); + const tokenHash = hashToken("fresh"); + await store.putUser({ + username: "alice", + passwordHash: "hash", + roles: ["viewer"], + }); + await store.putSession({ + tokenHash, + username: "alice", + roles: ["viewer"], + expiresAt: "2026-09-03T00:00:00.000Z", + }); + fake.reads = 0; + + await expect(store.getSession(tokenHash)).resolves.toMatchObject({ + tokenHash, + }); + fake.secrets.set( + objectName("user", "alice"), + secret("user", objectName("user", "alice"), { + username: "alice", + passwordHash: "hash", + roles: JSON.stringify(["viewer"]), + authVersion: "1", + disabled: "true", + }), + ); + await expect(store.getSession(tokenHash)).resolves.toBeUndefined(); + expect(fake.reads).toBe(4); + }); + test("persists authVersion and not copied roles in new sessions", async () => { const { fake, store } = setup(); await store.putUser({ diff --git a/tests/server/management.test.ts b/tests/server/management.test.ts index 2e31129..fdae5db 100644 --- a/tests/server/management.test.ts +++ b/tests/server/management.test.ts @@ -163,6 +163,79 @@ describe("server management service", () => { expect(deleted).toEqual(plan.stale); }); + test("emits safe per-resource apply, wait, delete, and failure progress", async () => { + const events: Array<{ + phase: string; + state: string; + resource: { name: string }; + }> = []; + const service = createManagementService( + dependencies({ + listDeployments: async () => [ + object("Deployment", "web", "web-uid") as V1Deployment, + ], + applyResource: async (resource) => resource, + deleteResource: async () => {}, + }), + ); + const execution = { + emit: async (event: (typeof events)[number]) => { + events.push(event); + }, + }; + await service.applyResources( + workspace, + [object("Service", "web", "service-uid")], + execution, + ); + await service.waitForResources(workspace, ["web"], undefined, execution); + await service.deleteResources( + workspace, + [ + { + apiVersion: "v1", + kind: "Service", + name: "old", + uid: "old-uid", + workspaceUid: workspace.uid, + }, + ], + execution, + ); + expect( + events.map( + ({ phase, state, resource }) => `${phase}:${resource.name}:${state}`, + ), + ).toEqual([ + "apply:web:started", + "apply:web:succeeded", + "wait:web:started", + "wait:web:succeeded", + "delete:old:started", + "delete:old:succeeded", + ]); + + const failing = createManagementService( + dependencies({ + applyResource: async () => { + throw new Error("provider failed"); + }, + }), + ); + await expect( + failing.applyResources( + workspace, + [object("Service", "broken", "broken-uid")], + execution, + ), + ).rejects.toThrow("provider failed"); + expect(events.at(-1)).toMatchObject({ + phase: "apply", + state: "failed", + resource: { name: "broken" }, + }); + }); + test("rejects namespace and resource ownership mismatches", async () => { const wrongNamespace = createManagementService( dependencies({ diff --git a/tests/server/operation-store.test.ts b/tests/server/operation-store.test.ts index dacc710..c6f3dd8 100644 --- a/tests/server/operation-store.test.ts +++ b/tests/server/operation-store.test.ts @@ -2,6 +2,7 @@ import { describe, expect, test } from "bun:test"; import { MemoryOperationStore, MemoryWorkspaceLeaseProvider, + MAX_OPERATION_EVENTS, OperationConflictError, OperationValidationError, recoverStaleOperations, @@ -153,6 +154,58 @@ describe("operation store", () => { ).rejects.toBeInstanceOf(OperationValidationError); }); + test("persists monotonic bounded progress events with cursors", async () => { + const store = new MemoryOperationStore(undefined, () => "event-operation"); + const operation = await store.create({ + workspaceId: "demo", + action: "resources.apply", + idempotencyKey: "events", + }); + await store.emit(operation.metadata.name, { + resource: { apiVersion: "v1", kind: "Secret", name: "credentials" }, + phase: "apply", + state: "started", + }); + const complete = await store.emit(operation.metadata.name, { + resource: { apiVersion: "v1", kind: "Secret", name: "credentials" }, + phase: "apply", + state: "succeeded", + }); + expect(complete.sequence).toBe(2); + expect(await store.events(operation.metadata.name, 1)).toMatchObject({ + items: [expect.objectContaining({ sequence: 2 })], + retainedFirstSequence: 1, + cursorGap: false, + }); + await store.transition(operation.metadata.name, "running"); + await store.transition(operation.metadata.name, "succeeded"); + expect(await store.events(operation.metadata.name, 1)).toMatchObject({ + items: [expect.objectContaining({ sequence: 2 })], + retainedFirstSequence: 1, + cursorGap: false, + }); + await expect( + store.events(operation.metadata.name, -1), + ).rejects.toBeInstanceOf(OperationValidationError); + for (let index = 0; index < MAX_OPERATION_EVENTS; index++) { + await store.emit(operation.metadata.name, { + resource: { apiVersion: "v1", kind: "Service", name: `web-${index}` }, + phase: "apply", + state: "succeeded", + }); + } + const bounded = await store.events(operation.metadata.name); + expect(bounded.items).toHaveLength(MAX_OPERATION_EVENTS); + expect(bounded.items[0]?.sequence).toBe(3); + expect(bounded).toMatchObject({ + retainedFirstSequence: 3, + cursorGap: true, + }); + expect(await store.events(operation.metadata.name, 2)).toMatchObject({ + cursorGap: false, + }); + }); + test("provides exclusive, renewable per-workspace leases", async () => { let now = 0; const leases = new MemoryWorkspaceLeaseProvider(() => now); diff --git a/tests/server/registry.test.ts b/tests/server/registry.test.ts index b1143a3..6c32476 100644 --- a/tests/server/registry.test.ts +++ b/tests/server/registry.test.ts @@ -54,6 +54,69 @@ describe("OCI registry digest resolution", () => { ]); }); + test("caches successful digest resolutions with deterministic expiry and bounds", async () => { + let now = 0; + let calls = 0; + const fetcher = async () => { + calls += 1; + return new Response("manifest", { + headers: { "docker-content-digest": `sha256:${"d".repeat(64)}` }, + }); + }; + const options = { + fetch: fetcher, + cacheTtlMs: 100, + cacheMaxEntries: 1, + clock: () => now, + }; + + await resolveRegistryDigest("cache.example.com/team/first:v1", options); + await resolveRegistryDigest("cache.example.com/team/first:v1", options); + expect(calls).toBe(1); + + now = 100; + await resolveRegistryDigest("cache.example.com/team/first:v1", options); + expect(calls).toBe(2); + + await resolveRegistryDigest("cache.example.com/team/second:v1", options); + await resolveRegistryDigest("cache.example.com/team/first:v1", options); + expect(calls).toBe(4); + }); + + test("does not cache failed resolutions or share entries across credentials", async () => { + let calls = 0; + const failing = async () => { + calls += 1; + return new Response("unavailable", { status: 503 }); + }; + const options = { fetch: failing, cacheTtlMs: 1_000 }; + await expect( + resolveRegistryDigest("failure.example.com/team/image:v1", options), + ).rejects.toThrow("503"); + await expect( + resolveRegistryDigest("failure.example.com/team/image:v1", options), + ).rejects.toThrow("503"); + expect(calls).toBe(2); + + const authorizedCalls: string[] = []; + const authorized = async ( + _input: string | URL | Request, + init?: RequestInit, + ) => { + authorizedCalls.push(new Headers(init?.headers).get("authorization")!); + return new Response("manifest", { + headers: { "docker-content-digest": `sha256:${"e".repeat(64)}` }, + }); + }; + for (const username of ["one", "two"]) { + await resolveRegistryDigest("credentials.example.com/team/image:v1", { + fetch: authorized, + credentials: { username, password: "password" }, + }); + } + expect(authorizedCalls).toHaveLength(2); + }); + test("resolves through a trusted internal registry origin", async () => { const expected = `sha256:${"c".repeat(64)}` as const; let requested = "";