import { createHash, randomUUID } from "node:crypto"; import { KUBER_API_VERSION, type ObjectMeta } from "./workspace-store"; 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"; export interface Operation { apiVersion: typeof KUBER_API_VERSION; kind: "Operation"; metadata: ObjectMeta & { workspaceUid?: string }; spec: { workspaceId: string; action: string; idempotencyKey: string; requestHash: string; }; status: { state: OperationState; startedAt?: string; finishedAt?: string; result?: unknown; error?: OperationError; events?: OperationEvent[]; }; } export interface CreateOperationInput { workspaceId: string; workspaceUid?: string; action: string; idempotencyKey: string; request?: unknown; } export class OperationValidationError extends Error { readonly code = "OPERATION_INVALID"; } export class OperationConflictError extends Error { readonly code = "OPERATION_CONFLICT"; } export class OperationNotFoundError extends Error { readonly code = "OPERATION_NOT_FOUND"; } /** createIdempotent must atomically index workspaceId + idempotencyKey. */ export interface OperationPersistence { createIdempotent(operation: Operation): Promise; get(id: string): Promise; list(workspaceId?: string): Promise; replace(operation: Operation, expectedResourceVersion: string): Promise; } /** The idempotent operation and whether this caller created its record. */ export interface CreatedOperation { operation: Operation; created: boolean; } export interface OperationStore { create(input: CreateOperationInput): Promise; createOrReuse(input: CreateOperationInput): Promise; /** * Atomically claims a pending operation for execution. A false result means * another request has already moved it out of pending. */ claimExecution(id: string): Promise; get(id: string): Promise; list(workspaceId?: string): Promise; transition( id: string, 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 { workspaceId: string; holder: string; expiresAt: string; renew(ttlMs?: number): Promise; release(): Promise; } export interface WorkspaceLeaseProvider { acquire( workspaceId: string, holder: string, ttlMs?: number, ): Promise; } /** * Marks any operation left in a non-terminal state by a previous process as * failed. On process restart no synchronous operation can legitimately still be * in flight (startup runs before serving), so pending/running records are stale. * Best-effort: individual conflicts or failures are logged and skipped so one * bad record cannot block recovery. */ export async function recoverStaleOperations( store: Pick, ): Promise { const operations = await store.list(); let recovered = 0; for (const operation of operations) { if ( operation.status.state !== "pending" && operation.status.state !== "running" ) continue; const name = operation.metadata.name; try { await store.transition(name, "failed", { error: { code: "OPERATION_INTERRUPTED", message: "Server restarted while the operation was in progress; retry to proceed", }, }); recovered += 1; } catch (error) { console.error( `Failed to mark stale operation '${name}' as failed during startup recovery`, error, ); } } return recovered; } const TRANSITIONS: Record> = { pending: new Set(["running", "failed", "cancelled"]), running: new Set(["succeeded", "failed", "cancelled"]), succeeded: new Set(), failed: new Set(), cancelled: new Set(), }; function clone(value: T): T { return structuredClone(value); } function canonicalJson(value: unknown): string { if (value === null || typeof value !== "object") { const encoded = JSON.stringify(value); if (encoded === undefined) throw new OperationValidationError( "Operation request must be JSON serializable", ); return encoded; } if (Array.isArray(value)) return `[${value.map(canonicalJson).join(",")}]`; return `{${Object.keys(value as Record) .sort() .map( (key) => `${JSON.stringify(key)}:${canonicalJson((value as Record)[key])}`, ) .join(",")}}`; } function hashRequest(input: CreateOperationInput): string { let request: string; try { request = canonicalJson(input.request ?? null); } catch { throw new OperationValidationError( "Operation request must be JSON serializable", ); } return createHash("sha256") .update(`${input.action}\0${request}`) .digest("hex"); } const SENSITIVE_KEY = /(?:password|passwd|token|secret|credential|private.?key|access.?key|connection.?string|database.?url)/i; /** Produces a persistence-safe result while retaining useful operation status. */ export function sanitizeOperationResult( value: unknown, action?: string, ): unknown { if (action === "databases.reconcile" || action === "storage.reconcile") { return { redacted: true }; } if (typeof value === "string") return redactString(value); if (Array.isArray(value)) return value.map((item) => sanitizeOperationResult(item)); if (!value || typeof value !== "object") return value; const record = value as Record; const isDataResource = record.kind === "Secret" || record.kind === "ConfigMap"; return Object.fromEntries( Object.entries(record).map(([key, item]) => [ key, SENSITIVE_KEY.test(key) || (isDataResource && (key === "data" || key === "stringData")) ? REDACTED : sanitizeOperationResult(item), ]), ); } /** Keep operation error codes stable while preventing provider details from escaping. */ export function sanitizeOperationError( error: OperationError, action?: string, ): OperationError { const message = action === "databases.reconcile" ? "Database reconciliation failed" : action === "storage.reconcile" ? "Storage reconciliation failed" : redactString(error.message); return { code: error.code, message }; } export class PersistentOperationStore implements OperationStore { constructor( private readonly persistence: OperationPersistence, private readonly now: () => Date = () => new Date(), private readonly uid: () => string = randomUUID, ) {} async create(input: CreateOperationInput): Promise { return (await this.createOrReuse(input)).operation; } async createOrReuse(input: CreateOperationInput): Promise { if (!input.workspaceId || !input.action.trim()) { throw new OperationValidationError( "Workspace ID and action are required", ); } const keyBytes = Buffer.byteLength(input.idempotencyKey); if (!input.idempotencyKey || keyBytes > MAX_IDEMPOTENCY_KEY_BYTES) { throw new OperationValidationError( `Idempotency key must be 1-${MAX_IDEMPOTENCY_KEY_BYTES} bytes`, ); } const now = this.now().toISOString(); const id = this.uid(); const operation: Operation = { apiVersion: KUBER_API_VERSION, kind: "Operation", metadata: { name: `operation-${id}`, uid: id, resourceVersion: "1", creationTimestamp: now, ...(input.workspaceUid && { workspaceUid: input.workspaceUid }), }, spec: { workspaceId: input.workspaceId, action: input.action, idempotencyKey: input.idempotencyKey, requestHash: hashRequest(input), }, status: { state: "pending" }, }; const stored = await this.persistence.createIdempotent(operation); if (stored.operation.spec.requestHash !== operation.spec.requestHash) { throw new OperationConflictError( "Idempotency key was already used for a different request", ); } return { operation: clone(stored.operation), created: stored.created }; } async get(id: string) { const operation = await this.persistence.get(id); return operation && clone(operation); } async claimExecution(id: string): Promise { try { return await this.transition(id, "running"); } catch (error) { if (!(error instanceof OperationConflictError)) throw error; return; } } async list(workspaceId?: string) { return clone(await this.persistence.list(workspaceId)); } async transition( id: string, state: OperationState, options: { result?: unknown; error?: OperationError } = {}, ): Promise { const current = await this.persistence.get(id); if (!current) throw new OperationNotFoundError(`Operation '${id}' not found`); if (!TRANSITIONS[current.status.state].has(state)) { throw new OperationConflictError( `Cannot transition operation from ${current.status.state} to ${state}`, ); } if (state === "succeeded" && options.error) { throw new OperationValidationError( "Successful operations cannot have errors", ); } if (state === "failed" && !options.error) { throw new OperationValidationError("Failed operations require an error"); } const timestamp = this.now().toISOString(); const operation: Operation = clone(current); operation.metadata.resourceVersion = String( Number(current.metadata.resourceVersion) + 1, ); operation.status = { state, ...(current.status.startedAt && { startedAt: current.status.startedAt }), ...(state === "running" && { startedAt: timestamp }), ...(["succeeded", "failed", "cancelled"].includes(state) && { finishedAt: timestamp, }), ...(options.result !== undefined && { result: clone( sanitizeOperationResult(options.result, current.spec.action), ), }), ...(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 { private readonly operations = new Map(); private readonly idempotency = new Map(); async createIdempotent(operation: Operation) { const key = `${operation.spec.workspaceId}\0${operation.spec.idempotencyKey}`; const existingId = this.idempotency.get(key); if (existingId) return { operation: clone(this.operations.get(existingId)!), created: false, }; this.operations.set(operation.metadata.name, clone(operation)); this.idempotency.set(key, operation.metadata.name); return { operation: clone(operation), created: true }; } async get(id: string) { const operation = this.operations.get(id); return operation && clone(operation); } async list(workspaceId?: string) { return [...this.operations.values()] .filter((operation) => workspaceId ? operation.spec.workspaceId === workspaceId : true, ) .sort((a, b) => a.metadata.creationTimestamp.localeCompare( b.metadata.creationTimestamp, ), ) .map(clone); } async replace(operation: Operation, expectedResourceVersion: string) { const current = this.operations.get(operation.metadata.name); if (!current) throw new OperationNotFoundError("Operation not found"); if (current.metadata.resourceVersion !== expectedResourceVersion) { throw new OperationConflictError("Operation was concurrently modified"); } this.operations.set(operation.metadata.name, clone(operation)); } } export class MemoryOperationStore extends PersistentOperationStore { constructor(now?: () => Date, uid?: () => string) { super(new MemoryOperationPersistence(), now, uid); } } export class MemoryWorkspaceLeaseProvider implements WorkspaceLeaseProvider { private readonly leases = new Map< string, { holder: string; expiresAt: number; token: string } >(); constructor(private readonly now: () => number = Date.now) {} async acquire(workspaceId: string, holder: string, ttlMs = 30_000) { if (!workspaceId || !holder || !Number.isFinite(ttlMs) || ttlMs <= 0) { throw new OperationValidationError("Invalid workspace lease request"); } const current = this.leases.get(workspaceId); if (current && current.expiresAt > this.now()) { return undefined; } const token = randomUUID(); this.leases.set(workspaceId, { holder, token, expiresAt: this.now() + ttlMs, }); const lease: WorkspaceLease = { workspaceId, holder, expiresAt: new Date(this.now() + ttlMs).toISOString(), renew: async (nextTtl = ttlMs) => { const active = this.leases.get(workspaceId); if ( !active || active.token !== token || active.expiresAt <= this.now() || nextTtl <= 0 ) { return false; } active.expiresAt = this.now() + nextTtl; lease.expiresAt = new Date(active.expiresAt).toISOString(); return true; }, release: async () => { if (this.leases.get(workspaceId)?.token === token) { this.leases.delete(workspaceId); } }, }; return lease; } }