import { basename, dirname, isAbsolute, relative, resolve } from "node:path"; import { cp, mkdir, mkdtemp, rm } from "node:fs/promises"; 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_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; context: string; dockerfile?: string; target?: string; buildArgs: string[]; }; type ProgressReporter = (message: string) => void | Promise; type BuildReporter = { progress?: ProgressReporter; stream?: Writable; }; type SpawnResult = { exitCode: number; stdout: string; 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 { return typeof stream === "number" ? undefined : stream; } function shellQuote(value: string): string { return `'${value.replaceAll("'", `'"'"'`)}'`; } function resolveBuildArgs(service: Service): string[] { if ( !service.build || typeof service.build === "string" || !service.build.args ) { return []; } if (Array.isArray(service.build.args)) return [...service.build.args]; return Object.entries(service.build.args) .filter(([, value]) => value !== null) .map(([key, value]) => `${key}=${String(value)}`); } function resolveBuildPlan( project: string, name: string, service: Service, cwd: string, repoRoot: string, buildRoot: string, ): BuildPlan | undefined { if (!service.build) return; const build = service.build; const contextInput = typeof build === "string" ? build : (build.context ?? "."); if (contextInput.includes("://")) { 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}`, ); } if (typeof build !== "string" && build.dockerfile_inline) { throw new Error(`dockerfile_inline is not supported for service ${name}`); } const dockerfilePath = typeof build === "string" || !build.dockerfile ? undefined : resolve(contextPath, build.dockerfile); const dockerfileRelative = dockerfilePath ? relative(repoRoot, dockerfilePath) : undefined; 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, 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() === getRemoteBuilder().split("@").at(-1); } async function pumpStream( stream: ReadableStream | null | undefined, onLine: (line: string) => void | Promise, ) { if (!stream) return; const reader = stream.getReader(); const decoder = new TextDecoder(); let buffer = ""; try { while (true) { const { done, value } = await reader.read(); if (done) break; buffer += decoder.decode(value, { stream: true }); let newline = buffer.indexOf("\n"); while (newline !== -1) { const line = buffer.slice(0, newline).replace(/\r$/, ""); buffer = buffer.slice(newline + 1); if (line) await onLine(line); newline = buffer.indexOf("\n"); } } buffer += decoder.decode(); const line = buffer.replace(/\r$/, ""); if (line) await onLine(line); } finally { reader.releaseLock(); } } async function runWithOutput( command: Bun.Subprocess, reporter?: BuildReporter, ) { 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) queueStreamLine(line); else if (reporter?.progress) await reporter.progress(line); }), pumpStream(toReadableStream(command.stderr), async (line) => { stderrLines.push(line); 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( summarizeCommandFailure(exitCode, stderrLines, stdoutLines), ); } return { exitCode, stdout: stdoutLines.join("\n"), stderr: stderrLines.join("\n"), } satisfies SpawnResult; } function ssh(script: string) { const remoteCommand = `bash -lc ${shellQuote(script)}`; return Bun.spawn(["ssh", getRemoteBuilder(), remoteCommand], { stdout: "pipe", stderr: "pipe", }); } function scp(localPath: string, remotePath: string) { 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()); } async function getHeadSha(repoRoot: string): Promise { return Bun.$.cwd(repoRoot)`git rev-parse HEAD`.text().then((e) => e.trim()); } async function getTrackedDiff(repoRoot: string): Promise { return Bun.$.cwd(repoRoot)`git diff --binary HEAD`.text(); } async function getUntrackedFiles(repoRoot: string): Promise { const result = Bun.spawn( ["git", "ls-files", "--others", "--exclude-standard"], { cwd: repoRoot, stdout: "pipe", stderr: "pipe", }, ); const output = await runWithOutput(result); return output.stdout .split("\n") .map((line) => line.trim()) .filter(Boolean); } async function getIgnoredDotenvFiles(repoRoot: string): Promise { const result = Bun.spawn( [ "git", "ls-files", "--others", "--ignored", "--exclude-standard", "--", ".env*", "**/.env*", ], { cwd: repoRoot, stdout: "pipe", stderr: "pipe", }, ); const output = await runWithOutput(result); return output.stdout .split("\n") .map((line) => line.trim()) .filter(Boolean); } async function getRemoteHead(remoteRepo: string): Promise { const result = ssh( `if [ -d ${shellQuote(`${remoteRepo}/.git`)} ]; then git -C ${shellQuote(remoteRepo)} rev-parse HEAD; fi`, ); const output = await runWithOutput(result); const head = output.stdout.trim(); return head || undefined; } async function syncRemoteRepo( repoRoot: string, remoteRepo: string, reporter?: BuildReporter, ) { const headSha = await getHeadSha(repoRoot); const remoteHead = await getRemoteHead(remoteRepo); const tempDir = await mkdtemp(`${tmpdir()}/kuber-build-`); try { if (remoteHead !== headSha) { await reporter?.progress?.("Transferring latest git bundle"); const bundlePath = `${tempDir}/repo.bundle`; 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(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, ); } else { 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); if (diff.trim().length > 0) { await reporter?.progress?.("Applying local git diff on remote builder"); const diffPath = `${tempDir}/repo.diff`; 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, ); } const untrackedFiles = await getUntrackedFiles(repoRoot); if (untrackedFiles.length > 0) { await reporter?.progress?.("Syncing untracked files to remote builder"); const untrackedDir = `${tempDir}/untracked`; for (const file of untrackedFiles) { const source = resolve(repoRoot, file); const target = resolve(untrackedDir, file); await mkdir(dirname(target), { recursive: true }); await cp(source, target, { force: true, recursive: true }); } 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, ); } const ignoredDotenvFiles = await getIgnoredDotenvFiles(repoRoot); await runWithOutput( ssh( [ "set -euo pipefail", `git -C ${shellQuote(remoteRepo)} clean -fdX -- .env* '**/.env*'`, ].join("; "), ), reporter, ); if (ignoredDotenvFiles.length === 0) return; await reporter?.progress?.("Syncing ignored .env* files to remote builder"); const dotenvDir = `${tempDir}/dotenv`; for (const file of ignoredDotenvFiles) { const source = resolve(repoRoot, file); const target = resolve(dotenvDir, file); await mkdir(dirname(target), { recursive: true }); await cp(source, target, { force: true }); } const dotenvArchive = `${tempDir}/dotenv.tar`; 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, ); } finally { await rm(tempDir, { recursive: true, force: true }); } } function getBuildPlans( project: string, compose: ComposeSpecification, cwd: string, repoRoot: string, remoteRepo: string, ): BuildPlan[] { return Object.entries(compose.services ?? {}).flatMap(([name, service]) => { const plan = resolveBuildPlan( project, name, service, cwd, repoRoot, remoteRepo, ); return plan ? [plan] : []; }); } async function buildRemote(plan: BuildPlan, reporter?: BuildReporter) { await reporter?.progress?.(`Building ${plan.name}`); const args = [ "docker", "build", "--push", "--progress=plain", "-t", plan.image, ...(plan.dockerfile ? ["-f", plan.dockerfile] : []), ...(plan.target ? ["--target", plan.target] : []), ...plan.buildArgs.flatMap((arg) => ["--build-arg", arg]), plan.context, ]; const command = args.map(shellQuote).join(" "); await runWithOutput(ssh(`set -euo pipefail; ${command}`), reporter); } async function buildLocal(plan: BuildPlan, reporter?: BuildReporter) { await reporter?.progress?.(`Building ${plan.name}`); const command = Bun.spawn( [ "docker", "build", "--push", "--progress=plain", "-t", plan.image, ...(plan.dockerfile ? ["-f", plan.dockerfile] : []), ...(plan.target ? ["--target", plan.target] : []), ...plan.buildArgs.flatMap((arg) => ["--build-arg", arg]), plan.context, ], { stdout: "pipe", stderr: "pipe", }, ); await runWithOutput(command, reporter); } export async function buildServices( project: string, compose: ComposeSpecification, cwd = process.cwd(), reporter?: BuildReporter, ): Promise { return withComposeArch(compose, async () => { if (!Object.values(compose.services ?? {}).some((service) => service.build)) { 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 (plans.length === 0) return 0; if (!localBuilder) { await syncRemoteRepo(repoRoot, buildRoot, reporter); } for (const plan of plans) { if (localBuilder) await buildLocal(plan, reporter); else await buildRemote(plan, reporter); } return plans.length; }); }