diff --git a/README.md b/README.md index 2d326cd..6eef6a4 100644 --- a/README.md +++ b/README.md @@ -108,7 +108,7 @@ For a permanent setup, write the generated script to a file and source it from y - `logs -f [deployment]`: follow logs continuously - `exec `: execute a command inside a running deployment pod - `db ls`: list managed Postgres claims declared in the current Compose file -- `db creds `: print the generated connection details for a managed Postgres claim +- `db creds `: reconcile and print the generated connection details for a managed Postgres claim - `s3 ls`: list managed S3 claims declared in the current Compose file - `s3 creds `: print all generated S3 environment variables for a service - `s3 ui `: print the Garage UI object-browser URL for a service bucket @@ -236,7 +236,12 @@ services: - postgresql:user/database ``` -This creates or reuses the managed CNPG role secret, reconciles the database resource, and injects `DATABASE_URL` into the generated app secret in Kubernetes. +This creates or reuses the managed CNPG role secret, reconciles the database +resource, and injects `DATABASE_URL` and +`REDIS_URL=redis://redis.database.svc.cluster.local` into the generated app +secret in Kubernetes. `REDIS_URL` is only added to services with a managed +Postgres claim. Running `kuber db creds ` performs the same focused +Secret, role, and Database reconciliation before printing credentials. ### Managed S3 @@ -297,7 +302,11 @@ This is also where generated values such as `DATABASE_URL` and the managed S3 en - `postgresql:...` is treated as a managed database claim, not as a filesystem mount - `s3:...` is treated as a managed object-storage claim, not as a filesystem mount -Named volumes also support kuber-specific storage hints. +Named volumes also support kuber-specific Longhorn storage hints. Kuber renders +each distinct placement policy as a deterministic, reusable Longhorn +`StorageClass`, then references that class from the PVC. The generated class +uses Longhorn's `numberOfReplicas`, `diskSelector`, and `dataLocality` +parameters; placement fields are never written directly to the PVC. Default behavior: @@ -329,8 +338,8 @@ services: Meaning: - `name(20Gi):/path` -> PVC size `20Gi` -- `name(20Gi on archive):/path` -> PVC size `20Gi`, `diskTag: ["archive"]`, `dataLocality: "none"` -- `name(20Gi on 1 archive):/path` -> PVC size `20Gi`, `replicaCount: 1`, `diskTag: ["archive"]`, `dataLocality: "none"` +- `name(20Gi on archive):/path` -> PVC size `20Gi`, with a StorageClass using `diskSelector: "archive"` and disabled data locality +- `name(20Gi on 1 archive):/path` -> PVC size `20Gi`, with a StorageClass using one replica and `diskSelector: "archive"` Explicit extensions are also supported. @@ -368,6 +377,9 @@ Precedence: - top-level `volumes..x-*` - fallback default `1Gi on 2 fast` +StorageClasses are cluster-scoped and content-addressed by policy. They are +shared across projects and intentionally retained when a project is removed. + ## Building Image builds run through the selected SSH builder. After a push, the registry's diff --git a/command/db.ts b/command/db.ts index fd59bd4..45e1860 100644 --- a/command/db.ts +++ b/command/db.ts @@ -5,7 +5,7 @@ import { DATABASE_PORT, buildDatabaseUrl, getComposePostgresClaims, - getRoleCredentials, + reconcilePostgresClaim, } from "../lib/database"; import { toTable } from "../lib/format"; @@ -43,7 +43,8 @@ const creds = defineCommand({ const service = args._[0]; if (!service) throw new Error("Service name is required"); - const claim = getComposePostgresClaims(await ctx().compose()).find( + const context = ctx(); + const claim = getComposePostgresClaims(await context.compose()).find( (entry) => entry.service === service, ); if (!claim) { @@ -52,7 +53,7 @@ const creds = defineCommand({ ); } - const credentials = await getRoleCredentials(claim.username); + const credentials = await reconcilePostgresClaim(context.project, claim); const url = buildDatabaseUrl(claim, credentials); console.log( diff --git a/command/logs.ts b/command/logs.ts index 505e9ce..3c2742c 100644 --- a/command/logs.ts +++ b/command/logs.ts @@ -2,6 +2,7 @@ import { defineCommand } from "citty"; import https from "node:https"; import { createLogger } from "../lib/logger"; import { core, kc } from "../lib/k8s"; +import { delay, parseRetryAfter } from "../lib/k8s-http"; import { getPodContainerName, listManagedDeployments, @@ -13,10 +14,6 @@ function formatError(error: unknown): string { return String(error); } -function delay(ms: number) { - return new Promise((resolve) => setTimeout(resolve, ms)); -} - async function streamResponseLines( stream: NodeJS.ReadableStream, onLine: (line: string) => void, @@ -49,7 +46,9 @@ async function followPodLogs( const cluster = kc.getCurrentCluster(); if (!cluster) throw new Error("No currently active cluster"); - const requestURL = new URL(`${cluster.server}/api/v1/namespaces/${namespace}/pods/${podName}/log`); + const requestURL = new URL( + `${cluster.server}/api/v1/namespaces/${namespace}/pods/${podName}/log`, + ); requestURL.searchParams.set("container", containerName); requestURL.searchParams.set("follow", "true"); @@ -63,30 +62,49 @@ async function followPodLogs( }; await kc.applyToHTTPSOptions(options); - await new Promise((resolve, reject) => { - const request = https.request(options, async (response) => { - if ((response.statusCode ?? 0) < 200 || (response.statusCode ?? 0) > 299) { - const chunks: Buffer[] = []; - response.on("data", (chunk) => { - chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk)); - }); - response.on("end", () => { - reject(Buffer.concat(chunks).toString("utf8") || `HTTP ${response.statusCode}`); - }); - return; - } + for (let attempt = 0; ; attempt += 1) { + const result = await new Promise<{ + statusCode: number; + retryAfter?: string; + body?: string; + }>((resolve, reject) => { + const request = https.request(options, async (response) => { + const statusCode = response.statusCode ?? 0; + if (statusCode < 200 || statusCode > 299) { + const chunks: Buffer[] = []; + response.on("data", (chunk) => { + chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk)); + }); + response.on("end", () => { + resolve({ + statusCode, + retryAfter: response.headers["retry-after"], + body: Buffer.concat(chunks).toString("utf8"), + }); + }); + return; + } - try { - await streamResponseLines(response, writeLine); - resolve(); - } catch (error) { - reject(error); - } + try { + await streamResponseLines(response, writeLine); + resolve({ statusCode }); + } catch (error) { + reject(error); + } + }); + + request.on("error", reject); + request.end(); }); - request.on("error", reject); - request.end(); - }); + if (result.statusCode >= 200 && result.statusCode <= 299) return; + if (result.statusCode !== 429 || attempt >= 3) { + throw new Error(result.body || `HTTP ${result.statusCode}`); + } + + const retryAfter = parseRetryAfter(result.retryAfter); + await delay(retryAfter ?? 250 * 2 ** attempt, signal); + } } async function followDeploymentLogs(name: string, signal: AbortSignal) { @@ -98,7 +116,11 @@ async function followDeploymentLogs(name: string, signal: AbortSignal) { for (const pod of pods) { const podName = pod.metadata?.name; - if (!podName || activePods.has(podName) || pod.status?.phase !== "Running") { + if ( + !podName || + activePods.has(podName) || + pod.status?.phase !== "Running" + ) { continue; } @@ -118,7 +140,7 @@ async function followDeploymentLogs(name: string, signal: AbortSignal) { }); } - await delay(2000); + await delay(2000, signal); } } @@ -127,7 +149,8 @@ async function logDeployment(name: string, follow: boolean) { if (!follow) { const pods = await listPodsForDeployment(name); - if (pods.length === 0) throw new Error(`No pods found for deployment ${name}`); + if (pods.length === 0) + throw new Error(`No pods found for deployment ${name}`); for (const pod of pods) { const podName = pod.metadata?.name; diff --git a/lib/apply.ts b/lib/apply.ts index c9e8b30..91cfe4a 100644 --- a/lib/apply.ts +++ b/lib/apply.ts @@ -9,13 +9,14 @@ const ResourceOrder = { Namespace: 0, GarageBucket: 1, GarageKey: 2, - PersistentVolumeClaim: 3, - Secret: 4, - ConfigMap: 5, - Service: 6, - Deployment: 7, - Ingress: 8, - IngressRoute: 9, + StorageClass: 3, + PersistentVolumeClaim: 4, + Secret: 5, + ConfigMap: 6, + Service: 7, + Deployment: 8, + Ingress: 9, + IngressRoute: 10, } as const; const ManagedResources = [ diff --git a/lib/convert.ts b/lib/convert.ts index 3e20fee..531c752 100644 --- a/lib/convert.ts +++ b/lib/convert.ts @@ -55,6 +55,20 @@ type NormalizedPort = { type ServiceVolume = NonNullable[number]; type ComposeVolumes = ComposeSpecification["volumes"]; +type LonghornStoragePolicy = { + diskTag?: string[]; + replicaCount?: number; + dataLocality?: "none"; +}; + +type LonghornStorageClass = KubernetesObject & { + provisioner: "driver.longhorn.io"; + allowVolumeExpansion: boolean; + reclaimPolicy: "Delete"; + volumeBindingMode: "Immediate"; + parameters: Record; +}; + const DefaultNamedVolumeStorage = { requestedStorage: "1Gi", diskTag: ["fast"], @@ -62,6 +76,32 @@ const DefaultNamedVolumeStorage = { dataLocality: "none" as const, }; +function getLonghornStoragePolicy( + mount: NormalizedMount, +): LonghornStoragePolicy | undefined { + if ( + !mount.diskTag && + mount.replicaCount === undefined && + !mount.dataLocality + ) { + return; + } + + return { + diskTag: mount.diskTag ? [...mount.diskTag].sort() : undefined, + replicaCount: mount.replicaCount, + dataLocality: mount.dataLocality, + }; +} + +function toLonghornStorageClassName(policy: LonghornStoragePolicy): string { + const digest = createHash("sha256") + .update(JSON.stringify(policy)) + .digest("hex") + .slice(0, 12); + return `kuber-longhorn-${digest}`; +} + function toPortNumber(value: number | string | undefined): number | undefined { if (typeof value === "number") return value; if (!value || value.includes("-")) return; @@ -149,7 +189,9 @@ function normalizeDiskTags(value: unknown): string[] | undefined { function normalizeReplicaCount(value: unknown): number | undefined { if (value === undefined || value === null || value === "") return; const count = Number(value); - return Number.isInteger(count) ? count : undefined; + return Number.isInteger(count) && count >= 1 && count <= 20 + ? count + : undefined; } function parseStorageSpec(value: string): { @@ -980,6 +1022,7 @@ export function volumesToPvc( if (!claimName || claims.has(claimName)) continue; + const policy = getLonghornStoragePolicy(mount); claims.set(claimName, { apiVersion: "v1", kind: "PersistentVolumeClaim", @@ -990,23 +1033,57 @@ export function volumesToPvc( }, spec: { accessModes: ["ReadWriteMany"], + storageClassName: policy + ? toLonghornStorageClassName(policy) + : undefined, resources: { requests: { storage: mount.requestedStorage ?? "1Gi", }, }, - ...(mount.diskTag ? { diskTag: mount.diskTag } : {}), - ...(mount.replicaCount !== undefined - ? { replicaCount: mount.replicaCount } - : {}), - ...(mount.dataLocality ? { dataLocality: mount.dataLocality } : {}), }, - } as V1PersistentVolumeClaim); + }); } return [...claims.values()]; } +export function volumesToStorageClasses( + service: Service, + cwd = process.cwd(), + volumes: ComposeVolumes = {}, +): LonghornStorageClass[] { + const classes = new Map(); + + for (const mount of toMounts(service, cwd, volumes)) { + const policy = getLonghornStoragePolicy(mount); + if (!policy) continue; + + const name = toLonghornStorageClassName(policy); + classes.set(name, { + apiVersion: "storage.k8s.io/v1", + kind: "StorageClass", + metadata: { + name, + labels: LABELS, + }, + provisioner: "driver.longhorn.io", + allowVolumeExpansion: true, + reclaimPolicy: "Delete", + volumeBindingMode: "Immediate", + parameters: { + ...(policy.diskTag ? { diskSelector: policy.diskTag.join(",") } : {}), + ...(policy.replicaCount !== undefined + ? { numberOfReplicas: String(policy.replicaCount) } + : {}), + ...(policy.dataLocality ? { dataLocality: "disabled" } : {}), + }, + }); + } + + return [...classes.values()]; +} + export function volumesToConfigMaps( project: string, service: Service, @@ -1110,6 +1187,14 @@ export async function composeToKubernetes( resources.set(getResourceKey(namespace), namespace); for (const [name, service] of Object.entries(compose.services ?? {})) { + for (const storageClass of volumesToStorageClasses( + service, + cwd, + compose.volumes, + )) { + resources.set(getResourceKey(storageClass), storageClass); + } + for (const pvc of volumesToPvc(project, service, cwd, compose.volumes)) { resources.set(getResourceKey(pvc), pvc); } diff --git a/lib/database.ts b/lib/database.ts index 0f78aae..85972fc 100644 --- a/lib/database.ts +++ b/lib/database.ts @@ -9,6 +9,7 @@ export const DATABASE_NAMESPACE = "database"; export const DATABASE_CLUSTER = "postgres"; export const DATABASE_HOST = `c.${DATABASE_NAMESPACE}.svc.cluster.local`; export const DATABASE_PORT = 5432; +export const REDIS_URL = "redis://redis.database.svc.cluster.local"; export const DATABASE_PROJECT_LABEL = "kuber.dev/project"; export const DATABASE_SERVICE_LABEL = "kuber.dev/service"; @@ -35,7 +36,7 @@ export type PostgresClaim = { secretName: string; }; -type RoleCredentials = { +export type RoleCredentials = { username: string; password: string; }; @@ -59,6 +60,16 @@ export function buildDatabaseUrl( return `postgresql://${encodeConnectionComponent(credentials.username)}:${encodeConnectionComponent(credentials.password)}@${DATABASE_HOST}:${DATABASE_PORT}/${encodeConnectionComponent(claim.database)}`; } +export function buildPostgresEnvironment( + claim: PostgresClaim, + credentials: RoleCredentials, +): Record { + return { + DATABASE_URL: buildDatabaseUrl(claim, credentials), + REDIS_URL, + }; +} + function parsePostgresVolumeString( entry: string, ): Omit | undefined { @@ -314,6 +325,16 @@ async function reconcileDatabases( } } +export async function reconcilePostgresClaim( + project: string, + claim: PostgresClaim, +): Promise { + const credentials = await ensureRoleSecret(claim); + await reconcileManagedRoles([claim]); + await reconcileDatabases(project, [claim]); + return credentials; +} + export async function reconcilePostgresClaims( project: string, compose: ComposeSpecification, @@ -337,12 +358,7 @@ export async function reconcilePostgresClaims( throw new Error(`Missing credentials for ${claim.secretName}`); } - return [ - claim.service, - { - DATABASE_URL: buildDatabaseUrl(claim, credentials), - }, - ]; + return [claim.service, buildPostgresEnvironment(claim, credentials)]; }), ); } diff --git a/lib/k8s-http.ts b/lib/k8s-http.ts new file mode 100644 index 0000000..ab2962c --- /dev/null +++ b/lib/k8s-http.ts @@ -0,0 +1,230 @@ +import type { + HttpLibrary, + RequestContext, +} from "@kubernetes/client-node/dist/gen/http/http.js"; +import { ResponseContext } from "@kubernetes/client-node/dist/gen/http/http.js"; +import { from } from "@kubernetes/client-node/dist/gen/rxjsStub.js"; +import http from "node:http"; +import https from "node:https"; + +type TransportOptions = { + maxConcurrent: number; + minIntervalMs: number; + maxRetries: number; + baseRetryMs: number; + maxRetryMs: number; + random: () => number; +}; + +const DefaultTransportOptions: TransportOptions = { + maxConcurrent: 4, + minIntervalMs: 100, + maxRetries: 3, + baseRetryMs: 250, + maxRetryMs: 5000, + random: Math.random, +}; + +function abortError(): Error { + const error = new Error("Request aborted"); + error.name = "AbortError"; + return error; +} + +export function delay(ms: number, signal?: AbortSignal): Promise { + if (signal?.aborted) return Promise.reject(abortError()); + + return new Promise((resolve, reject) => { + const timer = setTimeout(() => { + signal?.removeEventListener("abort", onAbort); + resolve(); + }, ms); + const onAbort = () => { + clearTimeout(timer); + reject(abortError()); + }; + signal?.addEventListener("abort", onAbort, { once: true }); + }); +} + +export function parseRetryAfter( + value: string | undefined, + now = Date.now(), +): number | undefined { + if (!value) return; + + const seconds = Number(value); + if (Number.isFinite(seconds) && seconds >= 0) return seconds * 1000; + + const date = Date.parse(value); + if (Number.isNaN(date)) return; + return Math.max(0, date - now); +} + +class RequestLimiter { + private active = 0; + private nextRequestAt = 0; + private readonly queue: Array<{ + signal?: AbortSignal; + onAbort?: () => void; + resolve: (release: () => void) => void; + reject: (error: Error) => void; + }> = []; + + constructor( + private readonly maxConcurrent: number, + private readonly minIntervalMs: number, + ) {} + + acquire(signal?: AbortSignal): Promise<() => void> { + if (signal?.aborted) return Promise.reject(abortError()); + + return new Promise((resolve, reject) => { + const entry: (typeof this.queue)[number] = { signal, resolve, reject }; + entry.onAbort = () => { + const index = this.queue.indexOf(entry); + if (index === -1) return; + this.queue.splice(index, 1); + reject(abortError()); + }; + signal?.addEventListener("abort", entry.onAbort, { once: true }); + this.queue.push(entry); + this.drain(); + }); + } + + private drain(): void { + while (this.active < this.maxConcurrent && this.queue.length > 0) { + const entry = this.queue.shift()!; + entry.signal?.removeEventListener("abort", entry.onAbort!); + if (entry.signal?.aborted) { + entry.reject(abortError()); + continue; + } + + this.active += 1; + const now = Date.now(); + const waitMs = Math.max(0, this.nextRequestAt - now); + this.nextRequestAt = + Math.max(now, this.nextRequestAt) + this.minIntervalMs; + + void delay(waitMs, entry.signal) + .then(() => { + let released = false; + entry.resolve(() => { + if (released) return; + released = true; + this.active -= 1; + this.drain(); + }); + }) + .catch((error) => { + this.active -= 1; + entry.reject( + error instanceof Error ? error : new Error(String(error)), + ); + this.drain(); + }); + } + } +} + +function sendOnce(request: RequestContext): Promise { + return new Promise((resolve, reject) => { + const url = new URL(request.getUrl()); + const transport = url.protocol === "http:" ? http : https; + const signal = request.getSignal(); + const req = transport.request( + url, + { + method: request.getHttpMethod().toString(), + headers: request.getHeaders(), + agent: request.getAgent() as never, + }, + (response) => { + const chunks: Buffer[] = []; + + response.on("data", (chunk) => { + chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk)); + }); + + response.on("end", () => { + signal?.removeEventListener("abort", onAbort); + const buffer = Buffer.concat(chunks); + const headers: Record = {}; + + for (const [key, value] of Object.entries(response.headers)) { + if (Array.isArray(value)) headers[key] = value.join(", "); + else if (value !== undefined) headers[key] = String(value); + } + + resolve( + new ResponseContext(response.statusCode ?? 0, headers, { + text: async () => buffer.toString("utf8"), + binary: async () => buffer, + }), + ); + }); + }, + ); + + const onAbort = () => req.destroy(abortError()); + if (signal?.aborted) onAbort(); + else signal?.addEventListener("abort", onAbort, { once: true }); + + req.on("error", (error) => { + signal?.removeEventListener("abort", onAbort); + reject(error); + }); + + const body = request.getBody(); + if (body !== undefined && body !== null) { + req.write(body as string | Uint8Array); + } + req.end(); + }); +} + +export function createKubernetesHttpLibrary( + overrides: Partial = {}, +): HttpLibrary { + const options = { ...DefaultTransportOptions, ...overrides }; + const limiter = new RequestLimiter( + options.maxConcurrent, + options.minIntervalMs, + ); + + return { + send(request) { + const result = (async () => { + for (let attempt = 0; ; attempt += 1) { + const release = await limiter.acquire(request.getSignal()); + let response: ResponseContext; + try { + response = await sendOnce(request); + } finally { + release(); + } + + if ( + response.httpStatusCode !== 429 || + attempt >= options.maxRetries + ) { + return response; + } + + const retryAfter = parseRetryAfter(response.headers["retry-after"]); + const exponential = Math.min( + options.maxRetryMs, + options.baseRetryMs * 2 ** attempt, + ); + const backoff = + retryAfter ?? exponential * (0.8 + options.random() * 0.4); + await delay(backoff, request.getSignal()); + } + })(); + + return from(result); + }, + }; +} diff --git a/lib/k8s.ts b/lib/k8s.ts index eb417c9..50b67b0 100644 --- a/lib/k8s.ts +++ b/lib/k8s.ts @@ -6,73 +6,12 @@ import { KubernetesObjectApi, Log, } from "@kubernetes/client-node"; -import { ResponseContext } from "@kubernetes/client-node/dist/gen/http/http.js"; -import { from } from "@kubernetes/client-node/dist/gen/rxjsStub.js"; -import http from "node:http"; -import https from "node:https"; +import { createKubernetesHttpLibrary } from "./k8s-http"; const kc = new KubeConfig(); kc.loadFromDefault(); -const bunHttpLibrary = { - send(request: { - getUrl(): string; - getHttpMethod(): { toString(): string }; - getBody(): unknown; - getHeaders(): Record; - getSignal(): AbortSignal | undefined; - getAgent(): http.Agent | https.Agent | undefined; - }) { - const result = new Promise((resolve, reject) => { - const url = new URL(request.getUrl()); - const transport = url.protocol === "http:" ? http : https; - const req = transport.request( - url, - { - method: request.getHttpMethod().toString(), - headers: request.getHeaders(), - agent: request.getAgent(), - }, - (response) => { - const chunks: Buffer[] = []; - - response.on("data", (chunk) => { - chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk)); - }); - - response.on("end", () => { - const buffer = Buffer.concat(chunks); - const headers: Record = {}; - - for (const [key, value] of Object.entries(response.headers)) { - if (Array.isArray(value)) headers[key] = value.join(", "); - else if (value !== undefined) headers[key] = String(value); - } - - resolve( - new ResponseContext(response.statusCode ?? 0, headers, { - text: async () => buffer.toString("utf8"), - binary: async () => buffer, - }), - ); - }); - }, - ); - - request.getSignal()?.addEventListener("abort", () => { - req.destroy(new Error("Request aborted")); - }); - - req.on("error", reject); - - const body = request.getBody(); - if (body !== undefined && body !== null) req.write(body as string | Uint8Array); - req.end(); - }); - - return from(result); - }, -}; +const bunHttpLibrary = createKubernetesHttpLibrary(); const makeApiClient = kc.makeApiClient.bind(kc); kc.makeApiClient = ((apiClientType) => { @@ -80,7 +19,8 @@ kc.makeApiClient = ((apiClientType) => { api?: { configuration?: { httpApi?: typeof bunHttpLibrary } }; configuration?: { httpApi?: typeof bunHttpLibrary }; }; - if (client.api?.configuration) client.api.configuration.httpApi = bunHttpLibrary; + if (client.api?.configuration) + client.api.configuration.httpApi = bunHttpLibrary; if (client.configuration) client.configuration.httpApi = bunHttpLibrary; return client; }) as typeof kc.makeApiClient; diff --git a/tests/lib/apply.test.ts b/tests/lib/apply.test.ts index 8ed6035..35c91f3 100644 --- a/tests/lib/apply.test.ts +++ b/tests/lib/apply.test.ts @@ -18,6 +18,7 @@ describe("resource ordering", () => { resource("ConfigMap"), resource("Secret"), resource("PersistentVolumeClaim"), + resource("StorageClass"), resource("Namespace"), resource("Ingress"), ]; @@ -25,6 +26,7 @@ describe("resource ordering", () => { "Namespace", "GarageBucket", "GarageKey", + "StorageClass", "PersistentVolumeClaim", "Secret", "ConfigMap", diff --git a/tests/lib/convert-storage-env.test.ts b/tests/lib/convert-storage-env.test.ts index b01996a..3101bc4 100644 --- a/tests/lib/convert-storage-env.test.ts +++ b/tests/lib/convert-storage-env.test.ts @@ -10,6 +10,7 @@ import { serviceToDeployment, volumesToConfigMaps, volumesToPvc, + volumesToStorageClasses, } from "../../lib/convert"; const temporaryDirectories: string[] = []; @@ -29,7 +30,7 @@ afterEach(async () => { }); describe("volume conversion", () => { - test("renders named volumes with default storage placement", () => { + test("renders named volumes with a default Longhorn storage class", () => { const claims = volumesToPvc("project", { volumes: ["data:/var/lib/data"], } as Service); @@ -38,11 +39,29 @@ describe("volume conversion", () => { metadata: { name: "project-data", namespace: "project" }, spec: { resources: { requests: { storage: "1Gi" } }, - diskTag: ["fast"], - replicaCount: 2, - dataLocality: "none", + storageClassName: expect.stringMatching( + /^kuber-longhorn-[a-f0-9]{12}$/, + ), }, }); + + expect( + volumesToStorageClasses({ volumes: ["data:/var/lib/data"] } as Service), + ).toEqual([ + expect.objectContaining({ + kind: "StorageClass", + metadata: { + name: claims[0]?.spec?.storageClassName, + labels: expect.any(Object), + }, + provisioner: "driver.longhorn.io", + parameters: { + diskSelector: "fast", + numberOfReplicas: "2", + dataLocality: "disabled", + }, + }), + ]); }); test("parses inline storage sizing and placement", () => { @@ -53,9 +72,18 @@ describe("volume conversion", () => { metadata: { name: "project-data" }, spec: { resources: { requests: { storage: "20Gi" } }, - diskTag: ["fast", "archive"], - replicaCount: 3, - dataLocality: "none", + storageClassName: expect.stringMatching(/^kuber-longhorn-/), + }, + }); + expect( + volumesToStorageClasses({ + volumes: ["data(20Gi on 3 fast,archive):/data"], + } as Service)[0], + ).toMatchObject({ + parameters: { + diskSelector: "archive,fast", + numberOfReplicas: "3", + dataLocality: "disabled", }, }); }); @@ -76,9 +104,7 @@ describe("volume conversion", () => { ); expect(claim?.spec).toMatchObject({ resources: { requests: { storage: "50Gi" } }, - diskTag: ["bulk", "archive"], - replicaCount: 3, - dataLocality: "none", + storageClassName: expect.stringMatching(/^kuber-longhorn-/), }); }); @@ -250,6 +276,9 @@ describe("whole Compose conversion", () => { expect( resources.filter((resource) => resource.kind === "PersistentVolumeClaim"), ).toHaveLength(1); + expect( + resources.filter((resource) => resource.kind === "StorageClass"), + ).toHaveLength(1); expect( resources.filter((resource) => resource.kind === "Deployment"), ).toHaveLength(2); diff --git a/tests/lib/database.test.ts b/tests/lib/database.test.ts index 5554273..7696456 100644 --- a/tests/lib/database.test.ts +++ b/tests/lib/database.test.ts @@ -1,11 +1,16 @@ -import { describe, expect, test } from "bun:test"; +import { afterEach, describe, expect, mock, spyOn, test } from "bun:test"; import type { ComposeSpecification, Service } from "../../schema/docker.d"; import { buildDatabaseUrl, + buildPostgresEnvironment, getComposePostgresClaims, getServicePostgresClaim, isPostgresVolumeEntry, + reconcilePostgresClaim, } from "../../lib/database"; +import { objectApi } from "../../lib/k8s"; + +afterEach(() => mock.restore()); describe("managed PostgreSQL claims", () => { test("supports short and explicit syntax", () => { @@ -94,4 +99,56 @@ describe("managed PostgreSQL claims", () => { "postgresql://user%40host:p%3Aa%2Fss%3F%23@c.database.svc.cluster.local:5432/my%2Fdatabase", ); }); + + test("injects Redis into services with managed PostgreSQL", () => { + expect( + buildPostgresEnvironment( + { + service: "app", + username: "app", + database: "app", + secretName: "postgres-app", + }, + { username: "app", password: "secret" }, + ), + ).toEqual({ + DATABASE_URL: + "postgresql://app:secret@c.database.svc.cluster.local:5432/app", + REDIS_URL: "redis://redis.database.svc.cluster.local", + }); + }); + + test("reconciles a selected claim before returning credentials", async () => { + const claim = { + service: "app", + username: "app", + database: "app", + secretName: "postgres-app", + }; + spyOn(objectApi, "read").mockImplementation(async (resource) => { + if (resource.kind === "Secret") { + return { + ...resource, + data: { + username: Buffer.from("app").toString("base64"), + password: Buffer.from("secret").toString("base64"), + }, + } as never; + } + return { ...resource, spec: { managed: { roles: [] } } } as never; + }); + const patch = spyOn(objectApi, "patch").mockImplementation( + async (resource) => resource as never, + ); + + expect(await reconcilePostgresClaim("project", claim)).toEqual({ + username: "app", + password: "secret", + }); + expect(patch.mock.calls.map(([resource]) => resource.kind)).toEqual([ + "Secret", + "Cluster", + "Database", + ]); + }); }); diff --git a/tests/lib/k8s-http.test.ts b/tests/lib/k8s-http.test.ts new file mode 100644 index 0000000..8b1b538 --- /dev/null +++ b/tests/lib/k8s-http.test.ts @@ -0,0 +1,103 @@ +import { afterEach, describe, expect, test } from "bun:test"; +import { + HttpMethod, + RequestContext, +} from "@kubernetes/client-node/dist/gen/http/http.js"; +import http from "node:http"; +import { + createKubernetesHttpLibrary, + delay, + parseRetryAfter, +} from "../../lib/k8s-http"; + +const servers: http.Server[] = []; + +async function serve( + handler: http.RequestListener, +): Promise<{ server: http.Server; url: string }> { + const server = http.createServer(handler); + servers.push(server); + await new Promise((resolve) => server.listen(0, "127.0.0.1", resolve)); + const address = server.address(); + if (!address || typeof address === "string") throw new Error("Missing port"); + return { server, url: `http://127.0.0.1:${address.port}` }; +} + +afterEach(async () => { + await Promise.all( + servers.splice(0).map( + (server) => + new Promise((resolve, reject) => { + server.close((error) => (error ? reject(error) : resolve())); + }), + ), + ); +}); + +describe("Kubernetes HTTP transport", () => { + test("parses Retry-After seconds and HTTP dates", () => { + expect(parseRetryAfter("2", 0)).toBe(2000); + expect(parseRetryAfter("Thu, 01 Jan 1970 00:00:03 GMT", 1000)).toBe(2000); + expect(parseRetryAfter("invalid", 0)).toBeUndefined(); + }); + + test("retries 429 responses and honors the retry limit", async () => { + let requests = 0; + const { url } = await serve((_request, response) => { + requests += 1; + response.statusCode = requests < 3 ? 429 : 200; + response.setHeader("Retry-After", "0"); + response.end(requests < 3 ? "slow down" : "ok"); + }); + const client = createKubernetesHttpLibrary({ + minIntervalMs: 0, + baseRetryMs: 0, + maxRetries: 3, + }); + + const response = await client + .send(new RequestContext(url, HttpMethod.GET)) + .toPromise(); + expect(response.httpStatusCode).toBe(200); + expect(requests).toBe(3); + + requests = 0; + const limitedClient = createKubernetesHttpLibrary({ + minIntervalMs: 0, + baseRetryMs: 0, + maxRetries: 1, + }); + const limitedResponse = await limitedClient + .send(new RequestContext(url, HttpMethod.GET)) + .toPromise(); + expect(limitedResponse.httpStatusCode).toBe(429); + expect(requests).toBe(2); + }); + + test("paces concurrent request starts", async () => { + const starts: number[] = []; + const { url } = await serve((_request, response) => { + starts.push(Date.now()); + response.end("ok"); + }); + const client = createKubernetesHttpLibrary({ + maxConcurrent: 3, + minIntervalMs: 25, + }); + + await Promise.all( + Array.from({ length: 3 }, () => + client.send(new RequestContext(url, HttpMethod.GET)).toPromise(), + ), + ); + expect(starts).toHaveLength(3); + expect(starts[2]! - starts[0]!).toBeGreaterThanOrEqual(35); + }); + + test("cancels retry waits", async () => { + const controller = new AbortController(); + const waiting = delay(10_000, controller.signal); + controller.abort(); + await expect(waiting).rejects.toMatchObject({ name: "AbortError" }); + }); +});