import { existsSync, mkdirSync, writeFileSync } from "node:fs"; import path from "node:path"; import { parse as parseToml } from "smol-toml"; import config from "../config.ts"; import { paths } from "../constants.ts"; import { db } from "../db/index.ts"; import { repoPath } from "./git.ts"; // --- Types --- interface CiVariableDef { default?: string; description?: string; } interface CiStepConfig { run_sh?: string; run_if?: string; 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; } 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?: string[]; 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; } // Reserved TOML table names that are not steps const RESERVED_TABLES = new Set(["on", "variables"]); // 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 = parseToml(tomlStr) as Record; } catch { return null; } const image = raw.image; if (typeof image !== "string" || !image) return null; const steps: CiStep[] = []; for (const [key, val] of Object.entries(raw)) { if (RESERVED_TABLES.has(key)) continue; if (typeof val !== "object" || val === null || Array.isArray(val)) continue; // It's a table section — treat as a step const stepCfg = val as Record; steps.push({ name: key, ...(stepCfg as CiStepConfig) }); } 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: Array.isArray(raw.cache) ? (raw.cache as string[]) : 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, }; } // --- 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" }, ); // Consume body to completion await resp.body?.cancel(); } 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); } } async function ensureVolume(volName: string, repoName: 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 }, }), }); await resp.body?.cancel(); } async function createContainer( runId: number, repoName: string, cfg: CiConfig, repoAbsPath: string, envVars: string[], ): Promise { const binds: string[] = [`${repoAbsPath}:/hearthforge-repo.git:ro`]; if (cfg.cache) { for (const cachePath of cfg.cache) { const volName = `hearthforge-ci-cache-${Buffer.from(`${repoName}:${cachePath}`).toString("base64url").slice(0, 24)}`; await ensureVolume(volName, repoName); 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 = ""; 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; log += text; if (Date.now() - lastSave >= 2000) { await onPartialLog(log); lastSave = Date.now(); } } const { text } = parseMuxFrames(buf); log += text; } else { const bodyBytes = new Uint8Array(await startResp.arrayBuffer()); const { text } = parseMuxFrames(bodyBytes); log = 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 }; } 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; 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"]) .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"); // 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, ); // Create step rows in DB for (const step of cfg.steps) { await db .insertInto("ci_steps") .values({ run_id: runId, name: step.name, status: "pending", }) .execute(); } // Pull image await pullImage(cfg.image); if (signal.aborted) throw new Error("Cancelled"); // Create + start container containerId = await createContainer( runId, repo.name, cfg, repoPath(repo.name), 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]); } // Clone project if requested if (cfg.clone_project_to && run.commit_sha) { await execInContainer( containerId, [ "sh", "-c", `git clone /hearthforge-repo.git ${cfg.clone_project_to} && git -C ${cfg.clone_project_to} checkout --detach ${run.commit_sha}`, ], cfg.work_dir, envArray, ); } // Execute steps const shell = cfg.shell ?? ["/bin/sh", "-c"]; let runFailed = false; for (const step of cfg.steps) { if (signal.aborted) { runFailed = true; break; } const stepRow = await db .selectFrom("ci_steps") .select("id") .where("run_id", "=", runId) .where("name", "=", step.name) .executeTakeFirst(); if (!stepRow) continue; const stepId = stepRow.id; // 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 if (step.clear && cfg.clone_project_to && run.commit_sha) { await execInContainer( containerId, [ "sh", "-c", `git -C ${cfg.clone_project_to} reset --hard ${run.commit_sha} && git -C ${cfg.clone_project_to} clean -fdx`, ], cfg.work_dir, envArray, ); } await db .updateTable("ci_steps") .set({ status: "running", started_at: now() }) .where("id", "=", stepId) .execute(); let stepLog = ""; let stepStatus: "success" | "failure" = "success"; 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); try { const { log, exitCode } = await execInContainer( containerId, [...shell, command], cfg.work_dir, envArray, timeoutSignal, async (partial) => { await db .updateTable("ci_steps") .set({ log: maskSecrets(partial, secretValues), }) .where("id", "=", stepId) .execute(); }, ); stepLog = maskSecrets(log, secretValues); if (exitCode !== 0) { stepStatus = "failure"; runFailed = true; } } catch (err) { stepLog = `Step failed: ${err instanceof Error ? err.message : String(err)}\n`; stepStatus = "failure"; runFailed = true; } } // Collect artifacts for this step if (!runFailed || stepStatus === "success") { 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(); if (runFailed) break; } // 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" : "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) { // Write error to a synthetic step if we have no steps yet const hasSteps = await db .selectFrom("ci_steps") .select("id") .where("run_id", "=", runId) .executeTakeFirst(); if (!hasSteps) { await db .insertInto("ci_steps") .values({ run_id: runId, name: "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(() => {}); // 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(); const countRow = await db .selectFrom("ci_runs") .select(db.fn.countAll().as("c")) .where("repo_id", "=", repo.id) .executeTakeFirstOrThrow(); await db .updateTable("ci_runs") .set({ repo_run_id: Number(countRow.c) }) .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; } /** Check if CI can connect to the container socket. */ export async function ciAvailable(): Promise { try { const socket = await getSocket(); const resp = await fetch("http://localhost/v1.47/info", { unix: socket, }); return resp.ok; } catch { return false; } }