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