diff --git a/README.md b/README.md index a76808d..a72e754 100644 --- a/README.md +++ b/README.md @@ -12,10 +12,10 @@ It reads a local Compose file, renders Kubernetes resources, applies them to the - Turns `env_file` entries into Kubernetes Secrets and mounts them through `envFrom` - Translates file mounts into ConfigMaps and directory/volume mounts into PVC-backed volumes - Builds images on the remote builder by syncing: - - committed git state - - tracked local diffs - - untracked files - - ignored `.env*` files + - committed git state + - tracked local diffs + - untracked files + - ignored `.env*` files - Supports host-based `ports` syntax that renders Kubernetes `Ingress` rules - Supports managed Postgres claims through special pseudo-volumes such as `postgresql:app` @@ -91,6 +91,22 @@ services: That produces a Kubernetes `Ingress` rule for `somedomain.astrxl.dev` pointing at the service port for container port `3000`. +Protected routes use the kuber dialect and render Traefik `IngressRoute` resources instead of plain Kubernetes `Ingress`: + +```yml +services: + app: + ports: + - db.astrxl.dev:4984:protected + - status.astrxl.dev:3001:protected(/dashboard,/socket.io) +``` + +Translation rules: + +- `host:port` -> Kubernetes `Ingress` +- `host:port:protected` -> Traefik `IngressRoute` with middleware `routing/cf-auth` and host-wide matching +- `host:port:protected(path1,path2,...)` -> Traefik `IngressRoute` with middleware `routing/cf-auth` and explicit `PathPrefix(...)` matches only + ### Managed Postgres You can declare a managed Postgres database with a pseudo-volume: @@ -129,6 +145,77 @@ This is also where generated values such as `DATABASE_URL` are injected. - `tmpfs` becomes `emptyDir` with memory backing - `postgresql:...` is treated as a managed database claim, not as a filesystem mount +Named volumes also support kuber-specific storage hints. + +Default behavior: + +```yml +# compose +services: + app: + volumes: + - myvolume:/data + +# effective kuber interpretation +services: + app: + volumes: + - myvolume(1Gi on 2 fast):/data +``` + +Short syntax: + +```yml +services: + app: + volumes: + - data(20Gi):/data + - archive(200Gi on archive):/archive + - cache(10Gi on 1 fast):/cache +``` + +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"` + +Explicit extensions are also supported. + +Top-level named volume: + +```yml +volumes: + data: + x-size: 20Gi + x-diskTag: [archive] + x-replicaCount: 1 + x-dataLocality: none +``` + +Long-form service mount: + +```yml +services: + app: + volumes: + - type: volume + source: data + target: /data + volume: + x-size: 20Gi + x-diskTag: [archive] + x-replicaCount: 1 + x-dataLocality: none +``` + +Precedence: + +- short syntax like `data(20Gi on 1 archive):/data` +- long-form `volume.x-*` +- top-level `volumes..x-*` +- fallback default `1Gi on 2 fast` + ## Building To build distributable binaries: @@ -149,4 +236,3 @@ bun build index.ts --target bun --minify --sourcemap --outdir dist - The remote build flow is optimized for local iteration, not for producing a perfectly clean export of the repository. - Managed database support is Kubernetes-only. It injects `DATABASE_URL` into the generated app secret and does not rewrite local `.env` files. - `kuber` operates on managed resources in the namespace matching the current directory name. - diff --git a/command/exec.ts b/command/exec.ts index 6f59faa..104c96a 100644 --- a/command/exec.ts +++ b/command/exec.ts @@ -2,6 +2,17 @@ import { defineCommand } from "citty"; import { getPodContainerName, getPodForDeployment } from "../lib/shared"; import { execClient } from "../lib/k8s"; +function prepareInteractiveStdin(): (() => void) | undefined { + if (!process.stdin.isTTY) return; + + process.stdin.setRawMode?.(true); + process.stdin.resume(); + return () => { + process.stdin.setRawMode?.(false); + process.stdin.pause(); + }; +} + export const exec = defineCommand({ meta: { name: "exec", @@ -23,21 +34,27 @@ export const exec = defineCommand({ const podName = pod.metadata?.name; if (!podName) throw new Error(`No pod name found for deployment ${deployment}`); - const socket = await execClient.exec( - namespace, - podName, - getPodContainerName(pod), - command, - process.stdout, - process.stderr, - process.stdin, - Boolean(process.stdin.isTTY), - ); + const restoreStdin = prepareInteractiveStdin(); - await new Promise((resolve, reject) => { - socket.onclose = () => resolve(); - socket.onerror = (event: { error?: unknown }) => - reject(event.error ?? new Error("Exec failed")); - }); + try { + const socket = await execClient.exec( + namespace, + podName, + getPodContainerName(pod), + command, + process.stdout, + process.stderr, + process.stdin, + Boolean(process.stdin.isTTY), + ); + + await new Promise((resolve, reject) => { + socket.onclose = () => resolve(); + socket.onerror = (event: { error?: unknown }) => + reject(event.error ?? new Error("Exec failed")); + }); + } finally { + restoreStdin?.(); + } }, }); diff --git a/command/export.ts b/command/export.ts new file mode 100644 index 0000000..eb36dd9 --- /dev/null +++ b/command/export.ts @@ -0,0 +1,34 @@ +import { defineCommand } from "citty"; +import { resolve } from "node:path"; +import { YAML } from "bun"; +import { sortResources } from "../lib/apply"; +import { ctx } from "../lib/context"; +import { renderResources } from "../lib/render"; + +function toDocumentYaml(resources: Parameters[0]): string { + return `${sortResources(resources) + .map((resource) => YAML.stringify(resource, null, 2).trimEnd()) + .join("\n---\n")}\n`; +} + +export const exportCommand = defineCommand({ + meta: { + name: "export", + description: "Write rendered manifests to a YAML file", + }, + args: { + output: { + type: "string", + alias: "o", + description: "Output YAML file path", + }, + }, + async run({ args }) { + const { project, compose, cwd } = ctx(); + const resources = await renderResources(project, await compose(), cwd); + const outputPath = resolve(cwd, args.output || `${project}.yaml`); + + await Bun.write(outputPath, toDocumentYaml(resources)); + console.log(outputPath); + }, +}); diff --git a/index.ts b/index.ts index a9ec2c5..4968197 100644 --- a/index.ts +++ b/index.ts @@ -1,6 +1,7 @@ import { defineCommand, runMain } from "citty"; import { db } from "./command/db"; import { down } from "./command/down"; +import { exportCommand } from "./command/export"; import { exec } from "./command/exec"; import { logs } from "./command/logs"; import { ps } from "./command/ps"; @@ -32,7 +33,18 @@ const main = defineCommand({ version: "1.0.0", description: "Docker Compose -> K8s translation layer", }, - subCommands: { db, down, exec, logs, ps, restart, start, stop, up }, + subCommands: { + db, + down, + export: exportCommand, + exec, + logs, + ps, + restart, + start, + stop, + up, + }, }); provideContext(() => runMain(main)).catch((e) => { diff --git a/lib/apply.ts b/lib/apply.ts index 9068ed7..fff84ec 100644 --- a/lib/apply.ts +++ b/lib/apply.ts @@ -13,6 +13,7 @@ const ResourceOrder = { Service: 4, Deployment: 5, Ingress: 6, + IngressRoute: 7, } as const; const ManagedResources = [ @@ -22,6 +23,7 @@ const ManagedResources = [ { apiVersion: "v1", kind: "Service" }, { apiVersion: "apps/v1", kind: "Deployment" }, { apiVersion: "networking.k8s.io/v1", kind: "Ingress" }, + { apiVersion: "traefik.io/v1alpha1", kind: "IngressRoute" }, ] as const; const ManagedBySelector = `app.kubernetes.io/managed-by=${LABELS["app.kubernetes.io/managed-by"]}`; @@ -51,7 +53,8 @@ export async function getNamespaceManagementStatus( metadata: { name: namespace }, }); - return existing.metadata?.labels?.["app.kubernetes.io/managed-by"] === ManagedByLabel + return existing.metadata?.labels?.["app.kubernetes.io/managed-by"] === + ManagedByLabel ? "managed" : "external"; } catch (error) { @@ -102,11 +105,15 @@ export async function applyResources( return applied; } -export async function deleteResource(resource: KubernetesObject): Promise { +export async function deleteResource( + resource: KubernetesObject, +): Promise { await objectApi.delete(resource); } -export async function deleteResources(resources: KubernetesObject[]): Promise { +export async function deleteResources( + resources: KubernetesObject[], +): Promise { for (const resource of sortResources(resources).reverse()) { await deleteResource(resource); } diff --git a/lib/arch.ts b/lib/arch.ts new file mode 100644 index 0000000..5b13525 --- /dev/null +++ b/lib/arch.ts @@ -0,0 +1,36 @@ +import { AsyncLocalStorage } from "node:async_hooks"; +import type { ComposeSpecification } from "../schema/docker.d"; + +export type ComposeArch = "amd64" | "arm64"; + +const composeArch = new AsyncLocalStorage(); + +export function resolveComposeArch(compose: ComposeSpecification): ComposeArch { + return compose["x-arch"] === "amd64" ? "amd64" : "arm64"; +} + +export function withComposeArch(compose: ComposeSpecification, fn: () => T): T { + return composeArch.run(resolveComposeArch(compose), fn); +} + +export function getComposeArch(): ComposeArch { + return composeArch.getStore() ?? "arm64"; +} + +export function getComposeArchPlacement(compose: ComposeSpecification) { + if (resolveComposeArch(compose) !== "amd64") return {}; + + return { + nodeSelector: { + "kubernetes.io/arch": "amd64", + }, + tolerations: [ + { + key: "arch", + operator: "Equal", + value: "amd64", + effect: "NoExecute", + }, + ], + }; +} diff --git a/lib/build.ts b/lib/build.ts index bb841e6..710c577 100644 --- a/lib/build.ts +++ b/lib/build.ts @@ -4,10 +4,19 @@ import { hostname, tmpdir } from "node:os"; import type { Writable } from "node:stream"; import type { ComposeSpecification, Service } from "../schema/docker.d"; import { IMAGE_REGISTRY } from "../const"; +import { getComposeArch, withComposeArch } from "./arch"; -const REMOTE_BUILDER = "kuber@astral"; +const REMOTE_BUILDER_ARM = "kuber@astral"; +const REMOTE_BUILDER_AMD = "kuber@astral-th"; const REMOTE_BUILD_ROOT = "kuber-build"; +const RemoteBuilder = { + amd64: REMOTE_BUILDER_AMD, + arm64: REMOTE_BUILDER_ARM, +} as const; + +type RemoteBuilder = (typeof RemoteBuilder)[keyof typeof RemoteBuilder]; + type BuildPlan = { name: string; image: string; @@ -29,6 +38,22 @@ type SpawnResult = { stderr: string; }; +function summarizeCommandFailure( + exitCode: number, + stderrLines: string[], + stdoutLines: string[], +): string { + const lines = (stderrLines.length > 0 ? stderrLines : stdoutLines) + .map((line) => line.trimEnd()) + .filter(Boolean); + if (lines.length === 0) return `Command failed with exit code ${exitCode}`; + + const preview = lines.slice(0, 8).join("\n"); + return lines.length > 8 + ? `${preview}\n... (${lines.length - 8} more lines)` + : preview; +} + function toReadableStream( stream: number | ReadableStream | undefined, ): ReadableStream | undefined { @@ -40,7 +65,11 @@ function shellQuote(value: string): string { } function resolveBuildArgs(service: Service): string[] { - if (!service.build || typeof service.build === "string" || !service.build.args) { + if ( + !service.build || + typeof service.build === "string" || + !service.build.args + ) { return []; } @@ -62,16 +91,21 @@ function resolveBuildPlan( if (!service.build) return; const build = service.build; - const contextInput = typeof build === "string" ? build : build.context ?? "."; + const contextInput = + typeof build === "string" ? build : (build.context ?? "."); if (contextInput.includes("://")) { - throw new Error(`Remote build context is not supported for service ${name}`); + throw new Error( + `Remote build context is not supported for service ${name}`, + ); } const contextPath = resolve(cwd, contextInput); const contextRelative = relative(repoRoot, contextPath); if (contextRelative.startsWith("..") || isAbsolute(contextRelative)) { - throw new Error(`Build context must stay inside the git repo for service ${name}`); + throw new Error( + `Build context must stay inside the git repo for service ${name}`, + ); } if (typeof build !== "string" && build.dockerfile_inline) { @@ -86,22 +120,35 @@ function resolveBuildPlan( ? relative(repoRoot, dockerfilePath) : undefined; - if (dockerfileRelative?.startsWith("..") || isAbsolute(dockerfileRelative ?? "")) { - throw new Error(`Dockerfile must stay inside the git repo for service ${name}`); + if ( + dockerfileRelative?.startsWith("..") || + isAbsolute(dockerfileRelative ?? "") + ) { + throw new Error( + `Dockerfile must stay inside the git repo for service ${name}`, + ); } return { name, image: `${IMAGE_REGISTRY}/kuber/${project}-${name}:latest`, context: `${buildRoot}/${contextRelative === "" ? "." : contextRelative}`, - dockerfile: dockerfileRelative ? `${buildRoot}/${dockerfileRelative}` : undefined, + dockerfile: dockerfileRelative + ? `${buildRoot}/${dockerfileRelative}` + : undefined, target: typeof build === "string" ? undefined : build.target, buildArgs: resolveBuildArgs(service), }; } +function getRemoteBuilder(): RemoteBuilder { + return getComposeArch() === "amd64" + ? RemoteBuilder.amd64 + : RemoteBuilder.arm64; +} + function isOnRemoteBuilder(): boolean { - return hostname() === REMOTE_BUILDER.split("@").at(-1); + return hostname() === getRemoteBuilder().split("@").at(-1); } async function pumpStream( @@ -144,24 +191,56 @@ async function runWithOutput( ) { const stdoutLines: string[] = []; const stderrLines: string[] = []; + let streamBuffer = ""; + let flushTimer: ReturnType | undefined; + + function flushStreamBuffer() { + if (!reporter?.stream || streamBuffer.length === 0) return; + reporter.stream.write(streamBuffer); + streamBuffer = ""; + } + + function queueStreamLine(line: string) { + streamBuffer += `${line}\n`; + if (streamBuffer.length >= 8192) { + if (flushTimer) { + clearTimeout(flushTimer); + flushTimer = undefined; + } + flushStreamBuffer(); + return; + } + + if (flushTimer) return; + flushTimer = setTimeout(() => { + flushTimer = undefined; + flushStreamBuffer(); + }, 33); + } await Promise.all([ pumpStream(toReadableStream(command.stdout), async (line) => { stdoutLines.push(line); - if (reporter?.stream) reporter.stream.write(`${line}\n`); + if (reporter?.stream) queueStreamLine(line); else if (reporter?.progress) await reporter.progress(line); }), pumpStream(toReadableStream(command.stderr), async (line) => { stderrLines.push(line); - if (reporter?.stream) reporter.stream.write(`${line}\n`); + if (reporter?.stream) queueStreamLine(line); else if (reporter?.progress) await reporter.progress(line); }), ]); + if (flushTimer) { + clearTimeout(flushTimer); + flushTimer = undefined; + } + flushStreamBuffer(); + const exitCode = await command.exited; if (exitCode !== 0) { throw new Error( - stderrLines.join("\n") || stdoutLines.join("\n") || `Command failed with exit code ${exitCode}`, + summarizeCommandFailure(exitCode, stderrLines, stdoutLines), ); } @@ -174,21 +253,23 @@ async function runWithOutput( function ssh(script: string) { const remoteCommand = `bash -lc ${shellQuote(script)}`; - return Bun.spawn(["ssh", REMOTE_BUILDER, remoteCommand], { + return Bun.spawn(["ssh", getRemoteBuilder(), remoteCommand], { stdout: "pipe", stderr: "pipe", }); } function scp(localPath: string, remotePath: string) { - return Bun.spawn(["scp", localPath, `${REMOTE_BUILDER}:${remotePath}`], { + return Bun.spawn(["scp", localPath, `${getRemoteBuilder()}:${remotePath}`], { stdout: "pipe", stderr: "pipe", }); } async function getRepoRoot(cwd: string): Promise { - return Bun.$.cwd(cwd)`git rev-parse --show-toplevel`.text().then((e) => e.trim()); + return Bun.$.cwd(cwd)`git rev-parse --show-toplevel` + .text() + .then((e) => e.trim()); } async function getHeadSha(repoRoot: string): Promise { @@ -243,13 +324,8 @@ async function getIgnoredDotenvFiles(repoRoot: string): Promise { } async function getRemoteHead(remoteRepo: string): Promise { - const result = Bun.spawn( - [ - "ssh", - REMOTE_BUILDER, - `bash -lc ${shellQuote(`if [ -d ${shellQuote(`${remoteRepo}/.git`)} ]; then git -C ${shellQuote(remoteRepo)} rev-parse HEAD; fi`)}`, - ], - { stdout: "pipe", stderr: "pipe" }, + const result = ssh( + `if [ -d ${shellQuote(`${remoteRepo}/.git`)} ]; then git -C ${shellQuote(remoteRepo)} rev-parse HEAD; fi`, ); const output = await runWithOutput(result); @@ -273,26 +349,35 @@ async function syncRemoteRepo( const remoteBundle = `${REMOTE_BUILD_ROOT}/${basename(repoRoot)}.bundle`; await Bun.$.cwd(repoRoot)`git bundle create ${bundlePath} HEAD`; - await runWithOutput(ssh(`mkdir -p ${shellQuote(REMOTE_BUILD_ROOT)}`), reporter); + await runWithOutput( + ssh(`mkdir -p ${shellQuote(REMOTE_BUILD_ROOT)}`), + reporter, + ); await runWithOutput(scp(bundlePath, remoteBundle), reporter); - await runWithOutput(ssh( - [ - "set -euo pipefail", - `mkdir -p ${shellQuote(remoteRepo)}`, - `if [ ! -d ${shellQuote(`${remoteRepo}/.git`)} ]; then git -C ${shellQuote(remoteRepo)} init; fi`, - `git -C ${shellQuote(remoteRepo)} fetch --force "$PWD/${remoteBundle}" HEAD`, - `git -C ${shellQuote(remoteRepo)} reset --hard FETCH_HEAD`, - `git -C ${shellQuote(remoteRepo)} clean -fd`, - ].join("; "), - ), reporter); + await runWithOutput( + ssh( + [ + "set -euo pipefail", + `mkdir -p ${shellQuote(remoteRepo)}`, + `if [ ! -d ${shellQuote(`${remoteRepo}/.git`)} ]; then git -C ${shellQuote(remoteRepo)} init; fi`, + `git -C ${shellQuote(remoteRepo)} fetch --force "$PWD/${remoteBundle}" HEAD`, + `git -C ${shellQuote(remoteRepo)} reset --hard FETCH_HEAD`, + `git -C ${shellQuote(remoteRepo)} clean -fd`, + ].join("; "), + ), + reporter, + ); } else { - await runWithOutput(ssh( - [ - "set -euo pipefail", - `git -C ${shellQuote(remoteRepo)} reset --hard HEAD`, - `git -C ${shellQuote(remoteRepo)} clean -fd`, - ].join("; "), - ), reporter); + await runWithOutput( + ssh( + [ + "set -euo pipefail", + `git -C ${shellQuote(remoteRepo)} reset --hard HEAD`, + `git -C ${shellQuote(remoteRepo)} clean -fd`, + ].join("; "), + ), + reporter, + ); } const diff = await getTrackedDiff(repoRoot); @@ -302,12 +387,15 @@ async function syncRemoteRepo( const remoteDiff = `${REMOTE_BUILD_ROOT}/${basename(repoRoot)}.diff`; await Bun.write(diffPath, diff); await runWithOutput(scp(diffPath, remoteDiff), reporter); - await runWithOutput(ssh( - [ - "set -euo pipefail", - `git -C ${shellQuote(remoteRepo)} apply --allow-binary-replacement "$PWD/${remoteDiff}"`, - ].join("; "), - ), reporter); + await runWithOutput( + ssh( + [ + "set -euo pipefail", + `git -C ${shellQuote(remoteRepo)} apply --allow-binary-replacement "$PWD/${remoteDiff}"`, + ].join("; "), + ), + reporter, + ); } const untrackedFiles = await getUntrackedFiles(repoRoot); @@ -325,22 +413,31 @@ async function syncRemoteRepo( const untrackedArchive = `${tempDir}/untracked.tar`; const remoteUntrackedArchive = `${REMOTE_BUILD_ROOT}/${basename(repoRoot)}-untracked.tar`; await Bun.$`tar -C ${untrackedDir} -cf ${untrackedArchive} .`; - await runWithOutput(scp(untrackedArchive, remoteUntrackedArchive), reporter); - await runWithOutput(ssh( - [ - "set -euo pipefail", - `tar -C ${shellQuote(remoteRepo)} -xf "$PWD/${remoteUntrackedArchive}"`, - ].join("; "), - ), reporter); + await runWithOutput( + scp(untrackedArchive, remoteUntrackedArchive), + reporter, + ); + await runWithOutput( + ssh( + [ + "set -euo pipefail", + `tar -C ${shellQuote(remoteRepo)} -xf "$PWD/${remoteUntrackedArchive}"`, + ].join("; "), + ), + reporter, + ); } const ignoredDotenvFiles = await getIgnoredDotenvFiles(repoRoot); - await runWithOutput(ssh( - [ - "set -euo pipefail", - `git -C ${shellQuote(remoteRepo)} clean -fdX -- .env* '**/.env*'`, - ].join("; "), - ), reporter); + await runWithOutput( + ssh( + [ + "set -euo pipefail", + `git -C ${shellQuote(remoteRepo)} clean -fdX -- .env* '**/.env*'`, + ].join("; "), + ), + reporter, + ); if (ignoredDotenvFiles.length === 0) return; @@ -358,12 +455,15 @@ async function syncRemoteRepo( const remoteDotenvArchive = `${REMOTE_BUILD_ROOT}/${basename(repoRoot)}-dotenv.tar`; await Bun.$`tar -C ${dotenvDir} -cf ${dotenvArchive} .`; await runWithOutput(scp(dotenvArchive, remoteDotenvArchive), reporter); - await runWithOutput(ssh( - [ - "set -euo pipefail", - `tar -C ${shellQuote(remoteRepo)} -xf "$PWD/${remoteDotenvArchive}"`, - ].join("; "), - ), reporter); + await runWithOutput( + ssh( + [ + "set -euo pipefail", + `tar -C ${shellQuote(remoteRepo)} -xf "$PWD/${remoteDotenvArchive}"`, + ].join("; "), + ), + reporter, + ); } finally { await rm(tempDir, { recursive: true, force: true }); } @@ -377,7 +477,14 @@ function getBuildPlans( remoteRepo: string, ): BuildPlan[] { return Object.entries(compose.services ?? {}).flatMap(([name, service]) => { - const plan = resolveBuildPlan(project, name, service, cwd, repoRoot, remoteRepo); + const plan = resolveBuildPlan( + project, + name, + service, + cwd, + repoRoot, + remoteRepo, + ); return plan ? [plan] : []; }); } @@ -431,23 +538,29 @@ export async function buildServices( cwd = process.cwd(), reporter?: BuildReporter, ): Promise { - const repoRoot = await getRepoRoot(cwd); - const localBuilder = isOnRemoteBuilder(); - const buildRoot = localBuilder - ? repoRoot - : `${REMOTE_BUILD_ROOT}/${basename(repoRoot)}`; - const plans = getBuildPlans(project, compose, cwd, repoRoot, buildRoot); + return withComposeArch(compose, async () => { + if (!Object.values(compose.services ?? {}).some((service) => service.build)) { + return 0; + } - if (plans.length === 0) return 0; + const repoRoot = await getRepoRoot(cwd); + const localBuilder = isOnRemoteBuilder(); + const buildRoot = localBuilder + ? repoRoot + : `${REMOTE_BUILD_ROOT}/${basename(repoRoot)}`; + const plans = getBuildPlans(project, compose, cwd, repoRoot, buildRoot); - if (!localBuilder) { - await syncRemoteRepo(repoRoot, buildRoot, reporter); - } + if (plans.length === 0) return 0; - for (const plan of plans) { - if (localBuilder) await buildLocal(plan, reporter); - else await buildRemote(plan, reporter); - } + if (!localBuilder) { + await syncRemoteRepo(repoRoot, buildRoot, reporter); + } - return plans.length; + for (const plan of plans) { + if (localBuilder) await buildLocal(plan, reporter); + else await buildRemote(plan, reporter); + } + + return plans.length; + }); } diff --git a/lib/convert.ts b/lib/convert.ts index 0026384..eee295f 100644 --- a/lib/convert.ts +++ b/lib/convert.ts @@ -19,6 +19,7 @@ import { basename, resolve } from "node:path"; import type { ComposeSpecification } from "../schema/docker.d"; import type { Service } from "../schema/docker.d"; import { IMAGE_REGISTRY, LABELS } from "../const"; +import { getComposeArchPlacement } from "./arch"; import { getServicePostgresClaim, isPostgresVolumeEntry } from "./database"; import { toEnvVars } from "./format"; import { deepMerge } from "./shared"; @@ -32,6 +33,10 @@ type NormalizedMount = { configKey?: string; subPath?: string; sizeLimit?: string; + requestedStorage?: string; + diskTag?: string[]; + replicaCount?: number; + dataLocality?: "none"; }; type NormalizedPort = { @@ -41,9 +46,19 @@ type NormalizedPort = { hostPort?: number; protocol: "TCP" | "UDP"; host?: string; + routingKind?: "ingress" | "ingressroute"; + paths?: string[]; }; type ServiceVolume = NonNullable[number]; +type ComposeVolumes = ComposeSpecification["volumes"]; + +const DefaultNamedVolumeStorage = { + requestedStorage: "1Gi", + diskTag: ["fast"], + replicaCount: 2, + dataLocality: "none" as const, +}; function toPortNumber(value: number | string | undefined): number | undefined { if (typeof value === "number") return value; @@ -105,6 +120,126 @@ function toKubeName(value: string): string { .slice(0, 63); } +function normalizeStorageSize(value: unknown): string | undefined { + if (value === undefined || value === null) return; + const size = String(value).trim(); + if (!size) return; + return size; +} + +function normalizeDiskTags(value: unknown): string[] | undefined { + if (value === undefined || value === null) return; + + const tags = (Array.isArray(value) ? value : [value]) + .flatMap((entry) => String(entry).split(",")) + .map((entry) => entry.trim()) + .filter(Boolean); + + return tags.length > 0 ? [...new Set(tags)] : undefined; +} + +function normalizeReplicaCount(value: unknown): number | undefined { + if (value === undefined || value === null || value === "") return; + const count = Number(value); + return Number.isInteger(count) ? count : undefined; +} + +function parseStorageSpec(value: string): { + requestedStorage?: string; + diskTag?: string[]; + replicaCount?: number; + dataLocality?: "none"; +} { + const [sizePart, placementPart] = value + .split(/\s+on\s+/i, 2) + .map((part) => part?.trim()); + + if (!placementPart) { + return { + requestedStorage: normalizeStorageSize(sizePart), + }; + } + + const placementMatch = /^(?:(?\d+)\s+)?(?.+)$/.exec( + placementPart, + ); + + return { + requestedStorage: normalizeStorageSize(sizePart), + diskTag: normalizeDiskTags(placementMatch?.groups?.tags), + replicaCount: normalizeReplicaCount(placementMatch?.groups?.replicas), + dataLocality: "none", + }; +} + +function parseNamedVolumeSize(source: string): + | { + source: string; + requestedStorage?: string; + diskTag?: string[]; + replicaCount?: number; + dataLocality?: "none"; + } + | undefined { + const match = /^(?[^()]+)\((?[^()]+)\)$/.exec(source.trim()); + if (!match?.groups) return; + + const name = match.groups.name?.trim(); + const { requestedStorage, diskTag, replicaCount, dataLocality } = + parseStorageSpec(match.groups.size ?? ""); + if (!name) return; + + return { + source: name, + requestedStorage, + diskTag, + replicaCount, + dataLocality, + }; +} + +function resolveVolumeSource(source: string | undefined): { + source: string | undefined; + requestedStorage?: string; + diskTag?: string[]; + replicaCount?: number; + dataLocality?: "none"; +} { + if (!source || isBindSource(source)) return { source }; + + const sized = parseNamedVolumeSize(source); + if (!sized) return { source }; + if (!sized.source) { + throw new Error(`Invalid volume source ${source}. Expected name(size).`); + } + + return sized; +} + +function getVolumeExtensionSize( + volume: { [key: string]: unknown } | undefined, +): string | undefined { + return normalizeStorageSize(volume?.["x-size"]); +} + +function getVolumeExtensionDiskTag( + volume: { [key: string]: unknown } | undefined, +): string[] | undefined { + return normalizeDiskTags(volume?.["x-diskTag"]); +} + +function getVolumeExtensionReplicaCount( + volume: { [key: string]: unknown } | undefined, +): number | undefined { + return normalizeReplicaCount(volume?.["x-replicaCount"]); +} + +function getVolumeExtensionDataLocality( + volume: { [key: string]: unknown } | undefined, +): "none" | undefined { + return volume?.["x-dataLocality"] === "none" ? "none" : undefined; +} + function toBindVolumeName(source: string): string { return toKubeName(`bind-${Bun.hash(source).toString(36)}`); } @@ -149,10 +284,13 @@ function parseStringMount( const maybeMode = parts.at(-1); const hasMode = maybeMode === "ro" || maybeMode === "rw"; const target = parts.at(hasMode ? -2 : -1); - const source = parts.slice(0, hasMode ? -2 : -1).join(":"); + const rawSource = parts.slice(0, hasMode ? -2 : -1).join(":"); if (!target) return; + const { source, requestedStorage, diskTag, replicaCount, dataLocality } = + resolveVolumeSource(rawSource || undefined); + return { name: source && isBindSource(source) @@ -168,6 +306,10 @@ function parseStringMount( target, configKey: source ? basename(source) : undefined, readOnly: hasMode ? maybeMode === "ro" : undefined, + requestedStorage, + diskTag, + replicaCount, + dataLocality, }; } @@ -187,35 +329,114 @@ function toMount( return; } + const { source, requestedStorage, diskTag, replicaCount, dataLocality } = + resolveVolumeSource(entry.source); + return { - name: entry.source + name: source ? entry.type === "bind" - ? toBindVolumeName(entry.source) - : `${entry.type}-${toKubeName(entry.source) || index}` + ? toBindVolumeName(source) + : `${entry.type}-${toKubeName(source) || index}` : `${entry.type}-${index}`, kind: - entry.type === "bind" && - entry.source && - isConfigFileSource(entry.source, cwd) + entry.type === "bind" && source && isConfigFileSource(source, cwd) ? "config-file" : entry.type, - source: entry.source, + source, target: entry.target, readOnly: toBoolean(entry.read_only), - configKey: entry.source ? basename(entry.source) : undefined, + configKey: source ? basename(source) : undefined, subPath: entry.volume?.subpath, sizeLimit: entry.type === "tmpfs" && entry.tmpfs?.size !== undefined ? String(entry.tmpfs.size) : undefined, + requestedStorage: requestedStorage ?? getVolumeExtensionSize(entry.volume), + diskTag: diskTag ?? getVolumeExtensionDiskTag(entry.volume), + replicaCount: replicaCount ?? getVolumeExtensionReplicaCount(entry.volume), + dataLocality: dataLocality ?? getVolumeExtensionDataLocality(entry.volume), }; } -function toMounts(service: Service, cwd = process.cwd()): NormalizedMount[] { +function getTopLevelVolumeSize( + volumes: ComposeVolumes, + source: string | undefined, +): string | undefined { + if (!source || !volumes) return; + + const volume = volumes[source] as { [key: string]: unknown } | undefined; + return getVolumeExtensionSize(volume); +} + +function getTopLevelVolumeDiskTag( + volumes: ComposeVolumes, + source: string | undefined, +): string[] | undefined { + if (!source || !volumes) return; + + const volume = volumes[source] as { [key: string]: unknown } | undefined; + return getVolumeExtensionDiskTag(volume); +} + +function getTopLevelVolumeReplicaCount( + volumes: ComposeVolumes, + source: string | undefined, +): number | undefined { + if (!source || !volumes) return; + + const volume = volumes[source] as { [key: string]: unknown } | undefined; + return getVolumeExtensionReplicaCount(volume); +} + +function getTopLevelVolumeDataLocality( + volumes: ComposeVolumes, + source: string | undefined, +): "none" | undefined { + if (!source || !volumes) return; + + const volume = volumes[source] as { [key: string]: unknown } | undefined; + return getVolumeExtensionDataLocality(volume); +} + +function toMounts( + service: Service, + cwd = process.cwd(), + volumes: ComposeVolumes = {}, +): NormalizedMount[] { const mounts = service.volumes?.flatMap((entry, index) => { const mount = toMount(entry, index, cwd); - return mount ? [mount] : []; + if (!mount) return []; + + return [ + { + ...mount, + requestedStorage: + mount.requestedStorage ?? + getTopLevelVolumeSize(volumes, mount.source) ?? + (mount.kind === "volume" && mount.source + ? DefaultNamedVolumeStorage.requestedStorage + : undefined), + diskTag: + mount.diskTag ?? + getTopLevelVolumeDiskTag(volumes, mount.source) ?? + (mount.kind === "volume" && mount.source + ? DefaultNamedVolumeStorage.diskTag + : undefined), + replicaCount: + mount.replicaCount ?? + getTopLevelVolumeReplicaCount(volumes, mount.source) ?? + (mount.kind === "volume" && mount.source + ? DefaultNamedVolumeStorage.replicaCount + : undefined), + dataLocality: + mount.dataLocality ?? + getTopLevelVolumeDataLocality(volumes, mount.source) ?? + (mount.kind === "volume" && mount.source + ? DefaultNamedVolumeStorage.dataLocality + : undefined), + }, + ]; }) ?? []; const tmpfs = (Array.isArray(service.tmpfs) ? service.tmpfs : [service.tmpfs]) @@ -300,6 +521,47 @@ function toVolumes( }); } +function normalizeProtectedPaths(value: string): string[] { + return [ + ...new Set( + value + .split(",") + .map((segment) => segment.trim()) + .filter(Boolean) + .map((segment) => (segment.startsWith("/") ? segment : `/${segment}`)), + ), + ]; +} + +function parsePortRoute(value: string): { + portSpec: string; + routingKind?: "ingress" | "ingressroute"; + paths?: string[]; +} { + const protectedMatch = + /^(?.+):protected(?:\((?[^)]*)\))?$/.exec(value); + if (!protectedMatch?.groups) return { portSpec: value }; + + const portSpec = protectedMatch.groups.portSpec?.trim(); + if (!portSpec) { + throw new Error(`Invalid protected port syntax ${value}.`); + } + + const rawPaths = protectedMatch.groups.paths?.trim(); + if (rawPaths === undefined) { + return { portSpec, routingKind: "ingressroute" }; + } + + const paths = normalizeProtectedPaths(rawPaths); + if (paths.length === 0) { + throw new Error( + `Invalid protected port syntax ${value}. Expected one or more paths.`, + ); + } + + return { portSpec, routingKind: "ingressroute", paths }; +} + function toPorts(service: Service): NormalizedPort[] { const ports = service.ports?.flatMap((entry, index) => { @@ -315,8 +577,10 @@ function toPorts(service: Service): NormalizedPort[] { } if (typeof entry === "string") { - const [portSpec, protocolSpec] = entry.split("/"); - if (!portSpec) return []; + const [rawPortSpec, protocolSpec] = entry.split("/"); + if (!rawPortSpec) return []; + + const { portSpec, routingKind, paths } = parsePortRoute(rawPortSpec); const segments = portSpec.split(":"); const target = toPortNumber(segments.at(-1)); @@ -336,6 +600,8 @@ function toPorts(service: Service): NormalizedPort[] { hostPort: published, host, protocol: toProtocol(protocolSpec), + routingKind: host ? (routingKind ?? "ingress") : undefined, + paths, }, ]; } @@ -353,6 +619,7 @@ function toPorts(service: Service): NormalizedPort[] { hostPort: published, host: toHostname(entry.host_ip), protocol: toProtocol(entry.protocol), + routingKind: toHostname(entry.host_ip) ? "ingress" : undefined, }, ]; }) ?? []; @@ -373,12 +640,7 @@ function toPorts(service: Service): NormalizedPort[] { ]; }) ?? []; - const deduped = new Map(); - for (const port of [...ports, ...exposedPorts]) { - deduped.set(`${port.containerPort}:${port.protocol}`, port); - } - - return [...deduped.values()]; + return [...ports, ...exposedPorts]; } function toContainerPorts( @@ -386,30 +648,49 @@ function toContainerPorts( ): V1ContainerPort[] | undefined { if (ports.length === 0) return; - return ports.map((port) => ({ - name: port.name, - containerPort: port.containerPort, - hostPort: port.hostPort, - protocol: port.protocol, - })); + const deduped = new Map(); + + for (const port of ports) { + deduped.set(`${port.containerPort}:${port.protocol}`, { + name: port.name, + containerPort: port.containerPort, + hostPort: port.hostPort, + protocol: port.protocol, + }); + } + + return [...deduped.values()]; } function toServicePorts(ports: NormalizedPort[]): V1ServicePort[] | undefined { if (ports.length === 0) return; - return ports.map((port) => ({ - name: port.name, - port: port.servicePort, - targetPort: port.containerPort, - protocol: port.protocol, - })); + const deduped = new Map(); + + for (const port of ports) { + deduped.set(`${port.servicePort}:${port.protocol}`, { + name: port.name, + port: port.servicePort, + targetPort: port.containerPort, + protocol: port.protocol, + }); + } + + return [...deduped.values()]; } function toIngressRules(name: string, ports: NormalizedPort[]) { const seenHosts = new Set(); return ports.flatMap((port) => { - if (!port.host || seenHosts.has(port.host)) return []; + if ( + !port.host || + port.routingKind === "ingressroute" || + seenHosts.has(port.host) + ) { + return []; + } + seenHosts.add(port.host); return [ @@ -423,7 +704,7 @@ function toIngressRules(name: string, ports: NormalizedPort[]) { backend: { service: { name, - port: { name: port.name }, + port: { number: port.servicePort }, }, }, }, @@ -434,6 +715,55 @@ function toIngressRules(name: string, ports: NormalizedPort[]) { }); } +function toIngressRoute( + project: string, + name: string, + ports: NormalizedPort[], +): KubernetesObject | undefined { + const routes = ports.flatMap((port) => { + if (!port.host || port.routingKind !== "ingressroute") return []; + + const matches = + port.paths && port.paths.length > 0 + ? port.paths.map( + (path) => `Host(\`${port.host}\`) && PathPrefix(\`${path}\`)`, + ) + : [`Host(\`${port.host}\`)`]; + + return matches.map((match) => ({ + kind: "Rule", + match, + middlewares: [ + { + name: "cf-auth", + namespace: "routing", + }, + ], + services: [ + { + name, + port: port.servicePort, + }, + ], + })); + }); + + if (routes.length === 0) return; + + return { + apiVersion: "traefik.io/v1alpha1", + kind: "IngressRoute", + metadata: { + name, + namespace: project, + labels: LABELS, + }, + spec: { + routes, + }, + } as KubernetesObject; +} + function toEnvFilePaths( envFile: Service["env_file"], ): { path: string; required: boolean }[] { @@ -502,11 +832,13 @@ async function readEnvFiles( export function serviceToDeployment( project: string, name: string, + compose: ComposeSpecification, service: Service, cwd = process.cwd(), extraEnv: Record = {}, + volumes: ComposeVolumes = {}, ): V1Deployment { - const mounts = toMounts(service, cwd); + const mounts = toMounts(service, cwd, volumes); const ports = toPorts(service); const hasEnvSecret = Boolean(service.env_file) || @@ -541,6 +873,7 @@ export function serviceToDeployment( }, spec: { restartPolicy: "Always", + ...getComposeArchPlacement(compose), volumes: toVolumes(project, mounts), containers: [ deepMerge( @@ -620,10 +953,11 @@ export function volumesToPvc( project: string, service: Service, cwd = process.cwd(), + volumes: ComposeVolumes = {}, ): V1PersistentVolumeClaim[] { const claims = new Map(); - for (const mount of toMounts(service, cwd)) { + for (const mount of toMounts(service, cwd, volumes)) { const claimName = mount.kind === "bind" ? mount.source @@ -647,11 +981,16 @@ export function volumesToPvc( accessModes: ["ReadWriteMany"], resources: { requests: { - storage: "1Gi", + 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()]; @@ -661,10 +1000,11 @@ export function volumesToConfigMaps( project: string, service: Service, cwd = process.cwd(), + volumes: ComposeVolumes = {}, ): V1ConfigMap[] { const configMaps = new Map(); - for (const mount of toMounts(service, cwd)) { + for (const mount of toMounts(service, cwd, volumes)) { if (mount.kind !== "config-file" || !mount.source || !mount.configKey) continue; @@ -751,11 +1091,16 @@ export async function composeToKubernetes( resources.set(getResourceKey(namespace), namespace); for (const [name, service] of Object.entries(compose.services ?? {})) { - for (const pvc of volumesToPvc(project, service, cwd)) { + for (const pvc of volumesToPvc(project, service, cwd, compose.volumes)) { resources.set(getResourceKey(pvc), pvc); } - for (const configMap of volumesToConfigMaps(project, service, cwd)) { + for (const configMap of volumesToConfigMaps( + project, + service, + cwd, + compose.volumes, + )) { resources.set(getResourceKey(configMap), configMap); } @@ -775,14 +1120,19 @@ export async function composeToKubernetes( const deployment = serviceToDeployment( project, name, + compose, service, cwd, serviceEnv[name] ?? {}, + compose.volumes, ); resources.set(getResourceKey(deployment), deployment); const ingress = serviceToIngress(project, name, service); if (ingress) resources.set(getResourceKey(ingress), ingress); + + const ingressRoute = toIngressRoute(project, name, toPorts(service)); + if (ingressRoute) resources.set(getResourceKey(ingressRoute), ingressRoute); } return [...resources.values()]; diff --git a/lib/database.ts b/lib/database.ts index 4c49a1f..ada0ec8 100644 --- a/lib/database.ts +++ b/lib/database.ts @@ -252,19 +252,25 @@ async function reconcileManagedRoles(claims: PostgresClaim[]): Promise { roles.set(claim.username, toManagedRole(claim)); } - await applyResource({ - apiVersion: "postgresql.cnpg.io/v1", - kind: "Cluster", - metadata: { - name: DATABASE_CLUSTER, - namespace: DATABASE_NAMESPACE, - }, - spec: { - managed: { - roles: [...roles.values()], + try { + await applyResource({ + apiVersion: "postgresql.cnpg.io/v1", + kind: "Cluster", + metadata: { + name: DATABASE_CLUSTER, + namespace: DATABASE_NAMESPACE, }, - }, - }); + spec: { + managed: { + roles: [...roles.values()], + }, + }, + }); + } catch (error) { + throw new Error( + `Failed to reconcile managed roles on ${DATABASE_NAMESPACE}/${DATABASE_CLUSTER}: ${error instanceof Error ? error.message : String(error)}`, + ); + } } async function reconcileDatabases( @@ -277,28 +283,34 @@ async function reconcileDatabases( } for (const claim of uniqueDatabases.values()) { - await applyResource({ - apiVersion: "postgresql.cnpg.io/v1", - kind: "Database", - metadata: { - name: claim.database, - namespace: DATABASE_NAMESPACE, - labels: { - ...LABELS, - [DATABASE_PROJECT_LABEL]: project, - [DATABASE_SERVICE_LABEL]: claim.service, + try { + await applyResource({ + apiVersion: "postgresql.cnpg.io/v1", + kind: "Database", + metadata: { + name: claim.database, + namespace: DATABASE_NAMESPACE, + labels: { + ...LABELS, + [DATABASE_PROJECT_LABEL]: project, + [DATABASE_SERVICE_LABEL]: claim.service, + }, }, - }, - spec: { - cluster: { - name: DATABASE_CLUSTER, + spec: { + cluster: { + name: DATABASE_CLUSTER, + }, + databaseReclaimPolicy: "retain", + ensure: "present", + name: claim.database, + owner: claim.username, }, - databaseReclaimPolicy: "retain", - ensure: "present", - name: claim.database, - owner: claim.username, - }, - }); + }); + } catch (error) { + throw new Error( + `Failed to reconcile database ${claim.database} owned by ${claim.username}: ${error instanceof Error ? error.message : String(error)}`, + ); + } } } diff --git a/lib/render.ts b/lib/render.ts new file mode 100644 index 0000000..2a73583 --- /dev/null +++ b/lib/render.ts @@ -0,0 +1,12 @@ +import type { ComposeSpecification } from "../schema/docker.d"; +import { composeToKubernetes, type KubernetesResource } from "./convert"; +import { reconcilePostgresClaims } from "./database"; + +export async function renderResources( + project: string, + compose: ComposeSpecification, + cwd = process.cwd(), +): Promise { + const serviceEnv = await reconcilePostgresClaims(project, compose); + return composeToKubernetes(project, compose, cwd, serviceEnv); +}