import { createHash } from "node:crypto"; import { existsSync, mkdirSync, writeFileSync } from "node:fs"; import path from "node:path"; import config from "../config.ts"; import { CI_MAX_LOG_BYTES, paths } from "../constants.ts"; import { db } from "../db/index.ts"; import { gitEnv, repoPath } from "./git.ts"; // --- Types --- interface CiVariableDef { default?: string; description?: string; } interface CiStepConfig { run_sh?: string; run_if?: string; /** Run even after an earlier step failed. */ always?: boolean; /** A non-zero exit marks the step "warning" and the run carries on. */ warn_on_fail?: boolean; clear?: boolean; timeout?: number; publish_file?: string | string[]; publish_tar?: string | string[]; publish_gzip?: string | string[]; publish_zip?: string | string[]; publish_zstd?: string | string[]; } export interface CiStep extends CiStepConfig { name: string; } /** * One cache path, normalised from either form `cache` accepts: a bare string, * or a table carrying a size cap. */ export interface CiCache { path: string; /** Bytes. Absent means unbounded. */ maxSize?: number; } /** One file or directory lifted out of another image, like COPY --from. */ export interface CiCopy { image: string; from: string; to: string; } export interface CiConfig { image: string; work_dir?: string; clone_project_to?: string; shell?: string[]; shell_setup?: string; timeout?: number; cpu_limit?: number; memory_limit?: string; cache?: CiCache[]; copy?: CiCopy[]; on?: { push?: string[] | boolean; tag?: boolean; manual?: boolean; }; variables?: Record; steps: CiStep[]; } export interface TriggerOpts { triggerSource: "push" | "tag" | "manual"; commitSha: string; commitBranch?: string; commitTag?: string; triggeredBy?: number; variableOverrides?: Record; } // In-memory map of running tasks for cancellation const runningTasks = new Map< number, { controller: AbortController; containerId?: string } >(); // FIFO queue of run IDs whose DB row is `status = "queued"`. We start the // next one in `pumpQueue` whenever `runningTasks.size` drops below // `CI_MAX_CONCURRENT`. `pumpQueue` is the SOLE writer of `runningTasks.set` // — callers that want to start a run push to `queuedRunIds` and call // `pumpQueue` synchronously. This makes the size check + slot reservation // atomic against concurrent callers (no await between). const queuedRunIds: number[] = []; async function promoteToPending(runId: number): Promise { await db .updateTable("ci_runs") .set({ status: "pending" }) .where("id", "=", runId) .execute(); } function pumpQueue(): void { while ( queuedRunIds.length > 0 && runningTasks.size < config.CI_MAX_CONCURRENT ) { const next = queuedRunIds.shift()!; const controller = new AbortController(); runningTasks.set(next, { controller }); // Promote queued→pending before spawning so executeRun's failure path // (which only matches pending/running) can still mark it failed if // it throws very early. We await the flip inside spawnRun's wrapper // so the order is: row=pending → executeRun starts → row=running. spawnRun(next, controller.signal, /* needsPromote */ true); } } /** Position (1-based) of this queued run within the queue, or null. */ export function ciQueuePosition(runId: number): number | null { const idx = queuedRunIds.indexOf(runId); return idx < 0 ? null : idx + 1; } // --- TOML Parsing --- export function parseCiConfig(tomlStr: string): CiConfig | null { let raw: Record; try { raw = Bun.TOML.parse(tomlStr) as Record; } catch { return null; } const image = raw.image; if (typeof image !== "string" || !image) return null; // Not a table per step: JavaScript enumerates integer-like keys first, so // a step named "2024" would jump to the front. An array keeps file order. if (!Array.isArray(raw.steps)) return null; const steps = raw.steps as CiStep[]; // A string or array element fails this the same way a missing name does. if (steps.some((s) => typeof s?.name !== "string" || !s.name)) return null; // `cache` takes a bare path or a table with a cap, so the common entry // stays one line. Normalised here, so nothing downstream sees the union. if (raw.cache !== undefined && !Array.isArray(raw.cache)) return null; const cache: CiCache[] = []; for (const entry of raw.cache ?? []) { if (typeof entry === "string") { if (!entry) return null; cache.push({ path: entry }); continue; } if (typeof entry !== "object" || entry === null) return null; const c = entry as Record; if (typeof c.path !== "string" || !c.path) return null; if (c.max_size === undefined) { cache.push({ path: c.path }); continue; } if (typeof c.max_size !== "string") return null; const maxSize = parseMemoryBytes(c.max_size); // parseMemoryBytes answers 0 for anything it cannot read. Left alone // that is a limit every cache exceeds, so the cache would be wiped // after every run and look broken rather than misconfigured. if (maxSize <= 0) return null; cache.push({ path: c.path, maxSize }); } if (raw.copy !== undefined && !Array.isArray(raw.copy)) return null; const copy = (raw.copy ?? []) as CiCopy[]; if ( copy.some( (c) => typeof c?.image !== "string" || !c.image || typeof c.from !== "string" || !c.from || typeof c.to !== "string" || !c.to, ) ) return null; const rawOn = raw.on as Record | undefined; const rawVars = raw.variables as | Record> | undefined; const variables: Record = {}; if (rawVars) { for (const [name, def] of Object.entries(rawVars)) { if (typeof def === "object" && def !== null) { variables[name] = { default: typeof def.default === "string" ? def.default : undefined, description: typeof def.description === "string" ? def.description : undefined, }; } } } return { image, work_dir: typeof raw.work_dir === "string" ? raw.work_dir : undefined, clone_project_to: typeof raw.clone_project_to === "string" ? raw.clone_project_to : undefined, shell: Array.isArray(raw.shell) ? (raw.shell as string[]) : undefined, shell_setup: typeof raw.shell_setup === "string" ? raw.shell_setup : undefined, timeout: typeof raw.timeout === "number" ? raw.timeout : undefined, cpu_limit: typeof raw.cpu_limit === "number" ? raw.cpu_limit : undefined, memory_limit: typeof raw.memory_limit === "string" ? raw.memory_limit : undefined, cache, copy: copy.length ? copy : undefined, on: rawOn ? { push: Array.isArray(rawOn.push) ? (rawOn.push as string[]) : typeof rawOn.push === "boolean" ? rawOn.push : undefined, tag: typeof rawOn.tag === "boolean" ? rawOn.tag : undefined, manual: typeof rawOn.manual === "boolean" ? rawOn.manual : undefined, } : undefined, variables, steps, }; } /** * Run-start checks kept out of parseCiConfig, which can only report one * generic failure. Returns an error string, or null when the config is usable. */ export function validateCiConfig(cfg: CiConfig): string | null { for (const c of cfg.copy ?? []) { // `to` is a directory and the basename is kept, so a `to` that repeats // the basename means someone expected a rename. Left alone it silently // produces to//, which only shows up as a missing tool // several steps later. if (path.basename(c.to) === path.basename(c.from)) { return `[[copy]] "to" is a directory, so ${c.from} lands at ${path.join(c.to, path.basename(c.from))}. Drop the last path segment from "to".`; } } if (!cfg.clone_project_to) return null; // The archive endpoint resolves `path` against /, while the execs that // create and clear the directory resolve against work_dir. A relative // value would name two different directories. if (!path.isAbsolute(cfg.clone_project_to)) { return `clone_project_to ("${cfg.clone_project_to}") must be an absolute path.`; } const clone = path.resolve(cfg.clone_project_to); for (const entry of cfg.cache ?? []) { const cachePath = path.resolve(entry.path); // Either nesting direction breaks. A cache below the checkout is // overwritten by the extract and deleted by `clear`. A cache above it // carries the previous run's tree back in. The extract then merges two // commits instead of replacing one. if ( cachePath === clone || cachePath.startsWith(`${clone}${path.sep}`) || clone.startsWith(`${cachePath}${path.sep}`) ) { return ( `cache path "${entry.path}" overlaps clone_project_to ("${cfg.clone_project_to}"). ` + "The checkout is extracted over that directory and a `clear` step " + "deletes it. Move the cache outside the clone directory." ); } } return null; } // --- Trigger matching --- export function shouldTriggerPush(cfg: CiConfig, branch: string): boolean { const pushCfg = cfg.on?.push; if (!pushCfg) return false; if (pushCfg === true) return true; if (Array.isArray(pushCfg)) { return pushCfg.some((pattern) => matchGlob(pattern, branch)); } return false; } export function shouldTriggerTag(cfg: CiConfig): boolean { return cfg.on?.tag === true; } function matchGlob(pattern: string, value: string): boolean { if (pattern === "*") return true; const re = new RegExp( `^${pattern.replace(/[.+^${}()|[\]\\]/g, "\\$&").replace(/\*/g, ".*")}$`, ); return re.test(value); } // --- Docker socket --- let resolvedSocket: string | null = null; async function getSocket(): Promise { if (resolvedSocket) return resolvedSocket; const candidates = config.CI_DOCKER_SOCKET ? [config.CI_DOCKER_SOCKET] : (() => { const uid = process.getuid?.(); return [ "/var/run/docker.sock", "/run/podman/podman.sock", ...(uid !== undefined ? [`/run/user/${uid}/podman/podman.sock`] : []), ]; })(); for (const s of candidates) { if (existsSync(s)) { resolvedSocket = s; return s; } } throw new Error( "No Docker/Podman socket found. Set CI_DOCKER_SOCKET env var.", ); } async function dockerFetch( endpoint: string, init?: RequestInit, ): Promise { const socket = await getSocket(); return fetch(`http://localhost/v1.47${endpoint}`, { ...init, unix: socket, }); } // --- Docker helpers --- function splitImageRef(image: string): { name: string; tag: string } { const lastColon = image.lastIndexOf(":"); if (lastColon < 0) return { name: image, tag: "latest" }; const possibleTag = image.slice(lastColon + 1); if (possibleTag.includes("/")) return { name: image, tag: "latest" }; return { name: image.slice(0, lastColon), tag: possibleTag }; } async function pullImage(image: string): Promise { const { name, tag } = splitImageRef(image); const resp = await dockerFetch( `/images/create?fromImage=${encodeURIComponent(name)}&tag=${encodeURIComponent(tag)}`, { method: "POST" }, ); if (!resp.ok) { const detail = (await resp.text().catch(() => "")).trim(); throw new Error( `Failed to pull image ${image}: HTTP ${resp.status}${detail ? ` ${detail}` : ""}`, ); } // The /images/create stream must be read to the end — the pull only // completes when the stream does. Cancelling it (the previous behavior) // aborted the pull, so createContainer could race a not-yet-present image. // Each line is a JSON progress object; a trailing {"error": …} means the // pull failed despite the HTTP 200. const body = await resp.text(); for (const line of body.split("\n")) { const trimmed = line.trim(); if (!trimmed) continue; let obj: { error?: string } | null = null; try { obj = JSON.parse(trimmed); } catch { continue; // non-JSON progress line — ignore } if (obj?.error) { throw new Error(`Failed to pull image ${image}: ${obj.error}`); } } } function formatBytes(n: number): string { const units = ["B", "K", "M", "G", "T"]; let i = 0; let v = n; while (v >= 1024 && i < units.length - 1) { v /= 1024; i++; } return `${i === 0 ? v : v.toFixed(1)}${units[i]}`; } function parseMemoryBytes(s: string): number { const m = s.match(/^(\d+(?:\.\d+)?)\s*([kmgKMG]?)b?$/); if (!m) return 0; const n = parseFloat(m[1] ?? "0"); switch ((m[2] ?? "").toLowerCase()) { case "k": return Math.floor(n * 1024); case "m": return Math.floor(n * 1024 * 1024); case "g": return Math.floor(n * 1024 * 1024 * 1024); default: return Math.floor(n); } } /** * Volume name for one cache path. * * Hashed, not encoded: base64 of the plain string keeps only its first bytes * once truncated, so `/ci/cache/target` and `/ci/build/project/target` used to * name the same volume and silently share it. */ function cacheVolumeName(repoName: string, cachePath: string): string { const digest = createHash("sha256") .update(`${repoName}:${cachePath}`) .digest("base64url"); return `hearthforge-ci-cache-${digest.slice(0, 24)}`; } async function ensureVolume( volName: string, repoName: string, cachePath: string, ): Promise { const resp = await dockerFetch("/volumes/create", { method: "POST", headers: { "Content-Type": "application/json" }, body: JSON.stringify({ Name: volName, Labels: { "com.hearthforge.repo": repoName, // The name is a digest, so without this nothing can say which // path a volume belongs to. "com.hearthforge.cache-path": cachePath, }, }), }); await resp.body?.cancel(); } /** * Delete this repo's cache volumes that the given config no longer mentions. * * Only safe to call for a default-branch run. The config is read per commit, * so pruning against a feature branch would delete the default branch's * volumes and the two would rebuild each other's caches forever. */ async function pruneStaleCaches( repoName: string, cache: CiCache[] | undefined, ): Promise { const keep = new Set( (cache ?? []).map((c) => cacheVolumeName(repoName, c.path)), ); const filters = encodeURIComponent( JSON.stringify({ label: [`com.hearthforge.repo=${repoName}`] }), ); const resp = await dockerFetch(`/volumes?filters=${filters}`); if (!resp.ok) { await resp.body?.cancel(); return; } const data = (await resp.json()) as { Volumes?: Array<{ Name: string }> }; for (const vol of data.Volumes ?? []) { if (keep.has(vol.Name)) continue; // A volume a concurrent run still holds refuses to go. That is fine, // the next default-branch run picks it up. const del = await dockerFetch(`/volumes/${vol.Name}`, { method: "DELETE", }); await del.body?.cancel(); } } /** * Delete cache volumes that outgrew their `max_size`. * * Sizes come from the daemon, not from a `du` in the CI container: the image * need not ship one, and a timed-out run has no container left to exec in. * Returns what was dropped, so the run can say so. * * A cache is only dropped on a size the daemon actually reported. Treating a * missing measurement as "over the limit" would clear the cache on every run * and read as caching being broken rather than misconfigured. */ async function enforceCacheLimits( repoName: string, cache: CiCache[] | undefined, ): Promise> { const capped = new Map( (cache ?? []) .filter((c) => c.maxSize !== undefined) .map((c) => [ cacheVolumeName(repoName, c.path), { path: c.path, maxSize: c.maxSize! }, ]), ); if (capped.size === 0) return []; const resp = await dockerFetch("/system/df"); if (!resp.ok) { await resp.body?.cancel(); return []; } const data = (await resp.json()) as { Volumes?: Array<{ Name: string; UsageData?: { Size: number; RefCount: number }; }>; }; const dropped: Array<{ path: string; size: number; maxSize: number }> = []; for (const vol of data.Volumes ?? []) { const cap = capped.get(vol.Name); if (!cap) continue; const size = vol.UsageData?.Size ?? -1; if (size < 0 || size <= cap.maxSize) continue; // A concurrent run still has it mounted, so the delete would fail. if ((vol.UsageData?.RefCount ?? 0) > 0) continue; const del = await dockerFetch(`/volumes/${vol.Name}`, { method: "DELETE", }); await del.body?.cancel(); if (del.ok) dropped.push({ ...cap, size }); } return dropped; } async function createContainer( runId: number, repoName: string, cfg: CiConfig, envVars: string[], ): Promise { // No bind for the repo: the Docker daemon resolves bind sources on the // host, where HearthForge's own paths do not exist. uploadRepo copies it in. const binds: string[] = []; for (const { path: cachePath } of cfg.cache ?? []) { const volName = cacheVolumeName(repoName, cachePath); await ensureVolume(volName, repoName, cachePath); binds.push(`${volName}:${cachePath}`); } const hostConfig: Record = { Binds: binds }; if (cfg.cpu_limit) { hostConfig.NanoCpus = Math.floor(cfg.cpu_limit * 1e9); } if (cfg.memory_limit) { hostConfig.Memory = parseMemoryBytes(cfg.memory_limit); } const body = JSON.stringify({ Image: cfg.image, Cmd: ["sleep", "infinity"], Env: envVars, WorkingDir: cfg.work_dir ?? "/", HostConfig: hostConfig, }); const resp = await dockerFetch( `/containers/create?name=hearthforge-ci-${runId}`, { method: "POST", headers: { "Content-Type": "application/json" }, body, }, ); if (!resp.ok) { const text = await resp.text(); throw new Error(`Failed to create container: ${resp.status} ${text}`); } const data = (await resp.json()) as { Id: string }; return data.Id; } async function startContainer(containerId: string): Promise { const resp = await dockerFetch(`/containers/${containerId}/start`, { method: "POST", }); if (!resp.ok && resp.status !== 304) { throw new Error(`Failed to start container: ${resp.status}`); } await resp.body?.cancel(); } interface ExecResult { log: string; exitCode: number; } const dec = new TextDecoder(); function parseMuxFrames(buf: Uint8Array): { text: string; remaining: Uint8Array; } { const chunks: string[] = []; let i = 0; while (i + 8 <= buf.length) { const view = new DataView(buf.buffer, buf.byteOffset + i, 8); const size = view.getUint32(4, false); if (i + 8 + size > buf.length) break; chunks.push(dec.decode(buf.slice(i + 8, i + 8 + size))); i += 8 + size; } const remaining = new Uint8Array(buf.length - i); if (i < buf.length) remaining.set(buf.subarray(i)); return { text: chunks.join(""), remaining }; } async function execInContainer( containerId: string, cmd: string[], workDir?: string, envVars?: string[], signal?: AbortSignal, onPartialLog?: (log: string) => Promise, ): Promise { // Create exec const execBody = JSON.stringify({ Cmd: cmd, AttachStdout: true, AttachStderr: true, ...(workDir ? { WorkingDir: workDir } : {}), ...(envVars ? { Env: envVars } : {}), }); const createResp = await dockerFetch(`/containers/${containerId}/exec`, { method: "POST", headers: { "Content-Type": "application/json" }, body: execBody, signal, }); if (!createResp.ok) { const text = await createResp.text(); throw new Error(`Failed to create exec: ${createResp.status} ${text}`); } const createData = (await createResp.json()) as { Id: string }; const execId = createData.Id; // Start exec and stream output const startResp = await dockerFetch(`/exec/${execId}/start`, { method: "POST", headers: { "Content-Type": "application/json" }, body: JSON.stringify({ Detach: false, Tty: false }), signal, }); let log = ""; let truncated = false; // Bound the buffered log: stop appending once we hit the cap (and note it // once) so a runaway step can't exhaust RAM or make each partial-flush // rewrite an ever-growing row. const appendLog = (text: string) => { if (truncated || !text) return; const room = CI_MAX_LOG_BYTES - log.length; if (text.length <= room) { log += text; } else { log += text.slice(0, Math.max(0, room)); log += `\n[log truncated at ${CI_MAX_LOG_BYTES} bytes]\n`; truncated = true; } }; if (onPartialLog && startResp.body) { const reader = startResp.body.getReader(); let buf = new Uint8Array(0); let lastSave = Date.now(); while (true) { const { done, value } = await reader.read(); if (done) break; const merged = new Uint8Array(buf.length + value.length); merged.set(buf); merged.set(value, buf.length); buf = merged; const { text, remaining } = parseMuxFrames(buf); buf = remaining; appendLog(text); // Once truncated the log no longer changes, so stop re-flushing it. if (!truncated && Date.now() - lastSave >= 2000) { await onPartialLog(log); lastSave = Date.now(); } } const { text } = parseMuxFrames(buf); appendLog(text); } else { const bodyBytes = new Uint8Array(await startResp.arrayBuffer()); const { text } = parseMuxFrames(bodyBytes); appendLog(text); } // Get exit code const inspectResp = await dockerFetch(`/exec/${execId}/json`); const inspectData = (await inspectResp.json()) as { ExitCode: number }; return { log, exitCode: inspectData.ExitCode ?? 1 }; } /** * Upload the commit's tree into the container at `destPath`. * * The checkout deliberately does not happen inside the container. That would * force every CI image to carry a git binary. The missing dependency would * only surface as a failed run. * * Not a bind mount either. The Docker daemon resolves bind sources in the * host filesystem. When HearthForge itself runs in a container, it finds * nothing at our DATA_DIR path and silently mounts an empty directory. * * `git archive` writes uid 0 and mode 0644/0755 into the tar headers, and * emits no entry for the archive root. The destination directory therefore * keeps the ownership and mode the container gave it. Piping a work tree * through `tar` instead would stamp the host's uid and 0700 onto it, which * locks out any image whose default user is not HearthForge's uid. * * No `.git` reaches the container. A step that needs history must fetch it. */ async function uploadCheckout( containerId: string, repoAbsPath: string, commitSha: string, destPath: string, signal: AbortSignal, ): Promise { const mk = await execInContainer(containerId, ["mkdir", "-p", destPath]); if (mk.exitCode !== 0) { throw new Error( `Failed to create ${destPath} in the container: ${mk.log.trim()}`, ); } // --end-of-options: commit_sha is unvalidated pkt-line text from the push. const archive = Bun.spawn( [ "git", "-C", repoAbsPath, "archive", "--format=tar", "--end-of-options", commitSha, ], { stdout: "pipe", stderr: "pipe", env: gitEnv, signal }, ); try { const resp = await dockerFetch( `/containers/${containerId}/archive?path=${encodeURIComponent(destPath)}`, { method: "PUT", headers: { "Content-Type": "application/x-tar" }, body: archive.stdout, signal, }, ); // The PUT stopped reading, so git would block writing into a full pipe. if (!resp.ok) archive.kill(); // stderr must be drained in the same turn as the wait. A git that fills // the pipe buffer blocks on the write, and awaiting `exited` first would // then hang the run for good. const [exitCode, err] = await Promise.all([ archive.exited, new Response(archive.stderr).text(), ]); if (!resp.ok) { const detail = (await resp.text().catch(() => "")).trim(); throw new Error( `Failed to upload the checkout: HTTP ${resp.status}${detail ? ` ${detail}` : ""}`, ); } await resp.body?.cancel(); if (exitCode !== 0) { throw new Error( `Failed to read ${commitSha}: git archive exited ${exitCode}${err.trim() ? ` ${err.trim()}` : ""}`, ); } } finally { // Reached on a thrown dockerFetch too, where nothing has reaped git. archive.kill(); await archive.exited; } } /** * Lift a file or directory out of another image, like `COPY --from`. * * The source container is created but never started: the archive endpoint * reads its filesystem either way. * * `to` is a directory and the source basename is preserved. Renaming would * mean rewriting tar headers in flight, which turns a stream into a parser. * A `/.` suffix on `from` copies the contents instead, as `docker cp` does. */ async function copyFromImage(containerId: string, spec: CiCopy): Promise { await pullImage(spec.image); // Unnamed on purpose. The startup sweep deletes by exact name, so a name // here would not get cleaned up, and a retry reuses the run id: one leaked // container would then fail every retry with a name conflict. const createResp = await dockerFetch("/containers/create", { method: "POST", headers: { "Content-Type": "application/json" }, body: JSON.stringify({ Image: spec.image, Cmd: ["true"] }), }); if (!createResp.ok) { const detail = (await createResp.text().catch(() => "")).trim(); throw new Error( `copy: cannot create a container from ${spec.image}: HTTP ${createResp.status}${detail ? ` ${detail}` : ""}`, ); } const sourceId = ((await createResp.json()) as { Id: string }).Id; try { // Docker rejects the upload unless the destination already exists. const mk = await execInContainer(containerId, ["mkdir", "-p", spec.to]); if (mk.exitCode !== 0) { throw new Error(`copy: cannot create ${spec.to}: ${mk.log.trim()}`); } const get = await dockerFetch( `/containers/${sourceId}/archive?path=${encodeURIComponent(spec.from)}`, ); if (!get.ok) { await get.body?.cancel(); throw new Error( `copy: cannot read ${spec.from} from ${spec.image}: HTTP ${get.status}`, ); } const put = await dockerFetch( `/containers/${containerId}/archive?path=${encodeURIComponent(spec.to)}`, { method: "PUT", headers: { "Content-Type": "application/x-tar" }, body: get.body, }, ); if (!put.ok) { await get.body?.cancel(); const detail = (await put.text().catch(() => "")).trim(); throw new Error( `copy: cannot write ${spec.from} to ${spec.to}: HTTP ${put.status}${detail ? ` ${detail}` : ""}`, ); } await put.body?.cancel(); } finally { await removeContainer(sourceId); } } async function removeContainer(containerId: string): Promise { try { const resp = await dockerFetch( `/containers/${containerId}?force=true`, { method: "DELETE" }, ); await resp.body?.cancel(); } catch { // Best-effort cleanup } } // --- Tar extraction --- function extractSingleFileFromTar(data: Uint8Array): Uint8Array | null { if (data.length < 512) return null; const dec = new TextDecoder(); const sizeOctal = dec .decode(data.slice(124, 136)) .replace(/\0/g, "") .trim(); const size = parseInt(sizeOctal, 8); if (Number.isNaN(size) || size < 0) return null; if (data.length < 512 + size) return null; return data.slice(512, 512 + size); } async function copyFileFromContainer( containerId: string, containerPath: string, ): Promise { const resp = await dockerFetch( `/containers/${containerId}/archive?path=${encodeURIComponent(containerPath)}`, ); if (!resp.ok) return null; const tarBytes = new Uint8Array(await resp.arrayBuffer()); return extractSingleFileFromTar(tarBytes); } // --- Secret masking --- function maskSecrets(text: string, secrets: string[]): string { for (const secret of secrets) { if (secret) text = text.split(secret).join("[MASKED]"); } return text; } // --- Env var building --- function buildEnvVars( runId: number, repoName: string, opts: TriggerOpts, cfg: CiConfig, secretValues: Array<{ name: string; value: string }>, ): { envArray: string[]; secretValues: string[] } { const vars: Record = { CI: "true", CI_PIPELINE_ID: String(runId), CI_REPO_NAME: repoName, CI_SERVER_URL: config.BASE_URL, CI_TRIGGER_SOURCE: opts.triggerSource, CI_COMMIT_SHA: opts.commitSha, CI_COMMIT_SHORT_SHA: opts.commitSha.slice(0, 8), CI_COMMIT_BRANCH: opts.commitBranch ?? "", CI_COMMIT_TAG: opts.commitTag ?? "", CI_COMMIT_REF_NAME: opts.commitTag ?? opts.commitBranch ?? "", }; // User-defined variable defaults if (cfg.variables) { for (const [name, def] of Object.entries(cfg.variables)) { if (def.default !== undefined) vars[name] = def.default; } } // Variable overrides from manual trigger if (opts.variableOverrides) { for (const [name, value] of Object.entries(opts.variableOverrides)) { vars[name] = value; } } // Secrets (injected but values tracked for masking) const secretVals: string[] = []; for (const { name, value } of secretValues) { vars[name] = value; secretVals.push(value); } const envArray = Object.entries(vars).map(([k, v]) => `${k}=${v}`); return { envArray, secretValues: secretVals }; } // --- Artifact collection --- async function collectArtifacts( runId: number, containerId: string, step: CiStep, _shell: string[], envVars: string[], ): Promise { const artifactDir = path.join(paths.CI_ARTIFACTS_DIR, String(runId)); mkdirSync(artifactDir, { recursive: true }); const toArray = (v: string | string[] | undefined): string[] => { if (!v) return []; return Array.isArray(v) ? v : [v]; }; // publish_file: copy directly out of container for (const srcPath of toArray(step.publish_file)) { const fileBytes = await copyFileFromContainer(containerId, srcPath); if (fileBytes) { const filename = path.basename(srcPath); const destPath = path.join(artifactDir, filename); writeFileSync(destPath, fileBytes); const stat = Bun.file(destPath); await db .insertInto("ci_artifacts") .values({ run_id: runId, filename, size: stat.size, }) .execute(); } } // Archive commands run with `path.dirname(srcPath)` as the working // directory and reference the source by its basename only, so the // user-controlled path never appears as part of an interpolated shell // string. Previously `publish_zip` used `sh -c "cd … && zip …"` with // raw interpolation — a step author who could write the CI TOML // could shell-inject through the source path. Today the only TOML // author is the admin, but this removes the implicit assumption. type ArchiveType = "tar" | "gzip" | "zip" | "zstd"; const archiveFormats: Array<{ type: ArchiveType; paths: string[]; ext: string; cmd: (basename: string, dst: string) => string[]; }> = [ { type: "tar", paths: toArray(step.publish_tar), ext: ".tar", cmd: (basename, dst) => ["tar", "-cf", dst, basename], }, { type: "gzip", paths: toArray(step.publish_gzip), ext: ".tar.gz", cmd: (basename, dst) => ["tar", "-czf", dst, basename], }, { type: "zstd", paths: toArray(step.publish_zstd), ext: ".tar.zst", cmd: (basename, dst) => ["tar", "--zstd", "-cf", dst, basename], }, { type: "zip", paths: toArray(step.publish_zip), ext: ".zip", cmd: (basename, dst) => ["zip", "-r", dst, basename], }, ]; let archiveIndex = 0; for (const { paths: archivePaths, ext, cmd } of archiveFormats) { for (const srcPath of archivePaths) { archiveIndex++; const tmpPath = `/tmp/hf-artifact-${runId}-${archiveIndex}${ext}`; // Run with the source's parent directory as the working // directory so each archive tool can reference the source // by its basename — no -C, no shell. const execResult = await execInContainer( containerId, cmd(path.basename(srcPath), tmpPath), path.dirname(srcPath), envVars, ).catch(() => null); if (!execResult || execResult.exitCode !== 0) continue; // Copy archive out const fileBytes = await copyFileFromContainer(containerId, tmpPath); if (!fileBytes) continue; const filename = `${path.basename(srcPath)}${ext}`; const destPath = path.join(artifactDir, filename); writeFileSync(destPath, fileBytes); const stat = Bun.file(destPath); await db .insertInto("ci_artifacts") .values({ run_id: runId, filename, size: stat.size, }) .execute(); } } } // --- Main execution --- async function executeRun(runId: number, signal: AbortSignal): Promise { const now = () => new Date().toISOString(); let containerId: string | undefined; // Set once the config is known, and only for a default-branch run. The // finally block cannot see cfg, and pruning off another branch is wrong. let pruneCaches: { repoName: string; cache?: CiCache[] } | null = null; let cacheLimits: { repoName: string; cache?: CiCache[] } | null = null; try { // Mark as running await db .updateTable("ci_runs") .set({ status: "running", started_at: now() }) .where("id", "=", runId) .execute(); // Load run details const run = await db .selectFrom("ci_runs") .selectAll() .where("id", "=", runId) .executeTakeFirst(); if (!run) throw new Error("Run not found"); const repo = await db .selectFrom("repositories") .select(["id", "name", "default_branch"]) .where("id", "=", run.repo_id) .executeTakeFirst(); if (!repo) throw new Error("Repo not found"); // Read .hearthforge-ci.toml at the commit const tomlBuf = await import("./git.ts").then((g) => g.git.show(repo.name, run.commit_sha!, ".hearthforge-ci.toml"), ); if (!tomlBuf) throw new Error(".hearthforge-ci.toml not found at commit"); const cfg = parseCiConfig(tomlBuf.toString("utf-8")); if (!cfg) throw new Error("Failed to parse .hearthforge-ci.toml"); const invalid = validateCiConfig(cfg); if (invalid) throw new Error(`Invalid .hearthforge-ci.toml: ${invalid}`); // Size caps apply on any branch: an oversized volume is oversized // whoever noticed. Dropping stale volumes is default-branch only, // because the config that names them is read per commit. cacheLimits = { repoName: repo.name, cache: cfg.cache }; if (run.commit_branch && run.commit_branch === repo.default_branch) { pruneCaches = { repoName: repo.name, cache: cfg.cache }; } // Load secrets for log masking const secrets = await db .selectFrom("ci_secrets") .select(["name", "value"]) .where("repo_id", "=", repo.id) .execute(); const variableOverrides = run.variable_overrides ? (JSON.parse(run.variable_overrides) as Record) : {}; const { envArray, secretValues } = buildEnvVars( runId, repo.name, { triggerSource: run.trigger_source as TriggerOpts["triggerSource"], commitSha: run.commit_sha ?? "", commitBranch: run.commit_branch ?? undefined, commitTag: run.commit_tag ?? undefined, variableOverrides, }, cfg, secrets, ); // Step names need not be unique, so a name cannot identify a row. // Keep the ids in config order instead. const stepIds: number[] = []; for (const step of cfg.steps) { const row = await db .insertInto("ci_steps") .values({ run_id: runId, name: step.name, status: "pending", }) .returning("id") .executeTakeFirstOrThrow(); stepIds.push(row.id); } // Pull image await pullImage(cfg.image); if (signal.aborted) throw new Error("Cancelled"); // Create + start container containerId = await createContainer(runId, repo.name, cfg, envArray); runningTasks.get(runId)!.containerId = containerId; await startContainer(containerId); if (signal.aborted) throw new Error("Cancelled"); // Create work_dir if (cfg.work_dir) { await execInContainer(containerId, ["mkdir", "-p", cfg.work_dir]); } // Copies run before the clone, so a step can rely on the tools they // bring in, and so they can supply a shell the image lacks. for (const spec of cfg.copy ?? []) { await copyFromImage(containerId, spec); if (signal.aborted) throw new Error("Cancelled"); } if (cfg.clone_project_to && run.commit_sha) { await uploadCheckout( containerId, repoPath(repo.name), run.commit_sha, cfg.clone_project_to, signal, ); } // Execute steps const shell = cfg.shell ?? ["/bin/sh", "-c"]; let runFailed = false; let sawWarning = false; for (const [stepIndex, step] of cfg.steps.entries()) { if (signal.aborted) { runFailed = true; break; } // After a failure, the skipped steps stay pending and are marked // skipped below. if (runFailed && !step.always) continue; // A step timeout removes the container to kill the command, so // there is nothing left to run an `always` step in. if (!containerId) continue; const stepId = stepIds[stepIndex]!; // Check run_if condition if (step.run_if) { const { exitCode } = await execInContainer( containerId, [...shell, step.run_if], cfg.work_dir, envArray, ); if (exitCode !== 0) { await db .updateTable("ci_steps") .set({ status: "skipped", started_at: now(), finished_at: now(), log: "Skipped: condition not met", }) .where("id", "=", stepId) .execute(); continue; } } // Handle clear option // Drop the directory and re-extract. `git clean` would need git // in the image, and this also removes files git never tracked. if (step.clear && cfg.clone_project_to && run.commit_sha) { let clearError: string | null = null; try { // One exec, not two. `rm -rf` can delete the container's // WorkingDir, and every later exec then fails to chdir // before its command starts. This exec chdirs first. const reset = await execInContainer(containerId, [ ...shell, 'rm -rf "$1" && mkdir -p "$1"', "sh", cfg.clone_project_to, ]); if (reset.exitCode !== 0) { throw new Error(reset.log.trim()); } await uploadCheckout( containerId, repoPath(repo.name), run.commit_sha, cfg.clone_project_to, signal, ); } catch (err) { clearError = err instanceof Error ? err.message : String(err); } if (clearError !== null) { // Not thrown: the outer handler only records a message // when no step has failed yet, so after an earlier failure // it would vanish. The step row always survives. await db .updateTable("ci_steps") .set({ status: "failure", started_at: now(), finished_at: now(), log: `Failed to reset ${cfg.clone_project_to}: ${clearError}\n`, }) .where("id", "=", stepId) .execute(); runFailed = true; continue; } } await db .updateTable("ci_steps") .set({ status: "running", started_at: now() }) .where("id", "=", stepId) .execute(); let stepLog = ""; let stepStatus: "success" | "failure" | "warning" = "success"; // Applies to every way a step can fail, a timeout included. const onFail = step.warn_on_fail ? "warning" : "failure"; if (step.run_sh) { const command = cfg.shell_setup ? `${cfg.shell_setup}\n${step.run_sh}` : step.run_sh; const stepTimeout = step.timeout ?? cfg.timeout ?? config.CI_DEFAULT_TIMEOUT; const timeoutSignal = AbortSignal.timeout(stepTimeout * 1000); // Abort the exec stream on either a run cancellation or the // per-step timeout. Docker has no per-exec kill, so on abort we // force-remove the container (below), which kills the command // still running inside it. const stepSignal = AbortSignal.any([signal, timeoutSignal]); try { const { log, exitCode } = await execInContainer( containerId, [...shell, command], cfg.work_dir, envArray, stepSignal, async (partial) => { await db .updateTable("ci_steps") .set({ log: maskSecrets(partial, secretValues), }) .where("id", "=", stepId) .execute(); }, ); stepLog = maskSecrets(log, secretValues); if (exitCode !== 0) stepStatus = onFail; } catch (err) { // A run cancellation is reported as "cancelled" by the outer // catch — don't relabel it as a step failure here. if (signal.aborted) throw err; if (timeoutSignal.aborted) { // Not subject to warn_on_fail. The container is about // to be destroyed, so every later step is skipped no // matter what. A run that cannot continue is a failure. stepStatus = "failure"; stepLog = `Step timed out after ${stepTimeout}s\n`; // Kill the container now so the timed-out command stops // immediately rather than lingering until cleanup. await removeContainer(containerId); containerId = ""; } else { stepStatus = onFail; stepLog = `Step failed: ${err instanceof Error ? err.message : String(err)}\n`; } } } if (stepStatus === "failure") runFailed = true; if (stepStatus === "warning") sawWarning = true; // Collect artifacts for this step if (stepStatus !== "failure") { await collectArtifacts( runId, containerId, step, shell, envArray, ).catch(() => {}); } await db .updateTable("ci_steps") .set({ status: stepStatus, finished_at: now(), log: stepLog, }) .where("id", "=", stepId) .execute(); } // Mark remaining steps as skipped await db .updateTable("ci_steps") .set({ status: "skipped", started_at: now(), finished_at: now(), log: "Skipped: previous step failed", }) .where("run_id", "=", runId) .where("status", "=", "pending") .execute(); const finalStatus = runFailed ? "failure" : sawWarning ? "warning" : "success"; await db .updateTable("ci_runs") .set({ status: finalStatus, finished_at: now() }) .where("id", "=", runId) .execute(); } catch (err) { const errMsg = err instanceof Error ? err.message : String(err); const isDockerUnavailable = errMsg.includes("No Docker/Podman socket"); const status = signal.aborted ? "cancelled" : isDockerUnavailable ? "skipped" : "failure"; const skipLog = signal.aborted ? "Skipped: run was cancelled" : isDockerUnavailable ? "Skipped: Docker/Podman not available" : "Skipped: run failed"; if (!isDockerUnavailable) { // A failed step already carries its own log. Failures outside any // step (validation, upload, clone) would leave no reason at all. const failedStep = await db .selectFrom("ci_steps") .select("id") .where("run_id", "=", runId) .where("status", "=", "failure") .executeTakeFirst(); if (!failedStep) { await db .insertInto("ci_steps") .values({ run_id: runId, // Not "setup": a pipeline may well have its own step // by that name. name: "pipeline setup", status: "failure", started_at: new Date().toISOString(), finished_at: new Date().toISOString(), log: `Error: ${errMsg}\n`, }) .execute(); } } await db .updateTable("ci_runs") .set({ status, finished_at: new Date().toISOString() }) .where("id", "=", runId) .execute(); // Mark pending steps as skipped await db .updateTable("ci_steps") .set({ status: "skipped", finished_at: new Date().toISOString(), log: skipLog, }) .where("run_id", "=", runId) .where("status", "=", "pending") .execute(); } finally { if (containerId) await removeContainer(containerId); runningTasks.delete(runId); // Prune old history const run = await db .selectFrom("ci_runs") .select("repo_id") .where("id", "=", runId) .executeTakeFirst(); if (run) await pruneHistory(run.repo_id).catch(() => {}); if (pruneCaches) { await pruneStaleCaches( pruneCaches.repoName, pruneCaches.cache, ).catch(() => {}); } // After removeContainer above: the daemon will not delete a volume // that is still mounted. if (cacheLimits) { const dropped = await enforceCacheLimits( cacheLimits.repoName, cacheLimits.cache, ).catch(() => []); if (dropped.length) { // Recorded on the run, because the only other symptom is the // next build being mysteriously slow. await db .insertInto("ci_steps") .values({ run_id: runId, name: "cache", status: "success", started_at: new Date().toISOString(), finished_at: new Date().toISOString(), log: `${dropped .map( (d) => `Dropped cache ${d.path}: ${formatBytes(d.size)} over the ${formatBytes(d.maxSize)} limit`, ) .join("\n")}\n`, }) .execute() .catch(() => {}); } } // A slot just freed up — start the next queued run if any. pumpQueue(); } } // --- Public API --- export async function triggerRun( repoName: string, opts: TriggerOpts, ): Promise { const repo = await db .selectFrom("repositories") .select("id") .where("name", "=", repoName) .executeTakeFirst(); if (!repo) throw new Error("Repository not found"); const runId = await db .insertInto("ci_runs") .values({ repo_id: repo.id, triggered_by: opts.triggeredBy ?? null, trigger_source: opts.triggerSource, commit_sha: opts.commitSha, commit_branch: opts.commitBranch ?? null, commit_tag: opts.commitTag ?? null, status: "pending", variable_overrides: opts.variableOverrides ? JSON.stringify(opts.variableOverrides) : null, }) .returning("id") .executeTakeFirstOrThrow(); // Allocate the human-facing run number atomically from a per-repo counter. // A single upsert-and-increment can't collide under concurrent triggers and // never reuses a number after pruneHistory shrinks the run table — both of // which a COUNT(*)-based scheme suffered from. const counter = await db .insertInto("ci_run_counters") .values({ repo_id: repo.id, last_run_id: 1 }) .onConflict((oc) => oc.column("repo_id").doUpdateSet((eb) => ({ last_run_id: eb("ci_run_counters.last_run_id", "+", 1), })), ) .returning("last_run_id") .executeTakeFirstOrThrow(); await db .updateTable("ci_runs") .set({ repo_run_id: counter.last_run_id }) .where("id", "=", runId.id) .execute(); // Flip to "queued" BEFORE pushing onto queuedRunIds. Otherwise a // concurrent pumpQueue (from a finishing run) could observe our entry, // promoteToPending it, and start executeRun while our UPDATE is still // in flight — the late UPDATE would then clobber a running row back to // "queued". Only pumpQueue (the sole writer of runningTasks.set) may // mutate the row's status after the push. The TOCTOU concern is moot // here because pumpQueue cannot pump a runId that hasn't been pushed. if (runningTasks.size >= config.CI_MAX_CONCURRENT) { await db .updateTable("ci_runs") .set({ status: "queued" }) .where("id", "=", runId.id) .execute(); } queuedRunIds.push(runId.id); pumpQueue(); return runId.id; } // Wraps the fire-and-forget executeRun so unhandled exceptions (e.g. a throw // before/after its own try/finally) are logged and the run is reconciled to // "failure" instead of staying pending forever. function spawnRun( runId: number, signal: AbortSignal, needsPromote = false, ): void { void (async () => { try { if (needsPromote) await promoteToPending(runId); await executeRun(runId, signal); } catch (err) { console.error(`[ci] executeRun threw for run ${runId}:`, err); runningTasks.delete(runId); try { await db .updateTable("ci_runs") .set({ status: "failure", finished_at: new Date().toISOString(), }) .where("id", "=", runId) // Include "queued" so a row that was promoted-but-not-yet- // observed (or never promoted because promoteToPending threw) // still gets marked failed instead of stuck. .where("status", "in", ["pending", "running", "queued"]) .execute(); } catch (dbErr) { console.error( `[ci] failed to mark run ${runId} failed:`, dbErr, ); } pumpQueue(); } })(); } export async function retryRun( runId: number, retriedBy: number, ): Promise { const run = await db .selectFrom("ci_runs") .select("repo_id") .where("id", "=", runId) .executeTakeFirst(); if (!run) throw new Error("Run not found"); // Delete existing steps await db.deleteFrom("ci_steps").where("run_id", "=", runId).execute(); // Delete artifacts from disk and DB const artifactDir = path.join(paths.CI_ARTIFACTS_DIR, String(runId)); if (existsSync(artifactDir)) { await Bun.$`rm -rf ${artifactDir}`.quiet().nothrow(); } await db.deleteFrom("ci_artifacts").where("run_id", "=", runId).execute(); // Reset to "pending" first; if a slot isn't free, flip to "queued" // before enqueueing. Mirrors triggerRun: pumpQueue is the sole writer // of runningTasks.set, so we never reserve a slot directly here. Two // concurrent retryRun calls (or a retryRun racing triggerRun) all // funnel through pumpQueue, which serializes the size check against // slot reservation in a single synchronous turn. await db .updateTable("ci_runs") .set({ status: "pending", triggered_by: retriedBy, started_at: null, finished_at: null, }) .where("id", "=", runId) .execute(); if (runningTasks.size >= config.CI_MAX_CONCURRENT) { await db .updateTable("ci_runs") .set({ status: "queued" }) .where("id", "=", runId) .execute(); } queuedRunIds.push(runId); pumpQueue(); } export async function cancelRun(runId: number): Promise { const task = runningTasks.get(runId); if (task) { const { containerId } = task; task.controller.abort(); if (containerId) { await removeContainer(containerId).catch(() => {}); } } const queueIdx = queuedRunIds.indexOf(runId); if (queueIdx >= 0) queuedRunIds.splice(queueIdx, 1); await db .updateTable("ci_runs") .set({ status: "cancelled", finished_at: new Date().toISOString() }) .where("id", "=", runId) .where("status", "in", ["pending", "running", "queued"]) .execute(); } async function pruneHistory(repoId: number): Promise { const maxHistory = config.CI_MAX_HISTORY; const allRuns = await db .selectFrom("ci_runs") .select("id") .where("repo_id", "=", repoId) .orderBy("id", "desc") .execute(); if (allRuns.length <= maxHistory) return; const toDelete = allRuns.slice(maxHistory).map((r) => r.id); for (const runId of toDelete) { // Remove artifacts from disk const artifactDir = path.join(paths.CI_ARTIFACTS_DIR, String(runId)); if (existsSync(artifactDir)) { await Bun.$`rm -rf ${artifactDir}`.quiet().nothrow(); } } await db.deleteFrom("ci_runs").where("id", "in", toDelete).execute(); } export async function purgeRepoCaches(repoName: string): Promise { try { const filters = encodeURIComponent( JSON.stringify({ label: [`com.hearthforge.repo=${repoName}`] }), ); const resp = await dockerFetch(`/volumes?filters=${filters}`); if (!resp.ok) { await resp.body?.cancel(); return; } const data = (await resp.json()) as { Volumes?: Array<{ Name: string }>; }; for (const vol of data.Volumes ?? []) { const delResp = await dockerFetch(`/volumes/${vol.Name}`, { method: "DELETE", }); await delResp.body?.cancel(); } } catch { // Best-effort } } export async function cancelStaleRuns(): Promise { const now = new Date().toISOString(); const stale = await db .selectFrom("ci_runs") .select("id") .where("status", "in", ["pending", "running", "queued"]) .execute(); await Promise.allSettled( stale.map((r) => dockerFetch(`/containers/hearthforge-ci-${r.id}?force=true`, { method: "DELETE", }).then((res) => res.body?.cancel()), ), ); await db .updateTable("ci_runs") .set({ status: "cancelled", finished_at: now }) .where("status", "in", ["pending", "running", "queued"]) .execute(); await db .updateTable("ci_steps") .set({ status: "cancelled", finished_at: now, log: "Skipped: run was cancelled", }) .where("status", "in", ["pending", "running"]) .execute(); } /** Reset the cached socket path (used in tests to switch mock sockets). */ export function resetDockerSocket(): void { resolvedSocket = null; }