ci.ts
⎇
Raw
1import { createHash } from "node:crypto";
2import { existsSync, mkdirSync, writeFileSync } from "node:fs";
3import path from "node:path";
4import config from "../config.ts";
5import { CI_MAX_LOG_BYTES, paths } from "../constants.ts";
6import { db } from "../db/index.ts";
7import { gitEnv, repoPath } from "./git.ts";
8
9// --- Types ---
10
11interface CiVariableDef {
12 default?: string;
13 description?: string;
14}
15
16interface CiStepConfig {
17 run_sh?: string;
18 run_if?: string;
19 /** Run even after an earlier step failed. */
20 always?: boolean;
21 /** A non-zero exit marks the step "warning" and the run carries on. */
22 warn_on_fail?: boolean;
23 clear?: boolean;
24 timeout?: number;
25 publish_file?: string | string[];
26 publish_tar?: string | string[];
27 publish_gzip?: string | string[];
28 publish_zip?: string | string[];
29 publish_zstd?: string | string[];
30}
31
32export interface CiStep extends CiStepConfig {
33 name: string;
34}
35
36/**
37 * One cache path, normalised from either form `cache` accepts: a bare string,
38 * or a table carrying a size cap.
39 */
40export interface CiCache {
41 path: string;
42 /** Bytes. Absent means unbounded. */
43 maxSize?: number;
44}
45
46/** One file or directory lifted out of another image, like COPY --from. */
47export interface CiCopy {
48 image: string;
49 from: string;
50 to: string;
51}
52
53export interface CiConfig {
54 image: string;
55 work_dir?: string;
56 clone_project_to?: string;
57 shell?: string[];
58 shell_setup?: string;
59 timeout?: number;
60 cpu_limit?: number;
61 memory_limit?: string;
62 cache?: CiCache[];
63 copy?: CiCopy[];
64 on?: {
65 push?: string[] | boolean;
66 tag?: boolean;
67 manual?: boolean;
68 };
69 variables?: Record<string, CiVariableDef>;
70 steps: CiStep[];
71}
72
73export interface TriggerOpts {
74 triggerSource: "push" | "tag" | "manual";
75 commitSha: string;
76 commitBranch?: string;
77 commitTag?: string;
78 triggeredBy?: number;
79 variableOverrides?: Record<string, string>;
80}
81
82// In-memory map of running tasks for cancellation
83const runningTasks = new Map<
84 number,
85 { controller: AbortController; containerId?: string }
86>();
87
88// FIFO queue of run IDs whose DB row is `status = "queued"`. We start the
89// next one in `pumpQueue` whenever `runningTasks.size` drops below
90// `CI_MAX_CONCURRENT`. `pumpQueue` is the SOLE writer of `runningTasks.set`
91// — callers that want to start a run push to `queuedRunIds` and call
92// `pumpQueue` synchronously. This makes the size check + slot reservation
93// atomic against concurrent callers (no await between).
94const queuedRunIds: number[] = [];
95
96async function promoteToPending(runId: number): Promise<void> {
97 await db
98 .updateTable("ci_runs")
99 .set({ status: "pending" })
100 .where("id", "=", runId)
101 .execute();
102}
103
104function pumpQueue(): void {
105 while (
106 queuedRunIds.length > 0 &&
107 runningTasks.size < config.CI_MAX_CONCURRENT
108 ) {
109 const next = queuedRunIds.shift()!;
110 const controller = new AbortController();
111 runningTasks.set(next, { controller });
112 // Promote queued→pending before spawning so executeRun's failure path
113 // (which only matches pending/running) can still mark it failed if
114 // it throws very early. We await the flip inside spawnRun's wrapper
115 // so the order is: row=pending → executeRun starts → row=running.
116 spawnRun(next, controller.signal, /* needsPromote */ true);
117 }
118}
119
120/** Position (1-based) of this queued run within the queue, or null. */
121export function ciQueuePosition(runId: number): number | null {
122 const idx = queuedRunIds.indexOf(runId);
123 return idx < 0 ? null : idx + 1;
124}
125
126// --- TOML Parsing ---
127
128export function parseCiConfig(tomlStr: string): CiConfig | null {
129 let raw: Record<string, unknown>;
130 try {
131 raw = Bun.TOML.parse(tomlStr) as Record<string, unknown>;
132 } catch {
133 return null;
134 }
135
136 const image = raw.image;
137 if (typeof image !== "string" || !image) return null;
138
139 // Not a table per step: JavaScript enumerates integer-like keys first, so
140 // a step named "2024" would jump to the front. An array keeps file order.
141 if (!Array.isArray(raw.steps)) return null;
142 const steps = raw.steps as CiStep[];
143 // A string or array element fails this the same way a missing name does.
144 if (steps.some((s) => typeof s?.name !== "string" || !s.name)) return null;
145
146 // `cache` takes a bare path or a table with a cap, so the common entry
147 // stays one line. Normalised here, so nothing downstream sees the union.
148 if (raw.cache !== undefined && !Array.isArray(raw.cache)) return null;
149 const cache: CiCache[] = [];
150 for (const entry of raw.cache ?? []) {
151 if (typeof entry === "string") {
152 if (!entry) return null;
153 cache.push({ path: entry });
154 continue;
155 }
156 if (typeof entry !== "object" || entry === null) return null;
157 const c = entry as Record<string, unknown>;
158 if (typeof c.path !== "string" || !c.path) return null;
159 if (c.max_size === undefined) {
160 cache.push({ path: c.path });
161 continue;
162 }
163 if (typeof c.max_size !== "string") return null;
164 const maxSize = parseMemoryBytes(c.max_size);
165 // parseMemoryBytes answers 0 for anything it cannot read. Left alone
166 // that is a limit every cache exceeds, so the cache would be wiped
167 // after every run and look broken rather than misconfigured.
168 if (maxSize <= 0) return null;
169 cache.push({ path: c.path, maxSize });
170 }
171
172 if (raw.copy !== undefined && !Array.isArray(raw.copy)) return null;
173 const copy = (raw.copy ?? []) as CiCopy[];
174 if (
175 copy.some(
176 (c) =>
177 typeof c?.image !== "string" ||
178 !c.image ||
179 typeof c.from !== "string" ||
180 !c.from ||
181 typeof c.to !== "string" ||
182 !c.to,
183 )
184 )
185 return null;
186
187 const rawOn = raw.on as Record<string, unknown> | undefined;
188 const rawVars = raw.variables as
189 | Record<string, Record<string, unknown>>
190 | undefined;
191 const variables: Record<string, CiVariableDef> = {};
192 if (rawVars) {
193 for (const [name, def] of Object.entries(rawVars)) {
194 if (typeof def === "object" && def !== null) {
195 variables[name] = {
196 default:
197 typeof def.default === "string"
198 ? def.default
199 : undefined,
200 description:
201 typeof def.description === "string"
202 ? def.description
203 : undefined,
204 };
205 }
206 }
207 }
208
209 return {
210 image,
211 work_dir: typeof raw.work_dir === "string" ? raw.work_dir : undefined,
212 clone_project_to:
213 typeof raw.clone_project_to === "string"
214 ? raw.clone_project_to
215 : undefined,
216 shell: Array.isArray(raw.shell) ? (raw.shell as string[]) : undefined,
217 shell_setup:
218 typeof raw.shell_setup === "string" ? raw.shell_setup : undefined,
219 timeout: typeof raw.timeout === "number" ? raw.timeout : undefined,
220 cpu_limit:
221 typeof raw.cpu_limit === "number" ? raw.cpu_limit : undefined,
222 memory_limit:
223 typeof raw.memory_limit === "string" ? raw.memory_limit : undefined,
224 cache,
225 copy: copy.length ? copy : undefined,
226 on: rawOn
227 ? {
228 push: Array.isArray(rawOn.push)
229 ? (rawOn.push as string[])
230 : typeof rawOn.push === "boolean"
231 ? rawOn.push
232 : undefined,
233 tag: typeof rawOn.tag === "boolean" ? rawOn.tag : undefined,
234 manual:
235 typeof rawOn.manual === "boolean"
236 ? rawOn.manual
237 : undefined,
238 }
239 : undefined,
240 variables,
241 steps,
242 };
243}
244
245/**
246 * Run-start checks kept out of parseCiConfig, which can only report one
247 * generic failure. Returns an error string, or null when the config is usable.
248 */
249export function validateCiConfig(cfg: CiConfig): string | null {
250 for (const c of cfg.copy ?? []) {
251 // `to` is a directory and the basename is kept, so a `to` that repeats
252 // the basename means someone expected a rename. Left alone it silently
253 // produces to/<name>/<name>, which only shows up as a missing tool
254 // several steps later.
255 if (path.basename(c.to) === path.basename(c.from)) {
256 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".`;
257 }
258 }
259
260 if (!cfg.clone_project_to) return null;
261
262 // The archive endpoint resolves `path` against /, while the execs that
263 // create and clear the directory resolve against work_dir. A relative
264 // value would name two different directories.
265 if (!path.isAbsolute(cfg.clone_project_to)) {
266 return `clone_project_to ("${cfg.clone_project_to}") must be an absolute path.`;
267 }
268
269 const clone = path.resolve(cfg.clone_project_to);
270 for (const entry of cfg.cache ?? []) {
271 const cachePath = path.resolve(entry.path);
272 // Either nesting direction breaks. A cache below the checkout is
273 // overwritten by the extract and deleted by `clear`. A cache above it
274 // carries the previous run's tree back in. The extract then merges two
275 // commits instead of replacing one.
276 if (
277 cachePath === clone ||
278 cachePath.startsWith(`${clone}${path.sep}`) ||
279 clone.startsWith(`${cachePath}${path.sep}`)
280 ) {
281 return (
282 `cache path "${entry.path}" overlaps clone_project_to ("${cfg.clone_project_to}"). ` +
283 "The checkout is extracted over that directory and a `clear` step " +
284 "deletes it. Move the cache outside the clone directory."
285 );
286 }
287 }
288 return null;
289}
290
291// --- Trigger matching ---
292
293export function shouldTriggerPush(cfg: CiConfig, branch: string): boolean {
294 const pushCfg = cfg.on?.push;
295 if (!pushCfg) return false;
296 if (pushCfg === true) return true;
297 if (Array.isArray(pushCfg)) {
298 return pushCfg.some((pattern) => matchGlob(pattern, branch));
299 }
300 return false;
301}
302
303export function shouldTriggerTag(cfg: CiConfig): boolean {
304 return cfg.on?.tag === true;
305}
306
307function matchGlob(pattern: string, value: string): boolean {
308 if (pattern === "*") return true;
309 const re = new RegExp(
310 `^${pattern.replace(/[.+^${}()|[\]\\]/g, "\\$&").replace(/\*/g, ".*")}$`,
311 );
312 return re.test(value);
313}
314
315// --- Docker socket ---
316
317let resolvedSocket: string | null = null;
318
319async function getSocket(): Promise<string> {
320 if (resolvedSocket) return resolvedSocket;
321 const candidates = config.CI_DOCKER_SOCKET
322 ? [config.CI_DOCKER_SOCKET]
323 : (() => {
324 const uid = process.getuid?.();
325 return [
326 "/var/run/docker.sock",
327 "/run/podman/podman.sock",
328 ...(uid !== undefined
329 ? [`/run/user/${uid}/podman/podman.sock`]
330 : []),
331 ];
332 })();
333 for (const s of candidates) {
334 if (existsSync(s)) {
335 resolvedSocket = s;
336 return s;
337 }
338 }
339 throw new Error(
340 "No Docker/Podman socket found. Set CI_DOCKER_SOCKET env var.",
341 );
342}
343
344async function dockerFetch(
345 endpoint: string,
346 init?: RequestInit,
347): Promise<Response> {
348 const socket = await getSocket();
349 return fetch(`http://localhost/v1.47${endpoint}`, {
350 ...init,
351 unix: socket,
352 });
353}
354
355// --- Docker helpers ---
356
357function splitImageRef(image: string): { name: string; tag: string } {
358 const lastColon = image.lastIndexOf(":");
359 if (lastColon < 0) return { name: image, tag: "latest" };
360 const possibleTag = image.slice(lastColon + 1);
361 if (possibleTag.includes("/")) return { name: image, tag: "latest" };
362 return { name: image.slice(0, lastColon), tag: possibleTag };
363}
364
365async function pullImage(image: string): Promise<void> {
366 const { name, tag } = splitImageRef(image);
367 const resp = await dockerFetch(
368 `/images/create?fromImage=${encodeURIComponent(name)}&tag=${encodeURIComponent(tag)}`,
369 { method: "POST" },
370 );
371 if (!resp.ok) {
372 const detail = (await resp.text().catch(() => "")).trim();
373 throw new Error(
374 `Failed to pull image ${image}: HTTP ${resp.status}${detail ? ` ${detail}` : ""}`,
375 );
376 }
377 // The /images/create stream must be read to the end — the pull only
378 // completes when the stream does. Cancelling it (the previous behavior)
379 // aborted the pull, so createContainer could race a not-yet-present image.
380 // Each line is a JSON progress object; a trailing {"error": …} means the
381 // pull failed despite the HTTP 200.
382 const body = await resp.text();
383 for (const line of body.split("\n")) {
384 const trimmed = line.trim();
385 if (!trimmed) continue;
386 let obj: { error?: string } | null = null;
387 try {
388 obj = JSON.parse(trimmed);
389 } catch {
390 continue; // non-JSON progress line — ignore
391 }
392 if (obj?.error) {
393 throw new Error(`Failed to pull image ${image}: ${obj.error}`);
394 }
395 }
396}
397
398function formatBytes(n: number): string {
399 const units = ["B", "K", "M", "G", "T"];
400 let i = 0;
401 let v = n;
402 while (v >= 1024 && i < units.length - 1) {
403 v /= 1024;
404 i++;
405 }
406 return `${i === 0 ? v : v.toFixed(1)}${units[i]}`;
407}
408
409function parseMemoryBytes(s: string): number {
410 const m = s.match(/^(\d+(?:\.\d+)?)\s*([kmgKMG]?)b?$/);
411 if (!m) return 0;
412 const n = parseFloat(m[1] ?? "0");
413 switch ((m[2] ?? "").toLowerCase()) {
414 case "k":
415 return Math.floor(n * 1024);
416 case "m":
417 return Math.floor(n * 1024 * 1024);
418 case "g":
419 return Math.floor(n * 1024 * 1024 * 1024);
420 default:
421 return Math.floor(n);
422 }
423}
424
425/**
426 * Volume name for one cache path.
427 *
428 * Hashed, not encoded: base64 of the plain string keeps only its first bytes
429 * once truncated, so `/ci/cache/target` and `/ci/build/project/target` used to
430 * name the same volume and silently share it.
431 */
432function cacheVolumeName(repoName: string, cachePath: string): string {
433 const digest = createHash("sha256")
434 .update(`${repoName}:${cachePath}`)
435 .digest("base64url");
436 return `hearthforge-ci-cache-${digest.slice(0, 24)}`;
437}
438
439async function ensureVolume(
440 volName: string,
441 repoName: string,
442 cachePath: string,
443): Promise<void> {
444 const resp = await dockerFetch("/volumes/create", {
445 method: "POST",
446 headers: { "Content-Type": "application/json" },
447 body: JSON.stringify({
448 Name: volName,
449 Labels: {
450 "com.hearthforge.repo": repoName,
451 // The name is a digest, so without this nothing can say which
452 // path a volume belongs to.
453 "com.hearthforge.cache-path": cachePath,
454 },
455 }),
456 });
457 await resp.body?.cancel();
458}
459
460/**
461 * Delete this repo's cache volumes that the given config no longer mentions.
462 *
463 * Only safe to call for a default-branch run. The config is read per commit,
464 * so pruning against a feature branch would delete the default branch's
465 * volumes and the two would rebuild each other's caches forever.
466 */
467async function pruneStaleCaches(
468 repoName: string,
469 cache: CiCache[] | undefined,
470): Promise<void> {
471 const keep = new Set(
472 (cache ?? []).map((c) => cacheVolumeName(repoName, c.path)),
473 );
474 const filters = encodeURIComponent(
475 JSON.stringify({ label: [`com.hearthforge.repo=${repoName}`] }),
476 );
477 const resp = await dockerFetch(`/volumes?filters=${filters}`);
478 if (!resp.ok) {
479 await resp.body?.cancel();
480 return;
481 }
482 const data = (await resp.json()) as { Volumes?: Array<{ Name: string }> };
483 for (const vol of data.Volumes ?? []) {
484 if (keep.has(vol.Name)) continue;
485 // A volume a concurrent run still holds refuses to go. That is fine,
486 // the next default-branch run picks it up.
487 const del = await dockerFetch(`/volumes/${vol.Name}`, {
488 method: "DELETE",
489 });
490 await del.body?.cancel();
491 }
492}
493
494/**
495 * Delete cache volumes that outgrew their `max_size`.
496 *
497 * Sizes come from the daemon, not from a `du` in the CI container: the image
498 * need not ship one, and a timed-out run has no container left to exec in.
499 * Returns what was dropped, so the run can say so.
500 *
501 * A cache is only dropped on a size the daemon actually reported. Treating a
502 * missing measurement as "over the limit" would clear the cache on every run
503 * and read as caching being broken rather than misconfigured.
504 */
505async function enforceCacheLimits(
506 repoName: string,
507 cache: CiCache[] | undefined,
508): Promise<Array<{ path: string; size: number; maxSize: number }>> {
509 const capped = new Map(
510 (cache ?? [])
511 .filter((c) => c.maxSize !== undefined)
512 .map((c) => [
513 cacheVolumeName(repoName, c.path),
514 { path: c.path, maxSize: c.maxSize! },
515 ]),
516 );
517 if (capped.size === 0) return [];
518
519 const resp = await dockerFetch("/system/df");
520 if (!resp.ok) {
521 await resp.body?.cancel();
522 return [];
523 }
524 const data = (await resp.json()) as {
525 Volumes?: Array<{
526 Name: string;
527 UsageData?: { Size: number; RefCount: number };
528 }>;
529 };
530
531 const dropped: Array<{ path: string; size: number; maxSize: number }> = [];
532 for (const vol of data.Volumes ?? []) {
533 const cap = capped.get(vol.Name);
534 if (!cap) continue;
535 const size = vol.UsageData?.Size ?? -1;
536 if (size < 0 || size <= cap.maxSize) continue;
537 // A concurrent run still has it mounted, so the delete would fail.
538 if ((vol.UsageData?.RefCount ?? 0) > 0) continue;
539
540 const del = await dockerFetch(`/volumes/${vol.Name}`, {
541 method: "DELETE",
542 });
543 await del.body?.cancel();
544 if (del.ok) dropped.push({ ...cap, size });
545 }
546 return dropped;
547}
548
549async function createContainer(
550 runId: number,
551 repoName: string,
552 cfg: CiConfig,
553 envVars: string[],
554): Promise<string> {
555 // No bind for the repo: the Docker daemon resolves bind sources on the
556 // host, where HearthForge's own paths do not exist. uploadRepo copies it in.
557 const binds: string[] = [];
558 for (const { path: cachePath } of cfg.cache ?? []) {
559 const volName = cacheVolumeName(repoName, cachePath);
560 await ensureVolume(volName, repoName, cachePath);
561 binds.push(`${volName}:${cachePath}`);
562 }
563
564 const hostConfig: Record<string, unknown> = { Binds: binds };
565 if (cfg.cpu_limit) {
566 hostConfig.NanoCpus = Math.floor(cfg.cpu_limit * 1e9);
567 }
568 if (cfg.memory_limit) {
569 hostConfig.Memory = parseMemoryBytes(cfg.memory_limit);
570 }
571
572 const body = JSON.stringify({
573 Image: cfg.image,
574 Cmd: ["sleep", "infinity"],
575 Env: envVars,
576 WorkingDir: cfg.work_dir ?? "/",
577 HostConfig: hostConfig,
578 });
579
580 const resp = await dockerFetch(
581 `/containers/create?name=hearthforge-ci-${runId}`,
582 {
583 method: "POST",
584 headers: { "Content-Type": "application/json" },
585 body,
586 },
587 );
588 if (!resp.ok) {
589 const text = await resp.text();
590 throw new Error(`Failed to create container: ${resp.status} ${text}`);
591 }
592 const data = (await resp.json()) as { Id: string };
593 return data.Id;
594}
595
596async function startContainer(containerId: string): Promise<void> {
597 const resp = await dockerFetch(`/containers/${containerId}/start`, {
598 method: "POST",
599 });
600 if (!resp.ok && resp.status !== 304) {
601 throw new Error(`Failed to start container: ${resp.status}`);
602 }
603 await resp.body?.cancel();
604}
605
606interface ExecResult {
607 log: string;
608 exitCode: number;
609}
610
611const dec = new TextDecoder();
612
613function parseMuxFrames(buf: Uint8Array): {
614 text: string;
615 remaining: Uint8Array<ArrayBuffer>;
616} {
617 const chunks: string[] = [];
618 let i = 0;
619 while (i + 8 <= buf.length) {
620 const view = new DataView(buf.buffer, buf.byteOffset + i, 8);
621 const size = view.getUint32(4, false);
622 if (i + 8 + size > buf.length) break;
623 chunks.push(dec.decode(buf.slice(i + 8, i + 8 + size)));
624 i += 8 + size;
625 }
626 const remaining = new Uint8Array(buf.length - i);
627 if (i < buf.length) remaining.set(buf.subarray(i));
628 return { text: chunks.join(""), remaining };
629}
630
631async function execInContainer(
632 containerId: string,
633 cmd: string[],
634 workDir?: string,
635 envVars?: string[],
636 signal?: AbortSignal,
637 onPartialLog?: (log: string) => Promise<void>,
638): Promise<ExecResult> {
639 // Create exec
640 const execBody = JSON.stringify({
641 Cmd: cmd,
642 AttachStdout: true,
643 AttachStderr: true,
644 ...(workDir ? { WorkingDir: workDir } : {}),
645 ...(envVars ? { Env: envVars } : {}),
646 });
647 const createResp = await dockerFetch(`/containers/${containerId}/exec`, {
648 method: "POST",
649 headers: { "Content-Type": "application/json" },
650 body: execBody,
651 signal,
652 });
653 if (!createResp.ok) {
654 const text = await createResp.text();
655 throw new Error(`Failed to create exec: ${createResp.status} ${text}`);
656 }
657 const createData = (await createResp.json()) as { Id: string };
658 const execId = createData.Id;
659
660 // Start exec and stream output
661 const startResp = await dockerFetch(`/exec/${execId}/start`, {
662 method: "POST",
663 headers: { "Content-Type": "application/json" },
664 body: JSON.stringify({ Detach: false, Tty: false }),
665 signal,
666 });
667
668 let log = "";
669 let truncated = false;
670 // Bound the buffered log: stop appending once we hit the cap (and note it
671 // once) so a runaway step can't exhaust RAM or make each partial-flush
672 // rewrite an ever-growing row.
673 const appendLog = (text: string) => {
674 if (truncated || !text) return;
675 const room = CI_MAX_LOG_BYTES - log.length;
676 if (text.length <= room) {
677 log += text;
678 } else {
679 log += text.slice(0, Math.max(0, room));
680 log += `\n[log truncated at ${CI_MAX_LOG_BYTES} bytes]\n`;
681 truncated = true;
682 }
683 };
684
685 if (onPartialLog && startResp.body) {
686 const reader = startResp.body.getReader();
687 let buf = new Uint8Array(0);
688 let lastSave = Date.now();
689 while (true) {
690 const { done, value } = await reader.read();
691 if (done) break;
692 const merged = new Uint8Array(buf.length + value.length);
693 merged.set(buf);
694 merged.set(value, buf.length);
695 buf = merged;
696 const { text, remaining } = parseMuxFrames(buf);
697 buf = remaining;
698 appendLog(text);
699 // Once truncated the log no longer changes, so stop re-flushing it.
700 if (!truncated && Date.now() - lastSave >= 2000) {
701 await onPartialLog(log);
702 lastSave = Date.now();
703 }
704 }
705 const { text } = parseMuxFrames(buf);
706 appendLog(text);
707 } else {
708 const bodyBytes = new Uint8Array(await startResp.arrayBuffer());
709 const { text } = parseMuxFrames(bodyBytes);
710 appendLog(text);
711 }
712
713 // Get exit code
714 const inspectResp = await dockerFetch(`/exec/${execId}/json`);
715 const inspectData = (await inspectResp.json()) as { ExitCode: number };
716
717 return { log, exitCode: inspectData.ExitCode ?? 1 };
718}
719
720/**
721 * Upload the commit's tree into the container at `destPath`.
722 *
723 * The checkout deliberately does not happen inside the container. That would
724 * force every CI image to carry a git binary. The missing dependency would
725 * only surface as a failed run.
726 *
727 * Not a bind mount either. The Docker daemon resolves bind sources in the
728 * host filesystem. When HearthForge itself runs in a container, it finds
729 * nothing at our DATA_DIR path and silently mounts an empty directory.
730 *
731 * `git archive` writes uid 0 and mode 0644/0755 into the tar headers, and
732 * emits no entry for the archive root. The destination directory therefore
733 * keeps the ownership and mode the container gave it. Piping a work tree
734 * through `tar` instead would stamp the host's uid and 0700 onto it, which
735 * locks out any image whose default user is not HearthForge's uid.
736 *
737 * No `.git` reaches the container. A step that needs history must fetch it.
738 */
739async function uploadCheckout(
740 containerId: string,
741 repoAbsPath: string,
742 commitSha: string,
743 destPath: string,
744 signal: AbortSignal,
745): Promise<void> {
746 const mk = await execInContainer(containerId, ["mkdir", "-p", destPath]);
747 if (mk.exitCode !== 0) {
748 throw new Error(
749 `Failed to create ${destPath} in the container: ${mk.log.trim()}`,
750 );
751 }
752
753 // --end-of-options: commit_sha is unvalidated pkt-line text from the push.
754 const archive = Bun.spawn(
755 [
756 "git",
757 "-C",
758 repoAbsPath,
759 "archive",
760 "--format=tar",
761 "--end-of-options",
762 commitSha,
763 ],
764 { stdout: "pipe", stderr: "pipe", env: gitEnv, signal },
765 );
766
767 try {
768 const resp = await dockerFetch(
769 `/containers/${containerId}/archive?path=${encodeURIComponent(destPath)}`,
770 {
771 method: "PUT",
772 headers: { "Content-Type": "application/x-tar" },
773 body: archive.stdout,
774 signal,
775 },
776 );
777
778 // The PUT stopped reading, so git would block writing into a full pipe.
779 if (!resp.ok) archive.kill();
780
781 // stderr must be drained in the same turn as the wait. A git that fills
782 // the pipe buffer blocks on the write, and awaiting `exited` first would
783 // then hang the run for good.
784 const [exitCode, err] = await Promise.all([
785 archive.exited,
786 new Response(archive.stderr).text(),
787 ]);
788
789 if (!resp.ok) {
790 const detail = (await resp.text().catch(() => "")).trim();
791 throw new Error(
792 `Failed to upload the checkout: HTTP ${resp.status}${detail ? ` ${detail}` : ""}`,
793 );
794 }
795 await resp.body?.cancel();
796 if (exitCode !== 0) {
797 throw new Error(
798 `Failed to read ${commitSha}: git archive exited ${exitCode}${err.trim() ? ` ${err.trim()}` : ""}`,
799 );
800 }
801 } finally {
802 // Reached on a thrown dockerFetch too, where nothing has reaped git.
803 archive.kill();
804 await archive.exited;
805 }
806}
807
808/**
809 * Lift a file or directory out of another image, like `COPY --from`.
810 *
811 * The source container is created but never started: the archive endpoint
812 * reads its filesystem either way.
813 *
814 * `to` is a directory and the source basename is preserved. Renaming would
815 * mean rewriting tar headers in flight, which turns a stream into a parser.
816 * A `/.` suffix on `from` copies the contents instead, as `docker cp` does.
817 */
818async function copyFromImage(containerId: string, spec: CiCopy): Promise<void> {
819 await pullImage(spec.image);
820
821 // Unnamed on purpose. The startup sweep deletes by exact name, so a name
822 // here would not get cleaned up, and a retry reuses the run id: one leaked
823 // container would then fail every retry with a name conflict.
824 const createResp = await dockerFetch("/containers/create", {
825 method: "POST",
826 headers: { "Content-Type": "application/json" },
827 body: JSON.stringify({ Image: spec.image, Cmd: ["true"] }),
828 });
829 if (!createResp.ok) {
830 const detail = (await createResp.text().catch(() => "")).trim();
831 throw new Error(
832 `copy: cannot create a container from ${spec.image}: HTTP ${createResp.status}${detail ? ` ${detail}` : ""}`,
833 );
834 }
835 const sourceId = ((await createResp.json()) as { Id: string }).Id;
836
837 try {
838 // Docker rejects the upload unless the destination already exists.
839 const mk = await execInContainer(containerId, ["mkdir", "-p", spec.to]);
840 if (mk.exitCode !== 0) {
841 throw new Error(`copy: cannot create ${spec.to}: ${mk.log.trim()}`);
842 }
843
844 const get = await dockerFetch(
845 `/containers/${sourceId}/archive?path=${encodeURIComponent(spec.from)}`,
846 );
847 if (!get.ok) {
848 await get.body?.cancel();
849 throw new Error(
850 `copy: cannot read ${spec.from} from ${spec.image}: HTTP ${get.status}`,
851 );
852 }
853
854 const put = await dockerFetch(
855 `/containers/${containerId}/archive?path=${encodeURIComponent(spec.to)}`,
856 {
857 method: "PUT",
858 headers: { "Content-Type": "application/x-tar" },
859 body: get.body,
860 },
861 );
862 if (!put.ok) {
863 await get.body?.cancel();
864 const detail = (await put.text().catch(() => "")).trim();
865 throw new Error(
866 `copy: cannot write ${spec.from} to ${spec.to}: HTTP ${put.status}${detail ? ` ${detail}` : ""}`,
867 );
868 }
869 await put.body?.cancel();
870 } finally {
871 await removeContainer(sourceId);
872 }
873}
874
875async function removeContainer(containerId: string): Promise<void> {
876 try {
877 const resp = await dockerFetch(
878 `/containers/${containerId}?force=true`,
879 { method: "DELETE" },
880 );
881 await resp.body?.cancel();
882 } catch {
883 // Best-effort cleanup
884 }
885}
886
887// --- Tar extraction ---
888
889function extractSingleFileFromTar(data: Uint8Array): Uint8Array | null {
890 if (data.length < 512) return null;
891 const dec = new TextDecoder();
892 const sizeOctal = dec
893 .decode(data.slice(124, 136))
894 .replace(/\0/g, "")
895 .trim();
896 const size = parseInt(sizeOctal, 8);
897 if (Number.isNaN(size) || size < 0) return null;
898 if (data.length < 512 + size) return null;
899 return data.slice(512, 512 + size);
900}
901
902async function copyFileFromContainer(
903 containerId: string,
904 containerPath: string,
905): Promise<Uint8Array | null> {
906 const resp = await dockerFetch(
907 `/containers/${containerId}/archive?path=${encodeURIComponent(containerPath)}`,
908 );
909 if (!resp.ok) return null;
910 const tarBytes = new Uint8Array(await resp.arrayBuffer());
911 return extractSingleFileFromTar(tarBytes);
912}
913
914// --- Secret masking ---
915
916function maskSecrets(text: string, secrets: string[]): string {
917 for (const secret of secrets) {
918 if (secret) text = text.split(secret).join("[MASKED]");
919 }
920 return text;
921}
922
923// --- Env var building ---
924
925function buildEnvVars(
926 runId: number,
927 repoName: string,
928 opts: TriggerOpts,
929 cfg: CiConfig,
930 secretValues: Array<{ name: string; value: string }>,
931): { envArray: string[]; secretValues: string[] } {
932 const vars: Record<string, string> = {
933 CI: "true",
934 CI_PIPELINE_ID: String(runId),
935 CI_REPO_NAME: repoName,
936 CI_SERVER_URL: config.BASE_URL,
937 CI_TRIGGER_SOURCE: opts.triggerSource,
938 CI_COMMIT_SHA: opts.commitSha,
939 CI_COMMIT_SHORT_SHA: opts.commitSha.slice(0, 8),
940 CI_COMMIT_BRANCH: opts.commitBranch ?? "",
941 CI_COMMIT_TAG: opts.commitTag ?? "",
942 CI_COMMIT_REF_NAME: opts.commitTag ?? opts.commitBranch ?? "",
943 };
944
945 // User-defined variable defaults
946 if (cfg.variables) {
947 for (const [name, def] of Object.entries(cfg.variables)) {
948 if (def.default !== undefined) vars[name] = def.default;
949 }
950 }
951
952 // Variable overrides from manual trigger
953 if (opts.variableOverrides) {
954 for (const [name, value] of Object.entries(opts.variableOverrides)) {
955 vars[name] = value;
956 }
957 }
958
959 // Secrets (injected but values tracked for masking)
960 const secretVals: string[] = [];
961 for (const { name, value } of secretValues) {
962 vars[name] = value;
963 secretVals.push(value);
964 }
965
966 const envArray = Object.entries(vars).map(([k, v]) => `${k}=${v}`);
967 return { envArray, secretValues: secretVals };
968}
969
970// --- Artifact collection ---
971
972async function collectArtifacts(
973 runId: number,
974 containerId: string,
975 step: CiStep,
976 _shell: string[],
977 envVars: string[],
978): Promise<void> {
979 const artifactDir = path.join(paths.CI_ARTIFACTS_DIR, String(runId));
980 mkdirSync(artifactDir, { recursive: true });
981
982 const toArray = (v: string | string[] | undefined): string[] => {
983 if (!v) return [];
984 return Array.isArray(v) ? v : [v];
985 };
986
987 // publish_file: copy directly out of container
988 for (const srcPath of toArray(step.publish_file)) {
989 const fileBytes = await copyFileFromContainer(containerId, srcPath);
990 if (fileBytes) {
991 const filename = path.basename(srcPath);
992 const destPath = path.join(artifactDir, filename);
993 writeFileSync(destPath, fileBytes);
994 const stat = Bun.file(destPath);
995 await db
996 .insertInto("ci_artifacts")
997 .values({
998 run_id: runId,
999 filename,
1000 size: stat.size,
1001 })
1002 .execute();
1003 }
1004 }
1005
1006 // Archive commands run with `path.dirname(srcPath)` as the working
1007 // directory and reference the source by its basename only, so the
1008 // user-controlled path never appears as part of an interpolated shell
1009 // string. Previously `publish_zip` used `sh -c "cd … && zip …"` with
1010 // raw interpolation — a step author who could write the CI TOML
1011 // could shell-inject through the source path. Today the only TOML
1012 // author is the admin, but this removes the implicit assumption.
1013 type ArchiveType = "tar" | "gzip" | "zip" | "zstd";
1014 const archiveFormats: Array<{
1015 type: ArchiveType;
1016 paths: string[];
1017 ext: string;
1018 cmd: (basename: string, dst: string) => string[];
1019 }> = [
1020 {
1021 type: "tar",
1022 paths: toArray(step.publish_tar),
1023 ext: ".tar",
1024 cmd: (basename, dst) => ["tar", "-cf", dst, basename],
1025 },
1026 {
1027 type: "gzip",
1028 paths: toArray(step.publish_gzip),
1029 ext: ".tar.gz",
1030 cmd: (basename, dst) => ["tar", "-czf", dst, basename],
1031 },
1032 {
1033 type: "zstd",
1034 paths: toArray(step.publish_zstd),
1035 ext: ".tar.zst",
1036 cmd: (basename, dst) => ["tar", "--zstd", "-cf", dst, basename],
1037 },
1038 {
1039 type: "zip",
1040 paths: toArray(step.publish_zip),
1041 ext: ".zip",
1042 cmd: (basename, dst) => ["zip", "-r", dst, basename],
1043 },
1044 ];
1045
1046 let archiveIndex = 0;
1047 for (const { paths: archivePaths, ext, cmd } of archiveFormats) {
1048 for (const srcPath of archivePaths) {
1049 archiveIndex++;
1050 const tmpPath = `/tmp/hf-artifact-${runId}-${archiveIndex}${ext}`;
1051 // Run with the source's parent directory as the working
1052 // directory so each archive tool can reference the source
1053 // by its basename — no -C, no shell.
1054 const execResult = await execInContainer(
1055 containerId,
1056 cmd(path.basename(srcPath), tmpPath),
1057 path.dirname(srcPath),
1058 envVars,
1059 ).catch(() => null);
1060 if (!execResult || execResult.exitCode !== 0) continue;
1061
1062 // Copy archive out
1063 const fileBytes = await copyFileFromContainer(containerId, tmpPath);
1064 if (!fileBytes) continue;
1065
1066 const filename = `${path.basename(srcPath)}${ext}`;
1067 const destPath = path.join(artifactDir, filename);
1068 writeFileSync(destPath, fileBytes);
1069 const stat = Bun.file(destPath);
1070 await db
1071 .insertInto("ci_artifacts")
1072 .values({
1073 run_id: runId,
1074 filename,
1075 size: stat.size,
1076 })
1077 .execute();
1078 }
1079 }
1080}
1081
1082// --- Main execution ---
1083
1084async function executeRun(runId: number, signal: AbortSignal): Promise<void> {
1085 const now = () => new Date().toISOString();
1086 let containerId: string | undefined;
1087 // Set once the config is known, and only for a default-branch run. The
1088 // finally block cannot see cfg, and pruning off another branch is wrong.
1089 let pruneCaches: { repoName: string; cache?: CiCache[] } | null = null;
1090 let cacheLimits: { repoName: string; cache?: CiCache[] } | null = null;
1091
1092 try {
1093 // Mark as running
1094 await db
1095 .updateTable("ci_runs")
1096 .set({ status: "running", started_at: now() })
1097 .where("id", "=", runId)
1098 .execute();
1099
1100 // Load run details
1101 const run = await db
1102 .selectFrom("ci_runs")
1103 .selectAll()
1104 .where("id", "=", runId)
1105 .executeTakeFirst();
1106 if (!run) throw new Error("Run not found");
1107
1108 const repo = await db
1109 .selectFrom("repositories")
1110 .select(["id", "name", "default_branch"])
1111 .where("id", "=", run.repo_id)
1112 .executeTakeFirst();
1113 if (!repo) throw new Error("Repo not found");
1114
1115 // Read .hearthforge-ci.toml at the commit
1116 const tomlBuf = await import("./git.ts").then((g) =>
1117 g.git.show(repo.name, run.commit_sha!, ".hearthforge-ci.toml"),
1118 );
1119 if (!tomlBuf)
1120 throw new Error(".hearthforge-ci.toml not found at commit");
1121
1122 const cfg = parseCiConfig(tomlBuf.toString("utf-8"));
1123 if (!cfg) throw new Error("Failed to parse .hearthforge-ci.toml");
1124
1125 const invalid = validateCiConfig(cfg);
1126 if (invalid)
1127 throw new Error(`Invalid .hearthforge-ci.toml: ${invalid}`);
1128
1129 // Size caps apply on any branch: an oversized volume is oversized
1130 // whoever noticed. Dropping stale volumes is default-branch only,
1131 // because the config that names them is read per commit.
1132 cacheLimits = { repoName: repo.name, cache: cfg.cache };
1133 if (run.commit_branch && run.commit_branch === repo.default_branch) {
1134 pruneCaches = { repoName: repo.name, cache: cfg.cache };
1135 }
1136
1137 // Load secrets for log masking
1138 const secrets = await db
1139 .selectFrom("ci_secrets")
1140 .select(["name", "value"])
1141 .where("repo_id", "=", repo.id)
1142 .execute();
1143
1144 const variableOverrides = run.variable_overrides
1145 ? (JSON.parse(run.variable_overrides) as Record<string, string>)
1146 : {};
1147
1148 const { envArray, secretValues } = buildEnvVars(
1149 runId,
1150 repo.name,
1151 {
1152 triggerSource:
1153 run.trigger_source as TriggerOpts["triggerSource"],
1154 commitSha: run.commit_sha ?? "",
1155 commitBranch: run.commit_branch ?? undefined,
1156 commitTag: run.commit_tag ?? undefined,
1157 variableOverrides,
1158 },
1159 cfg,
1160 secrets,
1161 );
1162
1163 // Step names need not be unique, so a name cannot identify a row.
1164 // Keep the ids in config order instead.
1165 const stepIds: number[] = [];
1166 for (const step of cfg.steps) {
1167 const row = await db
1168 .insertInto("ci_steps")
1169 .values({
1170 run_id: runId,
1171 name: step.name,
1172 status: "pending",
1173 })
1174 .returning("id")
1175 .executeTakeFirstOrThrow();
1176 stepIds.push(row.id);
1177 }
1178
1179 // Pull image
1180 await pullImage(cfg.image);
1181 if (signal.aborted) throw new Error("Cancelled");
1182
1183 // Create + start container
1184 containerId = await createContainer(runId, repo.name, cfg, envArray);
1185 runningTasks.get(runId)!.containerId = containerId;
1186
1187 await startContainer(containerId);
1188 if (signal.aborted) throw new Error("Cancelled");
1189
1190 // Create work_dir
1191 if (cfg.work_dir) {
1192 await execInContainer(containerId, ["mkdir", "-p", cfg.work_dir]);
1193 }
1194
1195 // Copies run before the clone, so a step can rely on the tools they
1196 // bring in, and so they can supply a shell the image lacks.
1197 for (const spec of cfg.copy ?? []) {
1198 await copyFromImage(containerId, spec);
1199 if (signal.aborted) throw new Error("Cancelled");
1200 }
1201
1202 if (cfg.clone_project_to && run.commit_sha) {
1203 await uploadCheckout(
1204 containerId,
1205 repoPath(repo.name),
1206 run.commit_sha,
1207 cfg.clone_project_to,
1208 signal,
1209 );
1210 }
1211
1212 // Execute steps
1213 const shell = cfg.shell ?? ["/bin/sh", "-c"];
1214 let runFailed = false;
1215 let sawWarning = false;
1216
1217 for (const [stepIndex, step] of cfg.steps.entries()) {
1218 if (signal.aborted) {
1219 runFailed = true;
1220 break;
1221 }
1222
1223 // After a failure, the skipped steps stay pending and are marked
1224 // skipped below.
1225 if (runFailed && !step.always) continue;
1226
1227 // A step timeout removes the container to kill the command, so
1228 // there is nothing left to run an `always` step in.
1229 if (!containerId) continue;
1230
1231 const stepId = stepIds[stepIndex]!;
1232
1233 // Check run_if condition
1234 if (step.run_if) {
1235 const { exitCode } = await execInContainer(
1236 containerId,
1237 [...shell, step.run_if],
1238 cfg.work_dir,
1239 envArray,
1240 );
1241 if (exitCode !== 0) {
1242 await db
1243 .updateTable("ci_steps")
1244 .set({
1245 status: "skipped",
1246 started_at: now(),
1247 finished_at: now(),
1248 log: "Skipped: condition not met",
1249 })
1250 .where("id", "=", stepId)
1251 .execute();
1252 continue;
1253 }
1254 }
1255
1256 // Handle clear option
1257 // Drop the directory and re-extract. `git clean` would need git
1258 // in the image, and this also removes files git never tracked.
1259 if (step.clear && cfg.clone_project_to && run.commit_sha) {
1260 let clearError: string | null = null;
1261 try {
1262 // One exec, not two. `rm -rf` can delete the container's
1263 // WorkingDir, and every later exec then fails to chdir
1264 // before its command starts. This exec chdirs first.
1265 const reset = await execInContainer(containerId, [
1266 ...shell,
1267 'rm -rf "$1" && mkdir -p "$1"',
1268 "sh",
1269 cfg.clone_project_to,
1270 ]);
1271 if (reset.exitCode !== 0) {
1272 throw new Error(reset.log.trim());
1273 }
1274 await uploadCheckout(
1275 containerId,
1276 repoPath(repo.name),
1277 run.commit_sha,
1278 cfg.clone_project_to,
1279 signal,
1280 );
1281 } catch (err) {
1282 clearError =
1283 err instanceof Error ? err.message : String(err);
1284 }
1285 if (clearError !== null) {
1286 // Not thrown: the outer handler only records a message
1287 // when no step has failed yet, so after an earlier failure
1288 // it would vanish. The step row always survives.
1289 await db
1290 .updateTable("ci_steps")
1291 .set({
1292 status: "failure",
1293 started_at: now(),
1294 finished_at: now(),
1295 log: `Failed to reset ${cfg.clone_project_to}: ${clearError}\n`,
1296 })
1297 .where("id", "=", stepId)
1298 .execute();
1299 runFailed = true;
1300 continue;
1301 }
1302 }
1303
1304 await db
1305 .updateTable("ci_steps")
1306 .set({ status: "running", started_at: now() })
1307 .where("id", "=", stepId)
1308 .execute();
1309
1310 let stepLog = "";
1311 let stepStatus: "success" | "failure" | "warning" = "success";
1312 // Applies to every way a step can fail, a timeout included.
1313 const onFail = step.warn_on_fail ? "warning" : "failure";
1314
1315 if (step.run_sh) {
1316 const command = cfg.shell_setup
1317 ? `${cfg.shell_setup}\n${step.run_sh}`
1318 : step.run_sh;
1319
1320 const stepTimeout =
1321 step.timeout ?? cfg.timeout ?? config.CI_DEFAULT_TIMEOUT;
1322 const timeoutSignal = AbortSignal.timeout(stepTimeout * 1000);
1323 // Abort the exec stream on either a run cancellation or the
1324 // per-step timeout. Docker has no per-exec kill, so on abort we
1325 // force-remove the container (below), which kills the command
1326 // still running inside it.
1327 const stepSignal = AbortSignal.any([signal, timeoutSignal]);
1328
1329 try {
1330 const { log, exitCode } = await execInContainer(
1331 containerId,
1332 [...shell, command],
1333 cfg.work_dir,
1334 envArray,
1335 stepSignal,
1336 async (partial) => {
1337 await db
1338 .updateTable("ci_steps")
1339 .set({
1340 log: maskSecrets(partial, secretValues),
1341 })
1342 .where("id", "=", stepId)
1343 .execute();
1344 },
1345 );
1346 stepLog = maskSecrets(log, secretValues);
1347 if (exitCode !== 0) stepStatus = onFail;
1348 } catch (err) {
1349 // A run cancellation is reported as "cancelled" by the outer
1350 // catch — don't relabel it as a step failure here.
1351 if (signal.aborted) throw err;
1352 if (timeoutSignal.aborted) {
1353 // Not subject to warn_on_fail. The container is about
1354 // to be destroyed, so every later step is skipped no
1355 // matter what. A run that cannot continue is a failure.
1356 stepStatus = "failure";
1357 stepLog = `Step timed out after ${stepTimeout}s\n`;
1358 // Kill the container now so the timed-out command stops
1359 // immediately rather than lingering until cleanup.
1360 await removeContainer(containerId);
1361 containerId = "";
1362 } else {
1363 stepStatus = onFail;
1364 stepLog = `Step failed: ${err instanceof Error ? err.message : String(err)}\n`;
1365 }
1366 }
1367 }
1368
1369 if (stepStatus === "failure") runFailed = true;
1370 if (stepStatus === "warning") sawWarning = true;
1371
1372 // Collect artifacts for this step
1373 if (stepStatus !== "failure") {
1374 await collectArtifacts(
1375 runId,
1376 containerId,
1377 step,
1378 shell,
1379 envArray,
1380 ).catch(() => {});
1381 }
1382
1383 await db
1384 .updateTable("ci_steps")
1385 .set({
1386 status: stepStatus,
1387 finished_at: now(),
1388 log: stepLog,
1389 })
1390 .where("id", "=", stepId)
1391 .execute();
1392 }
1393
1394 // Mark remaining steps as skipped
1395 await db
1396 .updateTable("ci_steps")
1397 .set({
1398 status: "skipped",
1399 started_at: now(),
1400 finished_at: now(),
1401 log: "Skipped: previous step failed",
1402 })
1403 .where("run_id", "=", runId)
1404 .where("status", "=", "pending")
1405 .execute();
1406
1407 const finalStatus = runFailed
1408 ? "failure"
1409 : sawWarning
1410 ? "warning"
1411 : "success";
1412 await db
1413 .updateTable("ci_runs")
1414 .set({ status: finalStatus, finished_at: now() })
1415 .where("id", "=", runId)
1416 .execute();
1417 } catch (err) {
1418 const errMsg = err instanceof Error ? err.message : String(err);
1419 const isDockerUnavailable = errMsg.includes("No Docker/Podman socket");
1420 const status = signal.aborted
1421 ? "cancelled"
1422 : isDockerUnavailable
1423 ? "skipped"
1424 : "failure";
1425 const skipLog = signal.aborted
1426 ? "Skipped: run was cancelled"
1427 : isDockerUnavailable
1428 ? "Skipped: Docker/Podman not available"
1429 : "Skipped: run failed";
1430
1431 if (!isDockerUnavailable) {
1432 // A failed step already carries its own log. Failures outside any
1433 // step (validation, upload, clone) would leave no reason at all.
1434 const failedStep = await db
1435 .selectFrom("ci_steps")
1436 .select("id")
1437 .where("run_id", "=", runId)
1438 .where("status", "=", "failure")
1439 .executeTakeFirst();
1440 if (!failedStep) {
1441 await db
1442 .insertInto("ci_steps")
1443 .values({
1444 run_id: runId,
1445 // Not "setup": a pipeline may well have its own step
1446 // by that name.
1447 name: "pipeline setup",
1448 status: "failure",
1449 started_at: new Date().toISOString(),
1450 finished_at: new Date().toISOString(),
1451 log: `Error: ${errMsg}\n`,
1452 })
1453 .execute();
1454 }
1455 }
1456 await db
1457 .updateTable("ci_runs")
1458 .set({ status, finished_at: new Date().toISOString() })
1459 .where("id", "=", runId)
1460 .execute();
1461 // Mark pending steps as skipped
1462 await db
1463 .updateTable("ci_steps")
1464 .set({
1465 status: "skipped",
1466 finished_at: new Date().toISOString(),
1467 log: skipLog,
1468 })
1469 .where("run_id", "=", runId)
1470 .where("status", "=", "pending")
1471 .execute();
1472 } finally {
1473 if (containerId) await removeContainer(containerId);
1474 runningTasks.delete(runId);
1475 // Prune old history
1476 const run = await db
1477 .selectFrom("ci_runs")
1478 .select("repo_id")
1479 .where("id", "=", runId)
1480 .executeTakeFirst();
1481 if (run) await pruneHistory(run.repo_id).catch(() => {});
1482 if (pruneCaches) {
1483 await pruneStaleCaches(
1484 pruneCaches.repoName,
1485 pruneCaches.cache,
1486 ).catch(() => {});
1487 }
1488 // After removeContainer above: the daemon will not delete a volume
1489 // that is still mounted.
1490 if (cacheLimits) {
1491 const dropped = await enforceCacheLimits(
1492 cacheLimits.repoName,
1493 cacheLimits.cache,
1494 ).catch(() => []);
1495 if (dropped.length) {
1496 // Recorded on the run, because the only other symptom is the
1497 // next build being mysteriously slow.
1498 await db
1499 .insertInto("ci_steps")
1500 .values({
1501 run_id: runId,
1502 name: "cache",
1503 status: "success",
1504 started_at: new Date().toISOString(),
1505 finished_at: new Date().toISOString(),
1506 log: `${dropped
1507 .map(
1508 (d) =>
1509 `Dropped cache ${d.path}: ${formatBytes(d.size)} over the ${formatBytes(d.maxSize)} limit`,
1510 )
1511 .join("\n")}\n`,
1512 })
1513 .execute()
1514 .catch(() => {});
1515 }
1516 }
1517 // A slot just freed up — start the next queued run if any.
1518 pumpQueue();
1519 }
1520}
1521
1522// --- Public API ---
1523
1524export async function triggerRun(
1525 repoName: string,
1526 opts: TriggerOpts,
1527): Promise<number> {
1528 const repo = await db
1529 .selectFrom("repositories")
1530 .select("id")
1531 .where("name", "=", repoName)
1532 .executeTakeFirst();
1533 if (!repo) throw new Error("Repository not found");
1534
1535 const runId = await db
1536 .insertInto("ci_runs")
1537 .values({
1538 repo_id: repo.id,
1539 triggered_by: opts.triggeredBy ?? null,
1540 trigger_source: opts.triggerSource,
1541 commit_sha: opts.commitSha,
1542 commit_branch: opts.commitBranch ?? null,
1543 commit_tag: opts.commitTag ?? null,
1544 status: "pending",
1545 variable_overrides: opts.variableOverrides
1546 ? JSON.stringify(opts.variableOverrides)
1547 : null,
1548 })
1549 .returning("id")
1550 .executeTakeFirstOrThrow();
1551
1552 // Allocate the human-facing run number atomically from a per-repo counter.
1553 // A single upsert-and-increment can't collide under concurrent triggers and
1554 // never reuses a number after pruneHistory shrinks the run table — both of
1555 // which a COUNT(*)-based scheme suffered from.
1556 const counter = await db
1557 .insertInto("ci_run_counters")
1558 .values({ repo_id: repo.id, last_run_id: 1 })
1559 .onConflict((oc) =>
1560 oc.column("repo_id").doUpdateSet((eb) => ({
1561 last_run_id: eb("ci_run_counters.last_run_id", "+", 1),
1562 })),
1563 )
1564 .returning("last_run_id")
1565 .executeTakeFirstOrThrow();
1566 await db
1567 .updateTable("ci_runs")
1568 .set({ repo_run_id: counter.last_run_id })
1569 .where("id", "=", runId.id)
1570 .execute();
1571
1572 // Flip to "queued" BEFORE pushing onto queuedRunIds. Otherwise a
1573 // concurrent pumpQueue (from a finishing run) could observe our entry,
1574 // promoteToPending it, and start executeRun while our UPDATE is still
1575 // in flight — the late UPDATE would then clobber a running row back to
1576 // "queued". Only pumpQueue (the sole writer of runningTasks.set) may
1577 // mutate the row's status after the push. The TOCTOU concern is moot
1578 // here because pumpQueue cannot pump a runId that hasn't been pushed.
1579 if (runningTasks.size >= config.CI_MAX_CONCURRENT) {
1580 await db
1581 .updateTable("ci_runs")
1582 .set({ status: "queued" })
1583 .where("id", "=", runId.id)
1584 .execute();
1585 }
1586 queuedRunIds.push(runId.id);
1587 pumpQueue();
1588
1589 return runId.id;
1590}
1591
1592// Wraps the fire-and-forget executeRun so unhandled exceptions (e.g. a throw
1593// before/after its own try/finally) are logged and the run is reconciled to
1594// "failure" instead of staying pending forever.
1595function spawnRun(
1596 runId: number,
1597 signal: AbortSignal,
1598 needsPromote = false,
1599): void {
1600 void (async () => {
1601 try {
1602 if (needsPromote) await promoteToPending(runId);
1603 await executeRun(runId, signal);
1604 } catch (err) {
1605 console.error(`[ci] executeRun threw for run ${runId}:`, err);
1606 runningTasks.delete(runId);
1607 try {
1608 await db
1609 .updateTable("ci_runs")
1610 .set({
1611 status: "failure",
1612 finished_at: new Date().toISOString(),
1613 })
1614 .where("id", "=", runId)
1615 // Include "queued" so a row that was promoted-but-not-yet-
1616 // observed (or never promoted because promoteToPending threw)
1617 // still gets marked failed instead of stuck.
1618 .where("status", "in", ["pending", "running", "queued"])
1619 .execute();
1620 } catch (dbErr) {
1621 console.error(
1622 `[ci] failed to mark run ${runId} failed:`,
1623 dbErr,
1624 );
1625 }
1626 pumpQueue();
1627 }
1628 })();
1629}
1630
1631export async function retryRun(
1632 runId: number,
1633 retriedBy: number,
1634): Promise<void> {
1635 const run = await db
1636 .selectFrom("ci_runs")
1637 .select("repo_id")
1638 .where("id", "=", runId)
1639 .executeTakeFirst();
1640 if (!run) throw new Error("Run not found");
1641
1642 // Delete existing steps
1643 await db.deleteFrom("ci_steps").where("run_id", "=", runId).execute();
1644
1645 // Delete artifacts from disk and DB
1646 const artifactDir = path.join(paths.CI_ARTIFACTS_DIR, String(runId));
1647 if (existsSync(artifactDir)) {
1648 await Bun.$`rm -rf ${artifactDir}`.quiet().nothrow();
1649 }
1650 await db.deleteFrom("ci_artifacts").where("run_id", "=", runId).execute();
1651
1652 // Reset to "pending" first; if a slot isn't free, flip to "queued"
1653 // before enqueueing. Mirrors triggerRun: pumpQueue is the sole writer
1654 // of runningTasks.set, so we never reserve a slot directly here. Two
1655 // concurrent retryRun calls (or a retryRun racing triggerRun) all
1656 // funnel through pumpQueue, which serializes the size check against
1657 // slot reservation in a single synchronous turn.
1658 await db
1659 .updateTable("ci_runs")
1660 .set({
1661 status: "pending",
1662 triggered_by: retriedBy,
1663 started_at: null,
1664 finished_at: null,
1665 })
1666 .where("id", "=", runId)
1667 .execute();
1668
1669 if (runningTasks.size >= config.CI_MAX_CONCURRENT) {
1670 await db
1671 .updateTable("ci_runs")
1672 .set({ status: "queued" })
1673 .where("id", "=", runId)
1674 .execute();
1675 }
1676 queuedRunIds.push(runId);
1677 pumpQueue();
1678}
1679
1680export async function cancelRun(runId: number): Promise<void> {
1681 const task = runningTasks.get(runId);
1682 if (task) {
1683 const { containerId } = task;
1684 task.controller.abort();
1685 if (containerId) {
1686 await removeContainer(containerId).catch(() => {});
1687 }
1688 }
1689 const queueIdx = queuedRunIds.indexOf(runId);
1690 if (queueIdx >= 0) queuedRunIds.splice(queueIdx, 1);
1691 await db
1692 .updateTable("ci_runs")
1693 .set({ status: "cancelled", finished_at: new Date().toISOString() })
1694 .where("id", "=", runId)
1695 .where("status", "in", ["pending", "running", "queued"])
1696 .execute();
1697}
1698
1699async function pruneHistory(repoId: number): Promise<void> {
1700 const maxHistory = config.CI_MAX_HISTORY;
1701 const allRuns = await db
1702 .selectFrom("ci_runs")
1703 .select("id")
1704 .where("repo_id", "=", repoId)
1705 .orderBy("id", "desc")
1706 .execute();
1707
1708 if (allRuns.length <= maxHistory) return;
1709
1710 const toDelete = allRuns.slice(maxHistory).map((r) => r.id);
1711 for (const runId of toDelete) {
1712 // Remove artifacts from disk
1713 const artifactDir = path.join(paths.CI_ARTIFACTS_DIR, String(runId));
1714 if (existsSync(artifactDir)) {
1715 await Bun.$`rm -rf ${artifactDir}`.quiet().nothrow();
1716 }
1717 }
1718 await db.deleteFrom("ci_runs").where("id", "in", toDelete).execute();
1719}
1720
1721export async function purgeRepoCaches(repoName: string): Promise<void> {
1722 try {
1723 const filters = encodeURIComponent(
1724 JSON.stringify({ label: [`com.hearthforge.repo=${repoName}`] }),
1725 );
1726 const resp = await dockerFetch(`/volumes?filters=${filters}`);
1727 if (!resp.ok) {
1728 await resp.body?.cancel();
1729 return;
1730 }
1731 const data = (await resp.json()) as {
1732 Volumes?: Array<{ Name: string }>;
1733 };
1734 for (const vol of data.Volumes ?? []) {
1735 const delResp = await dockerFetch(`/volumes/${vol.Name}`, {
1736 method: "DELETE",
1737 });
1738 await delResp.body?.cancel();
1739 }
1740 } catch {
1741 // Best-effort
1742 }
1743}
1744
1745export async function cancelStaleRuns(): Promise<void> {
1746 const now = new Date().toISOString();
1747
1748 const stale = await db
1749 .selectFrom("ci_runs")
1750 .select("id")
1751 .where("status", "in", ["pending", "running", "queued"])
1752 .execute();
1753
1754 await Promise.allSettled(
1755 stale.map((r) =>
1756 dockerFetch(`/containers/hearthforge-ci-${r.id}?force=true`, {
1757 method: "DELETE",
1758 }).then((res) => res.body?.cancel()),
1759 ),
1760 );
1761
1762 await db
1763 .updateTable("ci_runs")
1764 .set({ status: "cancelled", finished_at: now })
1765 .where("status", "in", ["pending", "running", "queued"])
1766 .execute();
1767 await db
1768 .updateTable("ci_steps")
1769 .set({
1770 status: "cancelled",
1771 finished_at: now,
1772 log: "Skipped: run was cancelled",
1773 })
1774 .where("status", "in", ["pending", "running"])
1775 .execute();
1776}
1777
1778/** Reset the cached socket path (used in tests to switch mock sockets). */
1779export function resetDockerSocket(): void {
1780 resolvedSocket = null;
1781}
1782