import { createApp, cleanupExpiredSessions } from "./app"; import { logServerRequest, processLogger, safeLog } from "../lib/request-log"; import { RedactingAuditStore } from "./audit-store"; import { readFile } from "node:fs/promises"; import { IMAGE_REGISTRY } from "../const"; import { BuildController, buildImageName } from "./build-controller"; import { KubernetesBuildOperations, KubernetesBuildStore, buildObjectApi, } from "./build-kubernetes"; import { FilesystemCas } from "./cas"; import { KubernetesAuthStore } from "./kubernetes-store"; import { MAINTENANCE_ROUTE_NAME, MAINTENANCE_STATE_NAME, MAINTENANCE_STATE_NAMESPACE, MaintenanceService, } from "./maintenance"; import { PatchStrategy, type KubernetesObject } from "@kubernetes/client-node"; import { createKubernetesClients, createKubernetesLeaseObjects, createKubernetesManagementDependencies, KubernetesAuditPersistence, KubernetesOperationPersistence, KubernetesWorkspaceLeaseProvider, KubernetesWorkspacePersistence, KubernetesWorkspaceStore, KubernetesTrustStore, } from "./kubernetes-state"; import { execUpgradeMatch, handleExecUpgrade, WireExecSession, type ExecConnection, } from "./app"; import { KubernetesExec } from "./kubernetes-exec"; import { createExecService } from "./exec-service"; import { createManagementService } from "./management"; import { PersistentOperationStore, recoverStaleOperations, } from "./operation-store"; import { KubernetesLogs } from "./kubernetes-logs"; import { createLogService } from "./log-service"; import { resolveRegistryDigest, type RegistryCredentials } from "./registry"; import { BackgroundBuildReconciler, buildReconcileIntervalMs, buildReconcileTimeoutMs, } from "./build-reconciler"; const clients = createKubernetesClients(); const store = new KubernetesAuthStore(clients.objects); const workspaceStore = new KubernetesWorkspaceStore( new KubernetesWorkspacePersistence(clients.objects), ); const operationStore = new PersistentOperationStore( new KubernetesOperationPersistence(clients.objects), ); const auditStore = new RedactingAuditStore( new KubernetesAuditPersistence(clients.objects), ); const trustStore = new KubernetesTrustStore(clients.objects); const management = createManagementService( createKubernetesManagementDependencies(clients), ); const namespace = process.env.KUBER_SYSTEM_NAMESPACE?.trim() || "kuber-system"; const dataRoot = process.env.KUBER_DATA_ROOT?.trim() || "/data"; const registry = ( process.env.KUBER_BUILD_REGISTRY?.trim() || IMAGE_REGISTRY ).replace(/\/+$/, ""); const registryConfigPath = process.env.KUBER_REGISTRY_CONFIG?.trim() || "/etc/kuber/registry/config.json"; const registryResolveOrigin = process.env.KUBER_REGISTRY_RESOLVE_ORIGIN?.trim(); const internalRegistryHost = process.env.KUBER_INTERNAL_REGISTRY_HOST?.trim(); const internalRegistryInsecure = process.env.KUBER_INTERNAL_REGISTRY_INSECURE?.trim() === "true"; const pushImagePrefix = process.env.KUBER_PUSH_IMAGE_PREFIX?.trim() || "kuber/"; if (!process.env.KUBER_REGISTRY_SECRET?.trim()) { console.warn( "KUBER_REGISTRY_SECRET is unset; BuildKit will use anonymous registry access", ); } async function registryCredentials(): Promise { let config: { auths?: Record }; try { config = JSON.parse(await readFile(registryConfigPath, "utf8")); } catch (error) { if ((error as NodeJS.ErrnoException).code === "ENOENT") return; throw error; } const host = registry.split("/")[0]!; const entry = config.auths?.[host] ?? config.auths?.[`https://${host}`] ?? config.auths?.[`https://${host}/v1/`]; if (!entry?.auth) return; const separator = Buffer.from(entry.auth, "base64") .toString("utf8") .indexOf(":"); if (separator < 0) return; const value = Buffer.from(entry.auth, "base64").toString("utf8"); return { username: value.slice(0, separator), password: value.slice(separator + 1), }; } const buildStore = new KubernetesBuildStore( buildObjectApi(clients.objects), namespace, `${dataRoot}/uploads`, ); const leases = new KubernetesWorkspaceLeaseProvider( createKubernetesLeaseObjects(clients.coordination), namespace, ); const maintenance = new MaintenanceService( { async readState() { try { const value = (await clients.objects.read({ apiVersion: "v1", kind: "ConfigMap", metadata: { name: MAINTENANCE_STATE_NAME, namespace: MAINTENANCE_STATE_NAMESPACE, }, })) as { data?: Record }; const hosts = JSON.parse(value.data?.hosts ?? "[]"); return Array.isArray(hosts) && hosts.every((host) => typeof host === "string") ? { hosts } : { hosts: [] }; } catch (error) { if ( error && typeof error === "object" && (("code" in error && error.code === 404) || ("statusCode" in error && error.statusCode === 404)) ) return; throw error; } }, async writeState({ hosts }) { await clients.objects.patch( { apiVersion: "v1", kind: "ConfigMap", metadata: { name: MAINTENANCE_STATE_NAME, namespace: MAINTENANCE_STATE_NAMESPACE, labels: { "kuber.astrxl.dev/type": "maintenance" }, }, data: { hosts: JSON.stringify(hosts) }, } as KubernetesObject, undefined, undefined, "kuber-server", true, PatchStrategy.ServerSideApply, ); }, async routeExists() { try { await clients.objects.read({ apiVersion: "traefik.io/v1alpha1", kind: "IngressRoute", metadata: { name: MAINTENANCE_ROUTE_NAME, namespace: "routing" }, }); return true; } catch (error) { if ( error && typeof error === "object" && (("code" in error && error.code === 404) || ("statusCode" in error && error.statusCode === 404)) ) return false; throw error; } }, async apply(resource) { await clients.objects.patch( resource as KubernetesObject, undefined, undefined, "kuber-server", true, PatchStrategy.ServerSideApply, ); }, async deleteRoute() { try { await clients.objects.delete({ apiVersion: "traefik.io/v1alpha1", kind: "IngressRoute", metadata: { name: MAINTENANCE_ROUTE_NAME, namespace: "routing" }, }); } catch (error) { if ( !( error && typeof error === "object" && (("code" in error && error.code === 404) || ("statusCode" in error && error.statusCode === 404)) ) ) throw error; } }, }, new KubernetesWorkspaceLeaseProvider( createKubernetesLeaseObjects(clients.coordination), MAINTENANCE_STATE_NAMESPACE, ), ); const builds = new BuildController({ cas: new FilesystemCas(`${dataRoot}/cas`), store: buildStore, kubernetes: new KubernetesBuildOperations(clients.batch, clients.core), namespace, workspaceRoot: `${dataRoot}/workspaces`, workspaceClaimName: process.env.KUBER_BUILD_DATA_CLAIM?.trim() || "kuber-build-data", cacheImage: (request) => `${internalRegistryHost ?? registry}/kuber/cache-${request.project}-${request.service}${request.spec.builder === "buildpacks" ? `-cnb-${request.spec.architecture}` : ""}`, imageName: (request) => buildImageName(registry, request.project, request.service), pushImage: internalRegistryHost ? (request) => `${internalRegistryHost}/${pushImagePrefix}${request.project}-${request.service}:latest` : undefined, pushRegistryInsecure: internalRegistryHost ? internalRegistryInsecure : undefined, buildkitImage: process.env.KUBER_BUILDKIT_IMAGE?.trim() || undefined, buildpacksImage: process.env.KUBER_BUILDPACKS_IMAGE?.trim() || undefined, registrySecretName: process.env.KUBER_REGISTRY_SECRET?.trim() || undefined, maxLogBytes: Number(process.env.KUBER_BUILD_LOG_BYTES ?? 512 * 1024), resolveDigest: async (image) => resolveRegistryDigest(image, { credentials: await registryCredentials(), origin: registryResolveOrigin, insecure: registryResolveOrigin?.startsWith("http://"), }), onReconcileFailure: (buildId, error) => safeLog(processLogger, { event: "build.reconcile.failed", buildId, error, }), reconcileLeases: leases, }); const logs = createLogService( new KubernetesLogs(clients.config, clients.apps, clients.core), ); const execService = createExecService(new KubernetesExec(clients.config)); const bootstrapUsername = process.env.KUBER_BOOTSTRAP_USERNAME?.trim(); const bootstrapPassword = process.env.KUBER_BOOTSTRAP_PASSWORD; if (bootstrapUsername && bootstrapPassword) { const existing = await store.getUser(bootstrapUsername); if (!existing) { await store.putUser({ username: bootstrapUsername, passwordHash: await Bun.password.hash(bootstrapPassword, { algorithm: "argon2id", }), roles: ["admin"], }); console.log(`Created bootstrap user ${bootstrapUsername}`); } } try { const recovered = await recoverStaleOperations(operationStore); if (recovered > 0) console.log(`Marked ${recovered} stale operation(s) as failed`); } catch (error) { console.error("Startup operation recovery failed", error); } try { await maintenance.reconcile(); } catch (error) { console.error("Startup maintenance reconciliation failed", error); } const sessionCleanupIntervalMs = Number( process.env.KUBER_SESSION_CLEANUP_MS ?? 10 * 60 * 1000, ); const sessionCleanupTimer = setInterval( async () => { try { await cleanupExpiredSessions(store); } catch (error) { console.error("Expired authentication cleanup failed", error); } }, Number.isFinite(sessionCleanupIntervalMs) && sessionCleanupIntervalMs > 0 ? sessionCleanupIntervalMs : 10 * 60 * 1000, ); sessionCleanupTimer.unref?.(); const buildReconciler = new BackgroundBuildReconciler({ runner: builds, timeoutMs: buildReconcileTimeoutMs( process.env.KUBER_BUILD_RECONCILE_TIMEOUT_MS, ), onFailure: (error) => safeLog(processLogger, { event: "build.reconcile.scan.failed", error }), }); const buildReconcileTimer = setInterval( () => buildReconciler.tick(), buildReconcileIntervalMs(process.env.KUBER_BUILD_RECONCILE_MS), ); buildReconcileTimer.unref?.(); const app = createApp({ store, workspaceStore, operationStore, auditStore, trustStore, management, builds, logs, execService, maintenance, leases, resolveImage: async (project, service) => { const image = buildImageName(registry, project, service); const digest = await resolveRegistryDigest(image, { credentials: await registryCredentials(), origin: registryResolveOrigin, insecure: registryResolveOrigin?.startsWith("http://"), }); return { image: image.replace(/:latest$/, ""), digest, reference: `${image.replace(/:latest$/, "")}@${digest}`, }; }, allowedOrigins: ( process.env.KUBER_ALLOWED_ORIGINS ?? "https://kuber.astrxl.dev" ) .split(",") .map((origin) => origin.trim()) .filter(Boolean), }); type ExecConnectionData = { connection: ExecConnection; }; const server = Bun.serve({ port: Number(process.env.PORT ?? 3000), fetch(request, server) { const url = new URL(request.url); const workspaceId = execUpgradeMatch(url); if (workspaceId && request.method === "GET") { const requestId = request.headers.get("x-request-id") ?? crypto.randomUUID(); return logServerRequest(request, requestId, () => handleExecUpgrade( { store, workspaceStore, auditStore, }, request, workspaceId, (upgradeRequest, connection) => server.upgrade(upgradeRequest, { data: { connection } }), ), ); } return app(request); }, websocket: { open(ws) { const { connection } = ws.data; const session = new WireExecSession( { sendText: (data) => ws.sendText(data), close: (code, reason) => ws.close(code, reason), }, execService, connection, ); ws.data = { connection, session } as unknown as ExecConnectionData; }, message(ws, message) { const state = ws.data as unknown as { session?: WireExecSession; }; state.session?.receive( typeof message === "string" ? message : new TextDecoder().decode(message), ); }, close(ws) { const state = ws.data as unknown as { session?: WireExecSession; }; state.session?.close(); }, }, }); console.log(`kuber server listening on ${server.url}`);