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
11export interface 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 // Inserted before any config step, so ordering by id puts it first. It
1093 // covers everything up to the first step: config parsing, the image pull,
1094 // the container, the copies, and the checkout.
1095 // Not "setup": a pipeline may well have its own step by that name.
1096 const setupStepId = (
1097 await db
1098 .insertInto("ci_steps")
1099 .values({
1100 run_id: runId,
1101 name: "pipeline setup",
1102 status: "running",
1103 started_at: now(),
1104 })
1105 .returning("id")
1106 .executeTakeFirstOrThrow()
1107 ).id;
1108
1109 try {
1110 // Mark as running
1111 await db
1112 .updateTable("ci_runs")
1113 .set({ status: "running", started_at: now() })
1114 .where("id", "=", runId)
1115 .execute();
1116
1117 // Load run details
1118 const run = await db
1119 .selectFrom("ci_runs")
1120 .selectAll()
1121 .where("id", "=", runId)
1122 .executeTakeFirst();
1123 if (!run) throw new Error("Run not found");
1124
1125 const repo = await db
1126 .selectFrom("repositories")
1127 .select(["id", "name", "default_branch"])
1128 .where("id", "=", run.repo_id)
1129 .executeTakeFirst();
1130 if (!repo) throw new Error("Repo not found");
1131
1132 // Read .hearthforge-ci.toml at the commit
1133 const tomlBuf = await import("./git.ts").then((g) =>
1134 g.git.show(repo.name, run.commit_sha!, ".hearthforge-ci.toml"),
1135 );
1136 if (!tomlBuf)
1137 throw new Error(".hearthforge-ci.toml not found at commit");
1138
1139 const cfg = parseCiConfig(tomlBuf.toString("utf-8"));
1140 if (!cfg) throw new Error("Failed to parse .hearthforge-ci.toml");
1141
1142 const invalid = validateCiConfig(cfg);
1143 if (invalid)
1144 throw new Error(`Invalid .hearthforge-ci.toml: ${invalid}`);
1145
1146 // Size caps apply on any branch: an oversized volume is oversized
1147 // whoever noticed. Dropping stale volumes is default-branch only,
1148 // because the config that names them is read per commit.
1149 cacheLimits = { repoName: repo.name, cache: cfg.cache };
1150 if (run.commit_branch && run.commit_branch === repo.default_branch) {
1151 pruneCaches = { repoName: repo.name, cache: cfg.cache };
1152 }
1153
1154 // Load secrets for log masking
1155 const secrets = await db
1156 .selectFrom("ci_secrets")
1157 .select(["name", "value"])
1158 .where("repo_id", "=", repo.id)
1159 .execute();
1160
1161 const variableOverrides = run.variable_overrides
1162 ? (JSON.parse(run.variable_overrides) as Record<string, string>)
1163 : {};
1164
1165 const { envArray, secretValues } = buildEnvVars(
1166 runId,
1167 repo.name,
1168 {
1169 triggerSource:
1170 run.trigger_source as TriggerOpts["triggerSource"],
1171 commitSha: run.commit_sha ?? "",
1172 commitBranch: run.commit_branch ?? undefined,
1173 commitTag: run.commit_tag ?? undefined,
1174 variableOverrides,
1175 },
1176 cfg,
1177 secrets,
1178 );
1179
1180 // Step names need not be unique, so a name cannot identify a row.
1181 // Keep the ids in config order instead.
1182 const stepIds: number[] = [];
1183 for (const step of cfg.steps) {
1184 const row = await db
1185 .insertInto("ci_steps")
1186 .values({
1187 run_id: runId,
1188 name: step.name,
1189 status: "pending",
1190 })
1191 .returning("id")
1192 .executeTakeFirstOrThrow();
1193 stepIds.push(row.id);
1194 }
1195
1196 // Pull image
1197 await pullImage(cfg.image);
1198 if (signal.aborted) throw new Error("Cancelled");
1199
1200 // Create + start container
1201 containerId = await createContainer(runId, repo.name, cfg, envArray);
1202 runningTasks.get(runId)!.containerId = containerId;
1203
1204 await startContainer(containerId);
1205 if (signal.aborted) throw new Error("Cancelled");
1206
1207 // Create work_dir
1208 if (cfg.work_dir) {
1209 await execInContainer(containerId, ["mkdir", "-p", cfg.work_dir]);
1210 }
1211
1212 // Copies run before the clone, so a step can rely on the tools they
1213 // bring in, and so they can supply a shell the image lacks.
1214 for (const spec of cfg.copy ?? []) {
1215 await copyFromImage(containerId, spec);
1216 if (signal.aborted) throw new Error("Cancelled");
1217 }
1218
1219 if (cfg.clone_project_to && run.commit_sha) {
1220 await uploadCheckout(
1221 containerId,
1222 repoPath(repo.name),
1223 run.commit_sha,
1224 cfg.clone_project_to,
1225 signal,
1226 );
1227 }
1228
1229 await db
1230 .updateTable("ci_steps")
1231 .set({ status: "success", finished_at: now() })
1232 .where("id", "=", setupStepId)
1233 .execute();
1234
1235 // Execute steps
1236 const shell = cfg.shell ?? ["/bin/sh", "-c"];
1237 let runFailed = false;
1238 let sawWarning = false;
1239
1240 for (const [stepIndex, step] of cfg.steps.entries()) {
1241 if (signal.aborted) {
1242 runFailed = true;
1243 break;
1244 }
1245
1246 // After a failure, the skipped steps stay pending and are marked
1247 // skipped below.
1248 if (runFailed && !step.always) continue;
1249
1250 // A step timeout removes the container to kill the command, so
1251 // there is nothing left to run an `always` step in.
1252 if (!containerId) continue;
1253
1254 const stepId = stepIds[stepIndex]!;
1255
1256 // Check run_if condition
1257 if (step.run_if) {
1258 const { exitCode } = await execInContainer(
1259 containerId,
1260 [...shell, step.run_if],
1261 cfg.work_dir,
1262 envArray,
1263 );
1264 if (exitCode !== 0) {
1265 await db
1266 .updateTable("ci_steps")
1267 .set({
1268 status: "skipped",
1269 started_at: now(),
1270 finished_at: now(),
1271 log: "Skipped: condition not met",
1272 })
1273 .where("id", "=", stepId)
1274 .execute();
1275 continue;
1276 }
1277 }
1278
1279 // Handle clear option
1280 // Drop the directory and re-extract. `git clean` would need git
1281 // in the image, and this also removes files git never tracked.
1282 if (step.clear && cfg.clone_project_to && run.commit_sha) {
1283 let clearError: string | null = null;
1284 try {
1285 // One exec, not two. `rm -rf` can delete the container's
1286 // WorkingDir, and every later exec then fails to chdir
1287 // before its command starts. This exec chdirs first.
1288 const reset = await execInContainer(containerId, [
1289 ...shell,
1290 'rm -rf "$1" && mkdir -p "$1"',
1291 "sh",
1292 cfg.clone_project_to,
1293 ]);
1294 if (reset.exitCode !== 0) {
1295 throw new Error(reset.log.trim());
1296 }
1297 await uploadCheckout(
1298 containerId,
1299 repoPath(repo.name),
1300 run.commit_sha,
1301 cfg.clone_project_to,
1302 signal,
1303 );
1304 } catch (err) {
1305 clearError =
1306 err instanceof Error ? err.message : String(err);
1307 }
1308 if (clearError !== null) {
1309 // Not thrown: the outer handler only records a message
1310 // when no step has failed yet, so after an earlier failure
1311 // it would vanish. The step row always survives.
1312 await db
1313 .updateTable("ci_steps")
1314 .set({
1315 status: "failure",
1316 started_at: now(),
1317 finished_at: now(),
1318 log: `Failed to reset ${cfg.clone_project_to}: ${clearError}\n`,
1319 })
1320 .where("id", "=", stepId)
1321 .execute();
1322 runFailed = true;
1323 continue;
1324 }
1325 }
1326
1327 await db
1328 .updateTable("ci_steps")
1329 .set({ status: "running", started_at: now() })
1330 .where("id", "=", stepId)
1331 .execute();
1332
1333 let stepLog = "";
1334 let stepStatus: "success" | "failure" | "warning" = "success";
1335 // Applies to every way a step can fail, a timeout included.
1336 const onFail = step.warn_on_fail ? "warning" : "failure";
1337
1338 if (step.run_sh) {
1339 const command = cfg.shell_setup
1340 ? `${cfg.shell_setup}\n${step.run_sh}`
1341 : step.run_sh;
1342
1343 const stepTimeout =
1344 step.timeout ?? cfg.timeout ?? config.CI_DEFAULT_TIMEOUT;
1345 const timeoutSignal = AbortSignal.timeout(stepTimeout * 1000);
1346 // Abort the exec stream on either a run cancellation or the
1347 // per-step timeout. Docker has no per-exec kill, so on abort we
1348 // force-remove the container (below), which kills the command
1349 // still running inside it.
1350 const stepSignal = AbortSignal.any([signal, timeoutSignal]);
1351
1352 try {
1353 const { log, exitCode } = await execInContainer(
1354 containerId,
1355 [...shell, command],
1356 cfg.work_dir,
1357 envArray,
1358 stepSignal,
1359 async (partial) => {
1360 await db
1361 .updateTable("ci_steps")
1362 .set({
1363 log: maskSecrets(partial, secretValues),
1364 })
1365 .where("id", "=", stepId)
1366 .execute();
1367 },
1368 );
1369 stepLog = maskSecrets(log, secretValues);
1370 if (exitCode !== 0) stepStatus = onFail;
1371 } catch (err) {
1372 // A run cancellation is reported as "cancelled" by the outer
1373 // catch — don't relabel it as a step failure here.
1374 if (signal.aborted) throw err;
1375 if (timeoutSignal.aborted) {
1376 // Not subject to warn_on_fail. The container is about
1377 // to be destroyed, so every later step is skipped no
1378 // matter what. A run that cannot continue is a failure.
1379 stepStatus = "failure";
1380 stepLog = `Step timed out after ${stepTimeout}s\n`;
1381 // Kill the container now so the timed-out command stops
1382 // immediately rather than lingering until cleanup.
1383 await removeContainer(containerId);
1384 containerId = "";
1385 } else {
1386 stepStatus = onFail;
1387 stepLog = `Step failed: ${err instanceof Error ? err.message : String(err)}\n`;
1388 }
1389 }
1390 }
1391
1392 if (stepStatus === "failure") runFailed = true;
1393 if (stepStatus === "warning") sawWarning = true;
1394
1395 // Collect artifacts for this step
1396 if (stepStatus !== "failure") {
1397 await collectArtifacts(
1398 runId,
1399 containerId,
1400 step,
1401 shell,
1402 envArray,
1403 ).catch(() => {});
1404 }
1405
1406 await db
1407 .updateTable("ci_steps")
1408 .set({
1409 status: stepStatus,
1410 finished_at: now(),
1411 log: stepLog,
1412 })
1413 .where("id", "=", stepId)
1414 .execute();
1415 }
1416
1417 // Mark remaining steps as skipped
1418 await db
1419 .updateTable("ci_steps")
1420 .set({
1421 status: "skipped",
1422 started_at: now(),
1423 finished_at: now(),
1424 log: "Skipped: previous step failed",
1425 })
1426 .where("run_id", "=", runId)
1427 .where("status", "=", "pending")
1428 .execute();
1429
1430 const finalStatus = runFailed
1431 ? "failure"
1432 : sawWarning
1433 ? "warning"
1434 : "success";
1435 await db
1436 .updateTable("ci_runs")
1437 .set({ status: finalStatus, finished_at: now() })
1438 .where("id", "=", runId)
1439 .execute();
1440 } catch (err) {
1441 const errMsg = err instanceof Error ? err.message : String(err);
1442 const isDockerUnavailable = errMsg.includes("No Docker/Podman socket");
1443 const status = signal.aborted
1444 ? "cancelled"
1445 : isDockerUnavailable
1446 ? "skipped"
1447 : "failure";
1448 const skipLog = signal.aborted
1449 ? "Skipped: run was cancelled"
1450 : isDockerUnavailable
1451 ? "Skipped: Docker/Podman not available"
1452 : "Skipped: run failed";
1453
1454 // Only while setup still owns the run. Once it succeeds, a failure
1455 // belongs to a step, and that step carries its own log.
1456 await db
1457 .updateTable("ci_steps")
1458 .set({
1459 status: signal.aborted
1460 ? "cancelled"
1461 : isDockerUnavailable
1462 ? "skipped"
1463 : "failure",
1464 finished_at: new Date().toISOString(),
1465 log: isDockerUnavailable ? skipLog : `Error: ${errMsg}\n`,
1466 })
1467 .where("id", "=", setupStepId)
1468 .where("status", "=", "running")
1469 .execute();
1470 await db
1471 .updateTable("ci_runs")
1472 .set({ status, finished_at: new Date().toISOString() })
1473 .where("id", "=", runId)
1474 .execute();
1475 // Mark pending steps as skipped
1476 await db
1477 .updateTable("ci_steps")
1478 .set({
1479 status: "skipped",
1480 finished_at: new Date().toISOString(),
1481 log: skipLog,
1482 })
1483 .where("run_id", "=", runId)
1484 .where("status", "=", "pending")
1485 .execute();
1486 } finally {
1487 if (containerId) await removeContainer(containerId);
1488 runningTasks.delete(runId);
1489 // Prune old history
1490 const run = await db
1491 .selectFrom("ci_runs")
1492 .select("repo_id")
1493 .where("id", "=", runId)
1494 .executeTakeFirst();
1495 if (run) await pruneHistory(run.repo_id).catch(() => {});
1496 if (pruneCaches) {
1497 await pruneStaleCaches(
1498 pruneCaches.repoName,
1499 pruneCaches.cache,
1500 ).catch(() => {});
1501 }
1502 // After removeContainer above: the daemon will not delete a volume
1503 // that is still mounted.
1504 if (cacheLimits) {
1505 const dropped = await enforceCacheLimits(
1506 cacheLimits.repoName,
1507 cacheLimits.cache,
1508 ).catch(() => []);
1509 if (dropped.length) {
1510 // Recorded on the run, because the only other symptom is the
1511 // next build being mysteriously slow.
1512 await db
1513 .insertInto("ci_steps")
1514 .values({
1515 run_id: runId,
1516 name: "cache",
1517 status: "success",
1518 started_at: new Date().toISOString(),
1519 finished_at: new Date().toISOString(),
1520 log: `${dropped
1521 .map(
1522 (d) =>
1523 `Dropped cache ${d.path}: ${formatBytes(d.size)} over the ${formatBytes(d.maxSize)} limit`,
1524 )
1525 .join("\n")}\n`,
1526 })
1527 .execute()
1528 .catch(() => {});
1529 }
1530 }
1531 // A slot just freed up — start the next queued run if any.
1532 pumpQueue();
1533 }
1534}
1535
1536// --- Public API ---
1537
1538export async function triggerRun(
1539 repoName: string,
1540 opts: TriggerOpts,
1541): Promise<number> {
1542 const repo = await db
1543 .selectFrom("repositories")
1544 .select("id")
1545 .where("name", "=", repoName)
1546 .executeTakeFirst();
1547 if (!repo) throw new Error("Repository not found");
1548
1549 const runId = await db
1550 .insertInto("ci_runs")
1551 .values({
1552 repo_id: repo.id,
1553 triggered_by: opts.triggeredBy ?? null,
1554 trigger_source: opts.triggerSource,
1555 commit_sha: opts.commitSha,
1556 commit_branch: opts.commitBranch ?? null,
1557 commit_tag: opts.commitTag ?? null,
1558 status: "pending",
1559 variable_overrides: opts.variableOverrides
1560 ? JSON.stringify(opts.variableOverrides)
1561 : null,
1562 })
1563 .returning("id")
1564 .executeTakeFirstOrThrow();
1565
1566 // Allocate the human-facing run number atomically from a per-repo counter.
1567 // A single upsert-and-increment can't collide under concurrent triggers and
1568 // never reuses a number after pruneHistory shrinks the run table — both of
1569 // which a COUNT(*)-based scheme suffered from.
1570 const counter = await db
1571 .insertInto("ci_run_counters")
1572 .values({ repo_id: repo.id, last_run_id: 1 })
1573 .onConflict((oc) =>
1574 oc.column("repo_id").doUpdateSet((eb) => ({
1575 last_run_id: eb("ci_run_counters.last_run_id", "+", 1),
1576 })),
1577 )
1578 .returning("last_run_id")
1579 .executeTakeFirstOrThrow();
1580 await db
1581 .updateTable("ci_runs")
1582 .set({ repo_run_id: counter.last_run_id })
1583 .where("id", "=", runId.id)
1584 .execute();
1585
1586 // Flip to "queued" BEFORE pushing onto queuedRunIds. Otherwise a
1587 // concurrent pumpQueue (from a finishing run) could observe our entry,
1588 // promoteToPending it, and start executeRun while our UPDATE is still
1589 // in flight — the late UPDATE would then clobber a running row back to
1590 // "queued". Only pumpQueue (the sole writer of runningTasks.set) may
1591 // mutate the row's status after the push. The TOCTOU concern is moot
1592 // here because pumpQueue cannot pump a runId that hasn't been pushed.
1593 if (runningTasks.size >= config.CI_MAX_CONCURRENT) {
1594 await db
1595 .updateTable("ci_runs")
1596 .set({ status: "queued" })
1597 .where("id", "=", runId.id)
1598 .execute();
1599 }
1600 queuedRunIds.push(runId.id);
1601 pumpQueue();
1602
1603 return runId.id;
1604}
1605
1606// Wraps the fire-and-forget executeRun so unhandled exceptions (e.g. a throw
1607// before/after its own try/finally) are logged and the run is reconciled to
1608// "failure" instead of staying pending forever.
1609function spawnRun(
1610 runId: number,
1611 signal: AbortSignal,
1612 needsPromote = false,
1613): void {
1614 void (async () => {
1615 try {
1616 if (needsPromote) await promoteToPending(runId);
1617 await executeRun(runId, signal);
1618 } catch (err) {
1619 console.error(`[ci] executeRun threw for run ${runId}:`, err);
1620 runningTasks.delete(runId);
1621 try {
1622 await db
1623 .updateTable("ci_runs")
1624 .set({
1625 status: "failure",
1626 finished_at: new Date().toISOString(),
1627 })
1628 .where("id", "=", runId)
1629 // Include "queued" so a row that was promoted-but-not-yet-
1630 // observed (or never promoted because promoteToPending threw)
1631 // still gets marked failed instead of stuck.
1632 .where("status", "in", ["pending", "running", "queued"])
1633 .execute();
1634 } catch (dbErr) {
1635 console.error(
1636 `[ci] failed to mark run ${runId} failed:`,
1637 dbErr,
1638 );
1639 }
1640 pumpQueue();
1641 }
1642 })();
1643}
1644
1645export async function retryRun(
1646 runId: number,
1647 retriedBy: number,
1648): Promise<void> {
1649 const run = await db
1650 .selectFrom("ci_runs")
1651 .select("repo_id")
1652 .where("id", "=", runId)
1653 .executeTakeFirst();
1654 if (!run) throw new Error("Run not found");
1655
1656 // Delete existing steps
1657 await db.deleteFrom("ci_steps").where("run_id", "=", runId).execute();
1658
1659 // Delete artifacts from disk and DB
1660 const artifactDir = path.join(paths.CI_ARTIFACTS_DIR, String(runId));
1661 if (existsSync(artifactDir)) {
1662 await Bun.$`rm -rf ${artifactDir}`.quiet().nothrow();
1663 }
1664 await db.deleteFrom("ci_artifacts").where("run_id", "=", runId).execute();
1665
1666 // Reset to "pending" first; if a slot isn't free, flip to "queued"
1667 // before enqueueing. Mirrors triggerRun: pumpQueue is the sole writer
1668 // of runningTasks.set, so we never reserve a slot directly here. Two
1669 // concurrent retryRun calls (or a retryRun racing triggerRun) all
1670 // funnel through pumpQueue, which serializes the size check against
1671 // slot reservation in a single synchronous turn.
1672 await db
1673 .updateTable("ci_runs")
1674 .set({
1675 status: "pending",
1676 triggered_by: retriedBy,
1677 started_at: null,
1678 finished_at: null,
1679 })
1680 .where("id", "=", runId)
1681 .execute();
1682
1683 if (runningTasks.size >= config.CI_MAX_CONCURRENT) {
1684 await db
1685 .updateTable("ci_runs")
1686 .set({ status: "queued" })
1687 .where("id", "=", runId)
1688 .execute();
1689 }
1690 queuedRunIds.push(runId);
1691 pumpQueue();
1692}
1693
1694export async function cancelRun(runId: number): Promise<void> {
1695 const task = runningTasks.get(runId);
1696 if (task) {
1697 const { containerId } = task;
1698 task.controller.abort();
1699 if (containerId) {
1700 await removeContainer(containerId).catch(() => {});
1701 }
1702 }
1703 const queueIdx = queuedRunIds.indexOf(runId);
1704 if (queueIdx >= 0) queuedRunIds.splice(queueIdx, 1);
1705 await db
1706 .updateTable("ci_runs")
1707 .set({ status: "cancelled", finished_at: new Date().toISOString() })
1708 .where("id", "=", runId)
1709 .where("status", "in", ["pending", "running", "queued"])
1710 .execute();
1711}
1712
1713async function pruneHistory(repoId: number): Promise<void> {
1714 const maxHistory = config.CI_MAX_HISTORY;
1715 const allRuns = await db
1716 .selectFrom("ci_runs")
1717 .select("id")
1718 .where("repo_id", "=", repoId)
1719 .orderBy("id", "desc")
1720 .execute();
1721
1722 if (allRuns.length <= maxHistory) return;
1723
1724 const toDelete = allRuns.slice(maxHistory).map((r) => r.id);
1725 for (const runId of toDelete) {
1726 // Remove artifacts from disk
1727 const artifactDir = path.join(paths.CI_ARTIFACTS_DIR, String(runId));
1728 if (existsSync(artifactDir)) {
1729 await Bun.$`rm -rf ${artifactDir}`.quiet().nothrow();
1730 }
1731 }
1732 await db.deleteFrom("ci_runs").where("id", "in", toDelete).execute();
1733}
1734
1735export async function purgeRepoCaches(repoName: string): Promise<void> {
1736 try {
1737 const filters = encodeURIComponent(
1738 JSON.stringify({ label: [`com.hearthforge.repo=${repoName}`] }),
1739 );
1740 const resp = await dockerFetch(`/volumes?filters=${filters}`);
1741 if (!resp.ok) {
1742 await resp.body?.cancel();
1743 return;
1744 }
1745 const data = (await resp.json()) as {
1746 Volumes?: Array<{ Name: string }>;
1747 };
1748 for (const vol of data.Volumes ?? []) {
1749 const delResp = await dockerFetch(`/volumes/${vol.Name}`, {
1750 method: "DELETE",
1751 });
1752 await delResp.body?.cancel();
1753 }
1754 } catch {
1755 // Best-effort
1756 }
1757}
1758
1759export async function cancelStaleRuns(): Promise<void> {
1760 const now = new Date().toISOString();
1761
1762 const stale = await db
1763 .selectFrom("ci_runs")
1764 .select("id")
1765 .where("status", "in", ["pending", "running", "queued"])
1766 .execute();
1767
1768 await Promise.allSettled(
1769 stale.map((r) =>
1770 dockerFetch(`/containers/hearthforge-ci-${r.id}?force=true`, {
1771 method: "DELETE",
1772 }).then((res) => res.body?.cancel()),
1773 ),
1774 );
1775
1776 await db
1777 .updateTable("ci_runs")
1778 .set({ status: "cancelled", finished_at: now })
1779 .where("status", "in", ["pending", "running", "queued"])
1780 .execute();
1781 await db
1782 .updateTable("ci_steps")
1783 .set({
1784 status: "cancelled",
1785 finished_at: now,
1786 log: "Skipped: run was cancelled",
1787 })
1788 .where("status", "in", ["pending", "running"])
1789 .execute();
1790}
1791
1792/** Reset the cached socket path (used in tests to switch mock sockets). */
1793export function resetDockerSocket(): void {
1794 resolvedSocket = null;
1795}
1796