ci.ts
⎇
Raw
1import { existsSync, mkdirSync, writeFileSync } from "node:fs";
2import path from "node:path";
3import config from "../config.ts";
4import { CI_MAX_LOG_BYTES, paths } from "../constants.ts";
5import { db } from "../db/index.ts";
6import { repoPath } from "./git.ts";
7
8// --- Types ---
9
10interface CiVariableDef {
11 default?: string;
12 description?: string;
13}
14
15interface CiStepConfig {
16 run_sh?: string;
17 run_if?: string;
18 clear?: boolean;
19 timeout?: number;
20 publish_file?: string | string[];
21 publish_tar?: string | string[];
22 publish_gzip?: string | string[];
23 publish_zip?: string | string[];
24 publish_zstd?: string | string[];
25}
26
27export interface CiStep extends CiStepConfig {
28 name: string;
29}
30
31export interface CiConfig {
32 image: string;
33 work_dir?: string;
34 clone_project_to?: string;
35 shell?: string[];
36 shell_setup?: string;
37 timeout?: number;
38 cpu_limit?: number;
39 memory_limit?: string;
40 cache?: string[];
41 on?: {
42 push?: string[] | boolean;
43 tag?: boolean;
44 manual?: boolean;
45 };
46 variables?: Record<string, CiVariableDef>;
47 steps: CiStep[];
48}
49
50export interface TriggerOpts {
51 triggerSource: "push" | "tag" | "manual";
52 commitSha: string;
53 commitBranch?: string;
54 commitTag?: string;
55 triggeredBy?: number;
56 variableOverrides?: Record<string, string>;
57}
58
59// Reserved TOML table names that are not steps
60const RESERVED_TABLES = new Set(["on", "variables"]);
61
62const CONTAINER_REPO_PATH = "/hearthforge-repo.git";
63
64// In-memory map of running tasks for cancellation
65const runningTasks = new Map<
66 number,
67 { controller: AbortController; containerId?: string }
68>();
69
70// FIFO queue of run IDs whose DB row is `status = "queued"`. We start the
71// next one in `pumpQueue` whenever `runningTasks.size` drops below
72// `CI_MAX_CONCURRENT`. `pumpQueue` is the SOLE writer of `runningTasks.set`
73// — callers that want to start a run push to `queuedRunIds` and call
74// `pumpQueue` synchronously. This makes the size check + slot reservation
75// atomic against concurrent callers (no await between).
76const queuedRunIds: number[] = [];
77
78async function promoteToPending(runId: number): Promise<void> {
79 await db
80 .updateTable("ci_runs")
81 .set({ status: "pending" })
82 .where("id", "=", runId)
83 .execute();
84}
85
86function pumpQueue(): void {
87 while (
88 queuedRunIds.length > 0 &&
89 runningTasks.size < config.CI_MAX_CONCURRENT
90 ) {
91 const next = queuedRunIds.shift()!;
92 const controller = new AbortController();
93 runningTasks.set(next, { controller });
94 // Promote queued→pending before spawning so executeRun's failure path
95 // (which only matches pending/running) can still mark it failed if
96 // it throws very early. We await the flip inside spawnRun's wrapper
97 // so the order is: row=pending → executeRun starts → row=running.
98 spawnRun(next, controller.signal, /* needsPromote */ true);
99 }
100}
101
102/** Position (1-based) of this queued run within the queue, or null. */
103export function ciQueuePosition(runId: number): number | null {
104 const idx = queuedRunIds.indexOf(runId);
105 return idx < 0 ? null : idx + 1;
106}
107
108// --- TOML Parsing ---
109
110export function parseCiConfig(tomlStr: string): CiConfig | null {
111 let raw: Record<string, unknown>;
112 try {
113 raw = Bun.TOML.parse(tomlStr) as Record<string, unknown>;
114 } catch {
115 return null;
116 }
117
118 const image = raw.image;
119 if (typeof image !== "string" || !image) return null;
120
121 const steps: CiStep[] = [];
122 for (const [key, val] of Object.entries(raw)) {
123 if (RESERVED_TABLES.has(key)) continue;
124 if (typeof val !== "object" || val === null || Array.isArray(val))
125 continue;
126 // It's a table section — treat as a step
127 const stepCfg = val as Record<string, unknown>;
128 steps.push({ name: key, ...(stepCfg as CiStepConfig) });
129 }
130
131 const rawOn = raw.on as Record<string, unknown> | undefined;
132 const rawVars = raw.variables as
133 | Record<string, Record<string, unknown>>
134 | undefined;
135 const variables: Record<string, CiVariableDef> = {};
136 if (rawVars) {
137 for (const [name, def] of Object.entries(rawVars)) {
138 if (typeof def === "object" && def !== null) {
139 variables[name] = {
140 default:
141 typeof def.default === "string"
142 ? def.default
143 : undefined,
144 description:
145 typeof def.description === "string"
146 ? def.description
147 : undefined,
148 };
149 }
150 }
151 }
152
153 return {
154 image,
155 work_dir: typeof raw.work_dir === "string" ? raw.work_dir : undefined,
156 clone_project_to:
157 typeof raw.clone_project_to === "string"
158 ? raw.clone_project_to
159 : undefined,
160 shell: Array.isArray(raw.shell) ? (raw.shell as string[]) : undefined,
161 shell_setup:
162 typeof raw.shell_setup === "string" ? raw.shell_setup : undefined,
163 timeout: typeof raw.timeout === "number" ? raw.timeout : undefined,
164 cpu_limit:
165 typeof raw.cpu_limit === "number" ? raw.cpu_limit : undefined,
166 memory_limit:
167 typeof raw.memory_limit === "string" ? raw.memory_limit : undefined,
168 cache: Array.isArray(raw.cache) ? (raw.cache as string[]) : undefined,
169 on: rawOn
170 ? {
171 push: Array.isArray(rawOn.push)
172 ? (rawOn.push as string[])
173 : typeof rawOn.push === "boolean"
174 ? rawOn.push
175 : undefined,
176 tag: typeof rawOn.tag === "boolean" ? rawOn.tag : undefined,
177 manual:
178 typeof rawOn.manual === "boolean"
179 ? rawOn.manual
180 : undefined,
181 }
182 : undefined,
183 variables,
184 steps,
185 };
186}
187
188/**
189 * Run-start checks kept out of parseCiConfig, which can only report one
190 * generic failure. Returns an error string, or null when the config is usable.
191 */
192export function validateCiConfig(cfg: CiConfig): string | null {
193 if (!cfg.clone_project_to || !cfg.cache) return null;
194
195 const clone = path.resolve(cfg.clone_project_to);
196 for (const entry of cfg.cache) {
197 const cachePath = path.resolve(entry);
198 if (
199 cachePath === clone ||
200 cachePath.startsWith(`${clone}${path.sep}`)
201 ) {
202 return (
203 `cache path "${entry}" is inside clone_project_to ("${cfg.clone_project_to}"). ` +
204 "The cache volume is mounted before the clone runs, which leaves the " +
205 "directory non-empty, and git refuses to clone into it. Move the cache " +
206 "outside the clone directory."
207 );
208 }
209 }
210 return null;
211}
212
213// --- Trigger matching ---
214
215export function shouldTriggerPush(cfg: CiConfig, branch: string): boolean {
216 const pushCfg = cfg.on?.push;
217 if (!pushCfg) return false;
218 if (pushCfg === true) return true;
219 if (Array.isArray(pushCfg)) {
220 return pushCfg.some((pattern) => matchGlob(pattern, branch));
221 }
222 return false;
223}
224
225export function shouldTriggerTag(cfg: CiConfig): boolean {
226 return cfg.on?.tag === true;
227}
228
229function matchGlob(pattern: string, value: string): boolean {
230 if (pattern === "*") return true;
231 const re = new RegExp(
232 `^${pattern.replace(/[.+^${}()|[\]\\]/g, "\\$&").replace(/\*/g, ".*")}$`,
233 );
234 return re.test(value);
235}
236
237// --- Docker socket ---
238
239let resolvedSocket: string | null = null;
240
241async function getSocket(): Promise<string> {
242 if (resolvedSocket) return resolvedSocket;
243 const candidates = config.CI_DOCKER_SOCKET
244 ? [config.CI_DOCKER_SOCKET]
245 : (() => {
246 const uid = process.getuid?.();
247 return [
248 "/var/run/docker.sock",
249 "/run/podman/podman.sock",
250 ...(uid !== undefined
251 ? [`/run/user/${uid}/podman/podman.sock`]
252 : []),
253 ];
254 })();
255 for (const s of candidates) {
256 if (existsSync(s)) {
257 resolvedSocket = s;
258 return s;
259 }
260 }
261 throw new Error(
262 "No Docker/Podman socket found. Set CI_DOCKER_SOCKET env var.",
263 );
264}
265
266async function dockerFetch(
267 endpoint: string,
268 init?: RequestInit,
269): Promise<Response> {
270 const socket = await getSocket();
271 return fetch(`http://localhost/v1.47${endpoint}`, {
272 ...init,
273 unix: socket,
274 });
275}
276
277// --- Docker helpers ---
278
279function splitImageRef(image: string): { name: string; tag: string } {
280 const lastColon = image.lastIndexOf(":");
281 if (lastColon < 0) return { name: image, tag: "latest" };
282 const possibleTag = image.slice(lastColon + 1);
283 if (possibleTag.includes("/")) return { name: image, tag: "latest" };
284 return { name: image.slice(0, lastColon), tag: possibleTag };
285}
286
287async function pullImage(image: string): Promise<void> {
288 const { name, tag } = splitImageRef(image);
289 const resp = await dockerFetch(
290 `/images/create?fromImage=${encodeURIComponent(name)}&tag=${encodeURIComponent(tag)}`,
291 { method: "POST" },
292 );
293 if (!resp.ok) {
294 const detail = (await resp.text().catch(() => "")).trim();
295 throw new Error(
296 `Failed to pull image ${image}: HTTP ${resp.status}${detail ? ` ${detail}` : ""}`,
297 );
298 }
299 // The /images/create stream must be read to the end — the pull only
300 // completes when the stream does. Cancelling it (the previous behavior)
301 // aborted the pull, so createContainer could race a not-yet-present image.
302 // Each line is a JSON progress object; a trailing {"error": …} means the
303 // pull failed despite the HTTP 200.
304 const body = await resp.text();
305 for (const line of body.split("\n")) {
306 const trimmed = line.trim();
307 if (!trimmed) continue;
308 let obj: { error?: string } | null = null;
309 try {
310 obj = JSON.parse(trimmed);
311 } catch {
312 continue; // non-JSON progress line — ignore
313 }
314 if (obj?.error) {
315 throw new Error(`Failed to pull image ${image}: ${obj.error}`);
316 }
317 }
318}
319
320function parseMemoryBytes(s: string): number {
321 const m = s.match(/^(\d+(?:\.\d+)?)\s*([kmgKMG]?)b?$/);
322 if (!m) return 0;
323 const n = parseFloat(m[1] ?? "0");
324 switch ((m[2] ?? "").toLowerCase()) {
325 case "k":
326 return Math.floor(n * 1024);
327 case "m":
328 return Math.floor(n * 1024 * 1024);
329 case "g":
330 return Math.floor(n * 1024 * 1024 * 1024);
331 default:
332 return Math.floor(n);
333 }
334}
335
336async function ensureVolume(volName: string, repoName: string): Promise<void> {
337 const resp = await dockerFetch("/volumes/create", {
338 method: "POST",
339 headers: { "Content-Type": "application/json" },
340 body: JSON.stringify({
341 Name: volName,
342 Labels: { "com.hearthforge.repo": repoName },
343 }),
344 });
345 await resp.body?.cancel();
346}
347
348async function createContainer(
349 runId: number,
350 repoName: string,
351 cfg: CiConfig,
352 envVars: string[],
353): Promise<string> {
354 // No bind for the repo: the Docker daemon resolves bind sources on the
355 // host, where HearthForge's own paths do not exist. uploadRepo copies it in.
356 const binds: string[] = [];
357 if (cfg.cache) {
358 for (const cachePath of cfg.cache) {
359 const volName = `hearthforge-ci-cache-${Buffer.from(`${repoName}:${cachePath}`).toString("base64url").slice(0, 24)}`;
360 await ensureVolume(volName, repoName);
361 binds.push(`${volName}:${cachePath}`);
362 }
363 }
364
365 const hostConfig: Record<string, unknown> = { Binds: binds };
366 if (cfg.cpu_limit) {
367 hostConfig.NanoCpus = Math.floor(cfg.cpu_limit * 1e9);
368 }
369 if (cfg.memory_limit) {
370 hostConfig.Memory = parseMemoryBytes(cfg.memory_limit);
371 }
372
373 const body = JSON.stringify({
374 Image: cfg.image,
375 Cmd: ["sleep", "infinity"],
376 Env: envVars,
377 WorkingDir: cfg.work_dir ?? "/",
378 HostConfig: hostConfig,
379 });
380
381 const resp = await dockerFetch(
382 `/containers/create?name=hearthforge-ci-${runId}`,
383 {
384 method: "POST",
385 headers: { "Content-Type": "application/json" },
386 body,
387 },
388 );
389 if (!resp.ok) {
390 const text = await resp.text();
391 throw new Error(`Failed to create container: ${resp.status} ${text}`);
392 }
393 const data = (await resp.json()) as { Id: string };
394 return data.Id;
395}
396
397async function startContainer(containerId: string): Promise<void> {
398 const resp = await dockerFetch(`/containers/${containerId}/start`, {
399 method: "POST",
400 });
401 if (!resp.ok && resp.status !== 304) {
402 throw new Error(`Failed to start container: ${resp.status}`);
403 }
404 await resp.body?.cancel();
405}
406
407interface ExecResult {
408 log: string;
409 exitCode: number;
410}
411
412const dec = new TextDecoder();
413
414function parseMuxFrames(buf: Uint8Array): {
415 text: string;
416 remaining: Uint8Array<ArrayBuffer>;
417} {
418 const chunks: string[] = [];
419 let i = 0;
420 while (i + 8 <= buf.length) {
421 const view = new DataView(buf.buffer, buf.byteOffset + i, 8);
422 const size = view.getUint32(4, false);
423 if (i + 8 + size > buf.length) break;
424 chunks.push(dec.decode(buf.slice(i + 8, i + 8 + size)));
425 i += 8 + size;
426 }
427 const remaining = new Uint8Array(buf.length - i);
428 if (i < buf.length) remaining.set(buf.subarray(i));
429 return { text: chunks.join(""), remaining };
430}
431
432async function execInContainer(
433 containerId: string,
434 cmd: string[],
435 workDir?: string,
436 envVars?: string[],
437 signal?: AbortSignal,
438 onPartialLog?: (log: string) => Promise<void>,
439): Promise<ExecResult> {
440 // Create exec
441 const execBody = JSON.stringify({
442 Cmd: cmd,
443 AttachStdout: true,
444 AttachStderr: true,
445 ...(workDir ? { WorkingDir: workDir } : {}),
446 ...(envVars ? { Env: envVars } : {}),
447 });
448 const createResp = await dockerFetch(`/containers/${containerId}/exec`, {
449 method: "POST",
450 headers: { "Content-Type": "application/json" },
451 body: execBody,
452 signal,
453 });
454 if (!createResp.ok) {
455 const text = await createResp.text();
456 throw new Error(`Failed to create exec: ${createResp.status} ${text}`);
457 }
458 const createData = (await createResp.json()) as { Id: string };
459 const execId = createData.Id;
460
461 // Start exec and stream output
462 const startResp = await dockerFetch(`/exec/${execId}/start`, {
463 method: "POST",
464 headers: { "Content-Type": "application/json" },
465 body: JSON.stringify({ Detach: false, Tty: false }),
466 signal,
467 });
468
469 let log = "";
470 let truncated = false;
471 // Bound the buffered log: stop appending once we hit the cap (and note it
472 // once) so a runaway step can't exhaust RAM or make each partial-flush
473 // rewrite an ever-growing row.
474 const appendLog = (text: string) => {
475 if (truncated || !text) return;
476 const room = CI_MAX_LOG_BYTES - log.length;
477 if (text.length <= room) {
478 log += text;
479 } else {
480 log += text.slice(0, Math.max(0, room));
481 log += `\n[log truncated at ${CI_MAX_LOG_BYTES} bytes]\n`;
482 truncated = true;
483 }
484 };
485
486 if (onPartialLog && startResp.body) {
487 const reader = startResp.body.getReader();
488 let buf = new Uint8Array(0);
489 let lastSave = Date.now();
490 while (true) {
491 const { done, value } = await reader.read();
492 if (done) break;
493 const merged = new Uint8Array(buf.length + value.length);
494 merged.set(buf);
495 merged.set(value, buf.length);
496 buf = merged;
497 const { text, remaining } = parseMuxFrames(buf);
498 buf = remaining;
499 appendLog(text);
500 // Once truncated the log no longer changes, so stop re-flushing it.
501 if (!truncated && Date.now() - lastSave >= 2000) {
502 await onPartialLog(log);
503 lastSave = Date.now();
504 }
505 }
506 const { text } = parseMuxFrames(buf);
507 appendLog(text);
508 } else {
509 const bodyBytes = new Uint8Array(await startResp.arrayBuffer());
510 const { text } = parseMuxFrames(bodyBytes);
511 appendLog(text);
512 }
513
514 // Get exit code
515 const inspectResp = await dockerFetch(`/exec/${execId}/json`);
516 const inspectData = (await inspectResp.json()) as { ExitCode: number };
517
518 return { log, exitCode: inspectData.ExitCode ?? 1 };
519}
520
521/**
522 * Copy the bare repo into the container at CONTAINER_REPO_PATH.
523 *
524 * Not a bind mount: the Docker daemon resolves bind sources in the host
525 * filesystem, so when HearthForge itself runs in a container it finds nothing
526 * at our DATA_DIR path and silently mounts an empty directory instead.
527 */
528async function uploadRepo(
529 containerId: string,
530 repoAbsPath: string,
531): Promise<void> {
532 const mk = await execInContainer(containerId, [
533 "mkdir",
534 "-p",
535 CONTAINER_REPO_PATH,
536 ]);
537 if (mk.exitCode !== 0) {
538 throw new Error(
539 `Failed to create ${CONTAINER_REPO_PATH} in the container: ${mk.log.trim()}`,
540 );
541 }
542
543 // Flatten ownership: host uids mean nothing here and trip git's
544 // ownership check on the clone below.
545 const tar = Bun.spawn(
546 [
547 "tar",
548 "-cf",
549 "-",
550 "--owner=0",
551 "--group=0",
552 "--numeric-owner",
553 "-C",
554 repoAbsPath,
555 ".",
556 ],
557 { stdout: "pipe", stderr: "pipe" },
558 );
559
560 const resp = await dockerFetch(
561 `/containers/${containerId}/archive?path=${encodeURIComponent(CONTAINER_REPO_PATH)}`,
562 {
563 method: "PUT",
564 headers: { "Content-Type": "application/x-tar" },
565 body: tar.stdout,
566 },
567 );
568 const tarExit = await tar.exited;
569
570 if (!resp.ok) {
571 const detail = (await resp.text().catch(() => "")).trim();
572 throw new Error(
573 `Failed to upload the repository: HTTP ${resp.status}${detail ? ` ${detail}` : ""}`,
574 );
575 }
576 await resp.body?.cancel();
577 if (tarExit !== 0) {
578 const err = (await new Response(tar.stderr).text()).trim();
579 throw new Error(
580 `Failed to read the repository: tar exited ${tarExit}${err ? ` ${err}` : ""}`,
581 );
582 }
583}
584
585async function removeContainer(containerId: string): Promise<void> {
586 try {
587 const resp = await dockerFetch(
588 `/containers/${containerId}?force=true`,
589 { method: "DELETE" },
590 );
591 await resp.body?.cancel();
592 } catch {
593 // Best-effort cleanup
594 }
595}
596
597// --- Tar extraction ---
598
599function extractSingleFileFromTar(data: Uint8Array): Uint8Array | null {
600 if (data.length < 512) return null;
601 const dec = new TextDecoder();
602 const sizeOctal = dec
603 .decode(data.slice(124, 136))
604 .replace(/\0/g, "")
605 .trim();
606 const size = parseInt(sizeOctal, 8);
607 if (Number.isNaN(size) || size < 0) return null;
608 if (data.length < 512 + size) return null;
609 return data.slice(512, 512 + size);
610}
611
612async function copyFileFromContainer(
613 containerId: string,
614 containerPath: string,
615): Promise<Uint8Array | null> {
616 const resp = await dockerFetch(
617 `/containers/${containerId}/archive?path=${encodeURIComponent(containerPath)}`,
618 );
619 if (!resp.ok) return null;
620 const tarBytes = new Uint8Array(await resp.arrayBuffer());
621 return extractSingleFileFromTar(tarBytes);
622}
623
624// --- Secret masking ---
625
626function maskSecrets(text: string, secrets: string[]): string {
627 for (const secret of secrets) {
628 if (secret) text = text.split(secret).join("[MASKED]");
629 }
630 return text;
631}
632
633// --- Env var building ---
634
635function buildEnvVars(
636 runId: number,
637 repoName: string,
638 opts: TriggerOpts,
639 cfg: CiConfig,
640 secretValues: Array<{ name: string; value: string }>,
641): { envArray: string[]; secretValues: string[] } {
642 const vars: Record<string, string> = {
643 CI: "true",
644 CI_PIPELINE_ID: String(runId),
645 CI_REPO_NAME: repoName,
646 CI_SERVER_URL: config.BASE_URL,
647 CI_TRIGGER_SOURCE: opts.triggerSource,
648 CI_COMMIT_SHA: opts.commitSha,
649 CI_COMMIT_SHORT_SHA: opts.commitSha.slice(0, 8),
650 CI_COMMIT_BRANCH: opts.commitBranch ?? "",
651 CI_COMMIT_TAG: opts.commitTag ?? "",
652 CI_COMMIT_REF_NAME: opts.commitTag ?? opts.commitBranch ?? "",
653 };
654
655 // User-defined variable defaults
656 if (cfg.variables) {
657 for (const [name, def] of Object.entries(cfg.variables)) {
658 if (def.default !== undefined) vars[name] = def.default;
659 }
660 }
661
662 // Variable overrides from manual trigger
663 if (opts.variableOverrides) {
664 for (const [name, value] of Object.entries(opts.variableOverrides)) {
665 vars[name] = value;
666 }
667 }
668
669 // Secrets (injected but values tracked for masking)
670 const secretVals: string[] = [];
671 for (const { name, value } of secretValues) {
672 vars[name] = value;
673 secretVals.push(value);
674 }
675
676 const envArray = Object.entries(vars).map(([k, v]) => `${k}=${v}`);
677 return { envArray, secretValues: secretVals };
678}
679
680// --- Artifact collection ---
681
682async function collectArtifacts(
683 runId: number,
684 containerId: string,
685 step: CiStep,
686 _shell: string[],
687 envVars: string[],
688): Promise<void> {
689 const artifactDir = path.join(paths.CI_ARTIFACTS_DIR, String(runId));
690 mkdirSync(artifactDir, { recursive: true });
691
692 const toArray = (v: string | string[] | undefined): string[] => {
693 if (!v) return [];
694 return Array.isArray(v) ? v : [v];
695 };
696
697 // publish_file: copy directly out of container
698 for (const srcPath of toArray(step.publish_file)) {
699 const fileBytes = await copyFileFromContainer(containerId, srcPath);
700 if (fileBytes) {
701 const filename = path.basename(srcPath);
702 const destPath = path.join(artifactDir, filename);
703 writeFileSync(destPath, fileBytes);
704 const stat = Bun.file(destPath);
705 await db
706 .insertInto("ci_artifacts")
707 .values({
708 run_id: runId,
709 filename,
710 size: stat.size,
711 })
712 .execute();
713 }
714 }
715
716 // Archive commands run with `path.dirname(srcPath)` as the working
717 // directory and reference the source by its basename only, so the
718 // user-controlled path never appears as part of an interpolated shell
719 // string. Previously `publish_zip` used `sh -c "cd … && zip …"` with
720 // raw interpolation — a step author who could write the CI TOML
721 // could shell-inject through the source path. Today the only TOML
722 // author is the admin, but this removes the implicit assumption.
723 type ArchiveType = "tar" | "gzip" | "zip" | "zstd";
724 const archiveFormats: Array<{
725 type: ArchiveType;
726 paths: string[];
727 ext: string;
728 cmd: (basename: string, dst: string) => string[];
729 }> = [
730 {
731 type: "tar",
732 paths: toArray(step.publish_tar),
733 ext: ".tar",
734 cmd: (basename, dst) => ["tar", "-cf", dst, basename],
735 },
736 {
737 type: "gzip",
738 paths: toArray(step.publish_gzip),
739 ext: ".tar.gz",
740 cmd: (basename, dst) => ["tar", "-czf", dst, basename],
741 },
742 {
743 type: "zstd",
744 paths: toArray(step.publish_zstd),
745 ext: ".tar.zst",
746 cmd: (basename, dst) => ["tar", "--zstd", "-cf", dst, basename],
747 },
748 {
749 type: "zip",
750 paths: toArray(step.publish_zip),
751 ext: ".zip",
752 cmd: (basename, dst) => ["zip", "-r", dst, basename],
753 },
754 ];
755
756 let archiveIndex = 0;
757 for (const { paths: archivePaths, ext, cmd } of archiveFormats) {
758 for (const srcPath of archivePaths) {
759 archiveIndex++;
760 const tmpPath = `/tmp/hf-artifact-${runId}-${archiveIndex}${ext}`;
761 // Run with the source's parent directory as the working
762 // directory so each archive tool can reference the source
763 // by its basename — no -C, no shell.
764 const execResult = await execInContainer(
765 containerId,
766 cmd(path.basename(srcPath), tmpPath),
767 path.dirname(srcPath),
768 envVars,
769 ).catch(() => null);
770 if (!execResult || execResult.exitCode !== 0) continue;
771
772 // Copy archive out
773 const fileBytes = await copyFileFromContainer(containerId, tmpPath);
774 if (!fileBytes) continue;
775
776 const filename = `${path.basename(srcPath)}${ext}`;
777 const destPath = path.join(artifactDir, filename);
778 writeFileSync(destPath, fileBytes);
779 const stat = Bun.file(destPath);
780 await db
781 .insertInto("ci_artifacts")
782 .values({
783 run_id: runId,
784 filename,
785 size: stat.size,
786 })
787 .execute();
788 }
789 }
790}
791
792// --- Main execution ---
793
794async function executeRun(runId: number, signal: AbortSignal): Promise<void> {
795 const now = () => new Date().toISOString();
796 let containerId: string | undefined;
797
798 try {
799 // Mark as running
800 await db
801 .updateTable("ci_runs")
802 .set({ status: "running", started_at: now() })
803 .where("id", "=", runId)
804 .execute();
805
806 // Load run details
807 const run = await db
808 .selectFrom("ci_runs")
809 .selectAll()
810 .where("id", "=", runId)
811 .executeTakeFirst();
812 if (!run) throw new Error("Run not found");
813
814 const repo = await db
815 .selectFrom("repositories")
816 .select(["id", "name"])
817 .where("id", "=", run.repo_id)
818 .executeTakeFirst();
819 if (!repo) throw new Error("Repo not found");
820
821 // Read .hearthforge-ci.toml at the commit
822 const tomlBuf = await import("./git.ts").then((g) =>
823 g.git.show(repo.name, run.commit_sha!, ".hearthforge-ci.toml"),
824 );
825 if (!tomlBuf)
826 throw new Error(".hearthforge-ci.toml not found at commit");
827
828 const cfg = parseCiConfig(tomlBuf.toString("utf-8"));
829 if (!cfg) throw new Error("Failed to parse .hearthforge-ci.toml");
830
831 const invalid = validateCiConfig(cfg);
832 if (invalid)
833 throw new Error(`Invalid .hearthforge-ci.toml: ${invalid}`);
834
835 // Load secrets for log masking
836 const secrets = await db
837 .selectFrom("ci_secrets")
838 .select(["name", "value"])
839 .where("repo_id", "=", repo.id)
840 .execute();
841
842 const variableOverrides = run.variable_overrides
843 ? (JSON.parse(run.variable_overrides) as Record<string, string>)
844 : {};
845
846 const { envArray, secretValues } = buildEnvVars(
847 runId,
848 repo.name,
849 {
850 triggerSource:
851 run.trigger_source as TriggerOpts["triggerSource"],
852 commitSha: run.commit_sha ?? "",
853 commitBranch: run.commit_branch ?? undefined,
854 commitTag: run.commit_tag ?? undefined,
855 variableOverrides,
856 },
857 cfg,
858 secrets,
859 );
860
861 // Create step rows in DB
862 for (const step of cfg.steps) {
863 await db
864 .insertInto("ci_steps")
865 .values({
866 run_id: runId,
867 name: step.name,
868 status: "pending",
869 })
870 .execute();
871 }
872
873 // Pull image
874 await pullImage(cfg.image);
875 if (signal.aborted) throw new Error("Cancelled");
876
877 // Create + start container
878 containerId = await createContainer(runId, repo.name, cfg, envArray);
879 runningTasks.get(runId)!.containerId = containerId;
880
881 await startContainer(containerId);
882 if (signal.aborted) throw new Error("Cancelled");
883
884 // Create work_dir
885 if (cfg.work_dir) {
886 await execInContainer(containerId, ["mkdir", "-p", cfg.work_dir]);
887 }
888
889 // Clone project if requested
890 if (cfg.clone_project_to && run.commit_sha) {
891 await uploadRepo(containerId, repoPath(repo.name));
892 const clone = await execInContainer(
893 containerId,
894 [
895 "sh",
896 "-c",
897 // safe.directory: the upload lands as root, but the step
898 // user is whatever the image defaults to. Without this git
899 // refuses the repo as "dubious ownership".
900 `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}`,
901 ],
902 cfg.work_dir,
903 envArray,
904 );
905 if (clone.exitCode !== 0) {
906 throw new Error(
907 `Failed to clone the repository into ${cfg.clone_project_to}: ${clone.log.trim()}`,
908 );
909 }
910 }
911
912 // Execute steps
913 const shell = cfg.shell ?? ["/bin/sh", "-c"];
914 let runFailed = false;
915
916 for (const step of cfg.steps) {
917 if (signal.aborted) {
918 runFailed = true;
919 break;
920 }
921
922 const stepRow = await db
923 .selectFrom("ci_steps")
924 .select("id")
925 .where("run_id", "=", runId)
926 .where("name", "=", step.name)
927 .executeTakeFirst();
928 if (!stepRow) continue;
929 const stepId = stepRow.id;
930
931 // Check run_if condition
932 if (step.run_if) {
933 const { exitCode } = await execInContainer(
934 containerId,
935 [...shell, step.run_if],
936 cfg.work_dir,
937 envArray,
938 );
939 if (exitCode !== 0) {
940 await db
941 .updateTable("ci_steps")
942 .set({
943 status: "skipped",
944 started_at: now(),
945 finished_at: now(),
946 log: "Skipped: condition not met",
947 })
948 .where("id", "=", stepId)
949 .execute();
950 continue;
951 }
952 }
953
954 // Handle clear option
955 if (step.clear && cfg.clone_project_to && run.commit_sha) {
956 const cleared = await execInContainer(
957 containerId,
958 [
959 "sh",
960 "-c",
961 `git -C ${cfg.clone_project_to} reset --hard ${run.commit_sha} && git -C ${cfg.clone_project_to} clean -fdx`,
962 ],
963 cfg.work_dir,
964 envArray,
965 );
966 if (cleared.exitCode !== 0) {
967 throw new Error(
968 `Failed to reset ${cfg.clone_project_to} before step "${step.name}": ${cleared.log.trim()}`,
969 );
970 }
971 }
972
973 await db
974 .updateTable("ci_steps")
975 .set({ status: "running", started_at: now() })
976 .where("id", "=", stepId)
977 .execute();
978
979 let stepLog = "";
980 let stepStatus: "success" | "failure" = "success";
981
982 if (step.run_sh) {
983 const command = cfg.shell_setup
984 ? `${cfg.shell_setup}\n${step.run_sh}`
985 : step.run_sh;
986
987 const stepTimeout =
988 step.timeout ?? cfg.timeout ?? config.CI_DEFAULT_TIMEOUT;
989 const timeoutSignal = AbortSignal.timeout(stepTimeout * 1000);
990 // Abort the exec stream on either a run cancellation or the
991 // per-step timeout. Docker has no per-exec kill, so on abort we
992 // force-remove the container (below), which kills the command
993 // still running inside it.
994 const stepSignal = AbortSignal.any([signal, timeoutSignal]);
995
996 try {
997 const { log, exitCode } = await execInContainer(
998 containerId,
999 [...shell, command],
1000 cfg.work_dir,
1001 envArray,
1002 stepSignal,
1003 async (partial) => {
1004 await db
1005 .updateTable("ci_steps")
1006 .set({
1007 log: maskSecrets(partial, secretValues),
1008 })
1009 .where("id", "=", stepId)
1010 .execute();
1011 },
1012 );
1013 stepLog = maskSecrets(log, secretValues);
1014 if (exitCode !== 0) {
1015 stepStatus = "failure";
1016 runFailed = true;
1017 }
1018 } catch (err) {
1019 // A run cancellation is reported as "cancelled" by the outer
1020 // catch — don't relabel it as a step failure here.
1021 if (signal.aborted) throw err;
1022 stepStatus = "failure";
1023 runFailed = true;
1024 if (timeoutSignal.aborted) {
1025 stepLog = `Step timed out after ${stepTimeout}s\n`;
1026 // Kill the container now so the timed-out command stops
1027 // immediately rather than lingering until cleanup.
1028 await removeContainer(containerId);
1029 containerId = "";
1030 } else {
1031 stepLog = `Step failed: ${err instanceof Error ? err.message : String(err)}\n`;
1032 }
1033 }
1034 }
1035
1036 // Collect artifacts for this step
1037 if (!runFailed || stepStatus === "success") {
1038 await collectArtifacts(
1039 runId,
1040 containerId,
1041 step,
1042 shell,
1043 envArray,
1044 ).catch(() => {});
1045 }
1046
1047 await db
1048 .updateTable("ci_steps")
1049 .set({
1050 status: stepStatus,
1051 finished_at: now(),
1052 log: stepLog,
1053 })
1054 .where("id", "=", stepId)
1055 .execute();
1056
1057 if (runFailed) break;
1058 }
1059
1060 // Mark remaining steps as skipped
1061 await db
1062 .updateTable("ci_steps")
1063 .set({
1064 status: "skipped",
1065 started_at: now(),
1066 finished_at: now(),
1067 log: "Skipped: previous step failed",
1068 })
1069 .where("run_id", "=", runId)
1070 .where("status", "=", "pending")
1071 .execute();
1072
1073 const finalStatus = runFailed ? "failure" : "success";
1074 await db
1075 .updateTable("ci_runs")
1076 .set({ status: finalStatus, finished_at: now() })
1077 .where("id", "=", runId)
1078 .execute();
1079 } catch (err) {
1080 const errMsg = err instanceof Error ? err.message : String(err);
1081 const isDockerUnavailable = errMsg.includes("No Docker/Podman socket");
1082 const status = signal.aborted
1083 ? "cancelled"
1084 : isDockerUnavailable
1085 ? "skipped"
1086 : "failure";
1087 const skipLog = signal.aborted
1088 ? "Skipped: run was cancelled"
1089 : isDockerUnavailable
1090 ? "Skipped: Docker/Podman not available"
1091 : "Skipped: run failed";
1092
1093 if (!isDockerUnavailable) {
1094 // A failed step already carries its own log. Failures outside any
1095 // step (validation, upload, clone) would leave no reason at all.
1096 const failedStep = await db
1097 .selectFrom("ci_steps")
1098 .select("id")
1099 .where("run_id", "=", runId)
1100 .where("status", "=", "failure")
1101 .executeTakeFirst();
1102 if (!failedStep) {
1103 await db
1104 .insertInto("ci_steps")
1105 .values({
1106 run_id: runId,
1107 // Not "setup": a pipeline may well have its own step
1108 // by that name.
1109 name: "pipeline setup",
1110 status: "failure",
1111 started_at: new Date().toISOString(),
1112 finished_at: new Date().toISOString(),
1113 log: `Error: ${errMsg}\n`,
1114 })
1115 .execute();
1116 }
1117 }
1118 await db
1119 .updateTable("ci_runs")
1120 .set({ status, finished_at: new Date().toISOString() })
1121 .where("id", "=", runId)
1122 .execute();
1123 // Mark pending steps as skipped
1124 await db
1125 .updateTable("ci_steps")
1126 .set({
1127 status: "skipped",
1128 finished_at: new Date().toISOString(),
1129 log: skipLog,
1130 })
1131 .where("run_id", "=", runId)
1132 .where("status", "=", "pending")
1133 .execute();
1134 } finally {
1135 if (containerId) await removeContainer(containerId);
1136 runningTasks.delete(runId);
1137 // Prune old history
1138 const run = await db
1139 .selectFrom("ci_runs")
1140 .select("repo_id")
1141 .where("id", "=", runId)
1142 .executeTakeFirst();
1143 if (run) await pruneHistory(run.repo_id).catch(() => {});
1144 // A slot just freed up — start the next queued run if any.
1145 pumpQueue();
1146 }
1147}
1148
1149// --- Public API ---
1150
1151export async function triggerRun(
1152 repoName: string,
1153 opts: TriggerOpts,
1154): Promise<number> {
1155 const repo = await db
1156 .selectFrom("repositories")
1157 .select("id")
1158 .where("name", "=", repoName)
1159 .executeTakeFirst();
1160 if (!repo) throw new Error("Repository not found");
1161
1162 const runId = await db
1163 .insertInto("ci_runs")
1164 .values({
1165 repo_id: repo.id,
1166 triggered_by: opts.triggeredBy ?? null,
1167 trigger_source: opts.triggerSource,
1168 commit_sha: opts.commitSha,
1169 commit_branch: opts.commitBranch ?? null,
1170 commit_tag: opts.commitTag ?? null,
1171 status: "pending",
1172 variable_overrides: opts.variableOverrides
1173 ? JSON.stringify(opts.variableOverrides)
1174 : null,
1175 })
1176 .returning("id")
1177 .executeTakeFirstOrThrow();
1178
1179 // Allocate the human-facing run number atomically from a per-repo counter.
1180 // A single upsert-and-increment can't collide under concurrent triggers and
1181 // never reuses a number after pruneHistory shrinks the run table — both of
1182 // which a COUNT(*)-based scheme suffered from.
1183 const counter = await db
1184 .insertInto("ci_run_counters")
1185 .values({ repo_id: repo.id, last_run_id: 1 })
1186 .onConflict((oc) =>
1187 oc.column("repo_id").doUpdateSet((eb) => ({
1188 last_run_id: eb("ci_run_counters.last_run_id", "+", 1),
1189 })),
1190 )
1191 .returning("last_run_id")
1192 .executeTakeFirstOrThrow();
1193 await db
1194 .updateTable("ci_runs")
1195 .set({ repo_run_id: counter.last_run_id })
1196 .where("id", "=", runId.id)
1197 .execute();
1198
1199 // Flip to "queued" BEFORE pushing onto queuedRunIds. Otherwise a
1200 // concurrent pumpQueue (from a finishing run) could observe our entry,
1201 // promoteToPending it, and start executeRun while our UPDATE is still
1202 // in flight — the late UPDATE would then clobber a running row back to
1203 // "queued". Only pumpQueue (the sole writer of runningTasks.set) may
1204 // mutate the row's status after the push. The TOCTOU concern is moot
1205 // here because pumpQueue cannot pump a runId that hasn't been pushed.
1206 if (runningTasks.size >= config.CI_MAX_CONCURRENT) {
1207 await db
1208 .updateTable("ci_runs")
1209 .set({ status: "queued" })
1210 .where("id", "=", runId.id)
1211 .execute();
1212 }
1213 queuedRunIds.push(runId.id);
1214 pumpQueue();
1215
1216 return runId.id;
1217}
1218
1219// Wraps the fire-and-forget executeRun so unhandled exceptions (e.g. a throw
1220// before/after its own try/finally) are logged and the run is reconciled to
1221// "failure" instead of staying pending forever.
1222function spawnRun(
1223 runId: number,
1224 signal: AbortSignal,
1225 needsPromote = false,
1226): void {
1227 void (async () => {
1228 try {
1229 if (needsPromote) await promoteToPending(runId);
1230 await executeRun(runId, signal);
1231 } catch (err) {
1232 console.error(`[ci] executeRun threw for run ${runId}:`, err);
1233 runningTasks.delete(runId);
1234 try {
1235 await db
1236 .updateTable("ci_runs")
1237 .set({
1238 status: "failure",
1239 finished_at: new Date().toISOString(),
1240 })
1241 .where("id", "=", runId)
1242 // Include "queued" so a row that was promoted-but-not-yet-
1243 // observed (or never promoted because promoteToPending threw)
1244 // still gets marked failed instead of stuck.
1245 .where("status", "in", ["pending", "running", "queued"])
1246 .execute();
1247 } catch (dbErr) {
1248 console.error(
1249 `[ci] failed to mark run ${runId} failed:`,
1250 dbErr,
1251 );
1252 }
1253 pumpQueue();
1254 }
1255 })();
1256}
1257
1258export async function retryRun(
1259 runId: number,
1260 retriedBy: number,
1261): Promise<void> {
1262 const run = await db
1263 .selectFrom("ci_runs")
1264 .select("repo_id")
1265 .where("id", "=", runId)
1266 .executeTakeFirst();
1267 if (!run) throw new Error("Run not found");
1268
1269 // Delete existing steps
1270 await db.deleteFrom("ci_steps").where("run_id", "=", runId).execute();
1271
1272 // Delete artifacts from disk and DB
1273 const artifactDir = path.join(paths.CI_ARTIFACTS_DIR, String(runId));
1274 if (existsSync(artifactDir)) {
1275 await Bun.$`rm -rf ${artifactDir}`.quiet().nothrow();
1276 }
1277 await db.deleteFrom("ci_artifacts").where("run_id", "=", runId).execute();
1278
1279 // Reset to "pending" first; if a slot isn't free, flip to "queued"
1280 // before enqueueing. Mirrors triggerRun: pumpQueue is the sole writer
1281 // of runningTasks.set, so we never reserve a slot directly here. Two
1282 // concurrent retryRun calls (or a retryRun racing triggerRun) all
1283 // funnel through pumpQueue, which serializes the size check against
1284 // slot reservation in a single synchronous turn.
1285 await db
1286 .updateTable("ci_runs")
1287 .set({
1288 status: "pending",
1289 triggered_by: retriedBy,
1290 started_at: null,
1291 finished_at: null,
1292 })
1293 .where("id", "=", runId)
1294 .execute();
1295
1296 if (runningTasks.size >= config.CI_MAX_CONCURRENT) {
1297 await db
1298 .updateTable("ci_runs")
1299 .set({ status: "queued" })
1300 .where("id", "=", runId)
1301 .execute();
1302 }
1303 queuedRunIds.push(runId);
1304 pumpQueue();
1305}
1306
1307export async function cancelRun(runId: number): Promise<void> {
1308 const task = runningTasks.get(runId);
1309 if (task) {
1310 const { containerId } = task;
1311 task.controller.abort();
1312 if (containerId) {
1313 await removeContainer(containerId).catch(() => {});
1314 }
1315 }
1316 const queueIdx = queuedRunIds.indexOf(runId);
1317 if (queueIdx >= 0) queuedRunIds.splice(queueIdx, 1);
1318 await db
1319 .updateTable("ci_runs")
1320 .set({ status: "cancelled", finished_at: new Date().toISOString() })
1321 .where("id", "=", runId)
1322 .where("status", "in", ["pending", "running", "queued"])
1323 .execute();
1324}
1325
1326async function pruneHistory(repoId: number): Promise<void> {
1327 const maxHistory = config.CI_MAX_HISTORY;
1328 const allRuns = await db
1329 .selectFrom("ci_runs")
1330 .select("id")
1331 .where("repo_id", "=", repoId)
1332 .orderBy("id", "desc")
1333 .execute();
1334
1335 if (allRuns.length <= maxHistory) return;
1336
1337 const toDelete = allRuns.slice(maxHistory).map((r) => r.id);
1338 for (const runId of toDelete) {
1339 // Remove artifacts from disk
1340 const artifactDir = path.join(paths.CI_ARTIFACTS_DIR, String(runId));
1341 if (existsSync(artifactDir)) {
1342 await Bun.$`rm -rf ${artifactDir}`.quiet().nothrow();
1343 }
1344 }
1345 await db.deleteFrom("ci_runs").where("id", "in", toDelete).execute();
1346}
1347
1348export async function purgeRepoCaches(repoName: string): Promise<void> {
1349 try {
1350 const filters = encodeURIComponent(
1351 JSON.stringify({ label: [`com.hearthforge.repo=${repoName}`] }),
1352 );
1353 const resp = await dockerFetch(`/volumes?filters=${filters}`);
1354 if (!resp.ok) {
1355 await resp.body?.cancel();
1356 return;
1357 }
1358 const data = (await resp.json()) as {
1359 Volumes?: Array<{ Name: string }>;
1360 };
1361 for (const vol of data.Volumes ?? []) {
1362 const delResp = await dockerFetch(`/volumes/${vol.Name}`, {
1363 method: "DELETE",
1364 });
1365 await delResp.body?.cancel();
1366 }
1367 } catch {
1368 // Best-effort
1369 }
1370}
1371
1372export async function cancelStaleRuns(): Promise<void> {
1373 const now = new Date().toISOString();
1374 const stale = await db
1375 .selectFrom("ci_runs")
1376 .select("id")
1377 .where("status", "in", ["pending", "running", "queued"])
1378 .execute();
1379
1380 await Promise.allSettled(
1381 stale.map((r) =>
1382 dockerFetch(`/containers/hearthforge-ci-${r.id}?force=true`, {
1383 method: "DELETE",
1384 }).then((res) => res.body?.cancel()),
1385 ),
1386 );
1387
1388 await db
1389 .updateTable("ci_runs")
1390 .set({ status: "cancelled", finished_at: now })
1391 .where("status", "in", ["pending", "running", "queued"])
1392 .execute();
1393 await db
1394 .updateTable("ci_steps")
1395 .set({
1396 status: "cancelled",
1397 finished_at: now,
1398 log: "Skipped: run was cancelled",
1399 })
1400 .where("status", "in", ["pending", "running"])
1401 .execute();
1402}
1403
1404/** Reset the cached socket path (used in tests to switch mock sockets). */
1405export function resetDockerSocket(): void {
1406 resolvedSocket = null;
1407}
1408