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