run.go
⎇
Raw
1package ci
2
3import (
4 "context"
5 "database/sql"
6 "encoding/json"
7 "errors"
8 "fmt"
9 "log"
10 "net/http"
11 "os"
12 "os/exec"
13 "path"
14 "path/filepath"
15 "slices"
16 "sort"
17 "strconv"
18 "strings"
19 "sync"
20 "time"
21
22 "hearthforge/internal/config"
23 "hearthforge/internal/db"
24 "hearthforge/internal/gitcmd"
25)
26
27// ciMaxLogBytes caps the log stored per step.
28const ciMaxLogBytes = 2 * 1024 * 1024
29
30// TriggerOpts describes why a run starts and which commit it builds.
31type TriggerOpts struct {
32 // TriggerSource is "push", "tag" or "manual".
33 TriggerSource string
34 CommitSha string
35 CommitBranch string
36 CommitTag string
37 TriggeredBy int64
38 VariableOverrides map[string]string
39}
40
41// task is one run that holds a concurrency slot.
42type task struct {
43 cancel context.CancelFunc
44 containerID string
45}
46
47// Runner owns the CI queue and talks to the container engine.
48type Runner struct {
49 cfg *config.Config
50 db *db.DB
51
52 socketMu sync.Mutex
53 socket string
54 client *http.Client
55
56 mu sync.Mutex
57 running map[int64]*task
58 // queue holds run ids whose DB row is "queued". pumpQueue is the sole
59 // writer of `running`, so the slot check and the reservation happen
60 // together under one lock.
61 queue []int64
62}
63
64func New(cfg *config.Config, database *db.DB) *Runner {
65 return &Runner{cfg: cfg, db: database, running: map[int64]*task{}}
66}
67
68func (r *Runner) repoPath(name string) string {
69 return filepath.Join(r.cfg.ReposDir(), name+".git")
70}
71
72// ErrNoConfig means the commit carries no .hearthforge-ci.toml. Callers tell
73// it apart from a parse error to explain why CI cannot run.
74var ErrNoConfig = errors.New(".hearthforge-ci.toml not found at commit")
75
76// ConfigAt reads and parses .hearthforge-ci.toml at a commit or ref.
77func (r *Runner) ConfigAt(ctx context.Context, repoName, ref string) (*Config, error) {
78 cmd := exec.CommandContext(ctx, "git", "-C", r.repoPath(repoName), "show",
79 "--end-of-options", ref+":.hearthforge-ci.toml")
80 cmd.Env = gitcmd.Env()
81 out, err := cmd.Output()
82 if err != nil || len(out) == 0 {
83 return nil, ErrNoConfig
84 }
85 return ParseCiConfig(string(out))
86}
87
88// --- Queue ---
89
90// QueuePosition returns the 1-based place of a queued run, or 0.
91func (r *Runner) QueuePosition(runID int64) int {
92 r.mu.Lock()
93 defer r.mu.Unlock()
94 for i, id := range r.queue {
95 if id == runID {
96 return i + 1
97 }
98 }
99 return 0
100}
101
102// enqueue puts a run in the queue and starts it when a slot is free.
103// The row is flipped to "queued" under the same lock that decides there is no
104// slot, so two triggers that arrive together cannot both stay "pending".
105func (r *Runner) enqueue(ctx context.Context, runID int64) {
106 r.mu.Lock()
107 r.queue = append(r.queue, runID)
108 if len(r.running)+len(r.queue) > r.cfg.CIMaxConcurrent {
109 // pumpQueue is the only starter and it needs the same lock, so this
110 // update cannot clobber a row that just went running.
111 r.execSQL(ctx, `UPDATE ci_runs SET status = 'queued' WHERE id = ?`, runID)
112 }
113 r.mu.Unlock()
114 r.pumpQueue()
115}
116
117func (r *Runner) pumpQueue() {
118 for {
119 r.mu.Lock()
120 if len(r.queue) == 0 || len(r.running) >= r.cfg.CIMaxConcurrent {
121 r.mu.Unlock()
122 return
123 }
124 runID := r.queue[0]
125 r.queue = r.queue[1:]
126 // The run outlives the request that created it, so it gets its own
127 // context. CancelRun cancels it.
128 ctx, cancel := context.WithCancel(context.Background())
129 t := &task{cancel: cancel}
130 r.running[runID] = t
131 r.mu.Unlock()
132
133 go r.spawnRun(ctx, runID, t)
134 }
135}
136
137// spawnRun promotes the row to pending, runs it, and reconciles the row when
138// executeRun itself fails.
139func (r *Runner) spawnRun(ctx context.Context, runID int64, t *task) {
140 defer func() {
141 if rec := recover(); rec != nil {
142 log.Printf("[ci] run %d panicked: %v", runID, rec)
143 r.finishSlot(runID, t)
144 r.execSQL(context.Background(),
145 `UPDATE ci_runs SET status = 'failure', finished_at = ?
146 WHERE id = ? AND status IN ('pending','running','queued')`, db.NowISO(), runID)
147 r.pumpQueue()
148 }
149 }()
150 // Promote queued→pending before executing, so executeRun's failure path
151 // (which only matches pending/running/queued) can mark it failed.
152 r.execSQL(ctx, `UPDATE ci_runs SET status = 'pending' WHERE id = ?`, runID)
153 r.executeRun(ctx, runID, t)
154}
155
156// finishSlot releases the run's concurrency slot. It only drops the slot this
157// execution owns: a retry that started the same run again holds a different
158// task, and deleting that one would leak its slot forever.
159func (r *Runner) finishSlot(runID int64, t *task) {
160 r.mu.Lock()
161 if r.running[runID] == t {
162 delete(r.running, runID)
163 }
164 r.mu.Unlock()
165}
166
167func (r *Runner) setContainer(t *task, containerID string) {
168 r.mu.Lock()
169 t.containerID = containerID
170 r.mu.Unlock()
171}
172
173// --- Small DB helpers ---
174
175func (r *Runner) execSQL(ctx context.Context, query string, args ...any) {
176 if _, err := r.db.ExecContext(ctx, query, args...); err != nil {
177 log.Printf("[ci] sql failed: %v", err)
178 }
179}
180
181func (r *Runner) insertStep(ctx context.Context, runID int64, name, status string, started bool) int64 {
182 var startedAt any
183 if started {
184 startedAt = db.NowISO()
185 }
186 res, err := r.db.ExecContext(ctx,
187 `INSERT INTO ci_steps (run_id, name, status, started_at) VALUES (?, ?, ?, ?)`,
188 runID, name, status, startedAt)
189 if err != nil {
190 log.Printf("[ci] insert step failed: %v", err)
191 return 0
192 }
193 id, _ := res.LastInsertId()
194 return id
195}
196
197// --- Public API ---
198
199// TriggerRun records a run and starts it, or queues it when every slot is
200// taken. It returns the run id.
201func (r *Runner) TriggerRun(ctx context.Context, repoName string, opts TriggerOpts) (int64, error) {
202 var repoID int64
203 err := r.db.QueryRowContext(ctx, `SELECT id FROM repositories WHERE name = ?`, repoName).Scan(&repoID)
204 if err != nil {
205 return 0, errors.New("repository not found")
206 }
207
208 var overrides any
209 if len(opts.VariableOverrides) > 0 {
210 buf, err := json.Marshal(opts.VariableOverrides)
211 if err != nil {
212 return 0, err
213 }
214 overrides = string(buf)
215 }
216 var triggeredBy any
217 if opts.TriggeredBy > 0 {
218 triggeredBy = opts.TriggeredBy
219 }
220 res, err := r.db.ExecContext(ctx,
221 `INSERT INTO ci_runs (repo_id, triggered_by, trigger_source, commit_sha,
222 commit_branch, commit_tag, status, variable_overrides)
223 VALUES (?, ?, ?, ?, ?, ?, 'pending', ?)`,
224 repoID, triggeredBy, opts.TriggerSource, opts.CommitSha,
225 nullString(opts.CommitBranch), nullString(opts.CommitTag), overrides)
226 if err != nil {
227 return 0, err
228 }
229 runID, err := res.LastInsertId()
230 if err != nil {
231 return 0, err
232 }
233
234 // The human-facing run number comes from a per-repo counter. One
235 // upsert-and-increment cannot collide under concurrent triggers, and it
236 // never reuses a number after pruneHistory shrinks the table.
237 var repoRunID int64
238 err = r.db.QueryRowContext(ctx,
239 `INSERT INTO ci_run_counters (repo_id, last_run_id) VALUES (?, 1)
240 ON CONFLICT(repo_id) DO UPDATE SET last_run_id = last_run_id + 1
241 RETURNING last_run_id`, repoID).Scan(&repoRunID)
242 if err != nil {
243 return 0, err
244 }
245 r.execSQL(ctx, `UPDATE ci_runs SET repo_run_id = ? WHERE id = ?`, repoRunID, runID)
246
247 r.enqueue(ctx, runID)
248 return runID, nil
249}
250
251// ErrRunNotFinished is returned when a retry targets a run that is still
252// queued or executing.
253var ErrRunNotFinished = errors.New("run is not finished")
254
255// Active reports whether a run still holds its slot. A run frees it only
256// after its cleanup, so a finished status alone does not mean idle.
257func (r *Runner) Active(runID int64) bool {
258 r.mu.Lock()
259 defer r.mu.Unlock()
260 return r.running[runID] != nil
261}
262
263// RetryRun clears a finished run's steps and artifacts and runs it again.
264// The reset only applies to a run in a terminal status, so a retry cannot
265// hijack a run that is still executing.
266func (r *Runner) RetryRun(ctx context.Context, runID, retriedBy int64) error {
267 var repoID int64
268 if err := r.db.QueryRowContext(ctx, `SELECT repo_id FROM ci_runs WHERE id = ?`, runID).Scan(&repoID); err != nil {
269 return errors.New("run not found")
270 }
271 // CancelRun marks the row before the execution has cleaned up. Its
272 // cleanup would clobber the new attempt's rows and container.
273 if r.Active(runID) {
274 return ErrRunNotFinished
275 }
276 res, err := r.db.ExecContext(ctx,
277 `UPDATE ci_runs SET status = 'pending', triggered_by = ?, started_at = NULL, finished_at = NULL
278 WHERE id = ? AND status IN ('success','failure','warning','cancelled','skipped')`,
279 retriedBy, runID)
280 if err != nil {
281 return err
282 }
283 if n, err := res.RowsAffected(); err != nil {
284 return err
285 } else if n == 0 {
286 return ErrRunNotFinished
287 }
288 r.execSQL(ctx, `DELETE FROM ci_steps WHERE run_id = ?`, runID)
289 os.RemoveAll(r.artifactDir(runID))
290 r.execSQL(ctx, `DELETE FROM ci_artifacts WHERE run_id = ?`, runID)
291 r.enqueue(ctx, runID)
292 return nil
293}
294
295// CancelRun stops a running or queued run.
296func (r *Runner) CancelRun(ctx context.Context, runID int64) error {
297 r.mu.Lock()
298 t := r.running[runID]
299 var containerID string
300 if t != nil {
301 containerID = t.containerID
302 }
303 for i, id := range r.queue {
304 if id == runID {
305 r.queue = append(r.queue[:i], r.queue[i+1:]...)
306 break
307 }
308 }
309 r.mu.Unlock()
310
311 if t != nil {
312 t.cancel()
313 if containerID != "" {
314 // Docker has no per-exec kill, so removing the container is what
315 // stops the command running inside it.
316 r.removeContainer(ctx, containerID)
317 }
318 }
319 r.execSQL(ctx,
320 `UPDATE ci_runs SET status = 'cancelled', finished_at = ?
321 WHERE id = ? AND status IN ('pending','running','queued')`, db.NowISO(), runID)
322 return nil
323}
324
325// PurgeRepoCaches deletes this repo's cache volumes and returns how many
326// went. It fails when the engine is unreachable.
327func (r *Runner) PurgeRepoCaches(ctx context.Context, repoName string) (int, error) {
328 names, err := r.listRepoVolumes(ctx, repoName)
329 if err != nil {
330 return 0, err
331 }
332 removed := 0
333 for _, name := range names {
334 if r.removeVolume(ctx, name) {
335 removed++
336 }
337 }
338 return removed, nil
339}
340
341// StopRepo cancels the repo's active runs, waits for them to clean up, and
342// deletes their artifacts and cache volumes. Call it before the repo row goes:
343// the run ids come from the database.
344func (r *Runner) StopRepo(ctx context.Context, repoID int64, repoName string) {
345 ids, err := r.queryIDs(ctx, `SELECT id FROM ci_runs WHERE repo_id = ?`, repoID)
346 if err != nil {
347 log.Printf("[ci] stop repo %s: %v", repoName, err)
348 return
349 }
350 for _, id := range ids {
351 r.CancelRun(ctx, id)
352 }
353 // A cancelled run removes its container before it frees its slot, and a
354 // volume cannot be deleted while a container holds it.
355 for deadline := time.Now().Add(30 * time.Second); time.Now().Before(deadline); {
356 if !slices.ContainsFunc(ids, r.Active) {
357 break
358 }
359 time.Sleep(100 * time.Millisecond)
360 }
361 for _, id := range ids {
362 os.RemoveAll(r.artifactDir(id))
363 }
364 if _, err := r.PurgeRepoCaches(ctx, repoName); err != nil && !errors.Is(err, errNoSocket) {
365 log.Printf("[ci] purge caches of %s: %v", repoName, err)
366 }
367}
368
369func (r *Runner) queryIDs(ctx context.Context, q string, args ...any) ([]int64, error) {
370 rows, err := r.db.QueryContext(ctx, q, args...)
371 if err != nil {
372 return nil, err
373 }
374 defer rows.Close()
375 var ids []int64
376 for rows.Next() {
377 var id int64
378 if err := rows.Scan(&id); err != nil {
379 return nil, err
380 }
381 ids = append(ids, id)
382 }
383 return ids, rows.Err()
384}
385
386// CancelStaleRuns cleans up runs a restart interrupted. Their containers are
387// force-removed and the rows read "cancelled", so they can be retried.
388func (r *Runner) CancelStaleRuns(ctx context.Context) error {
389 ids, err := r.queryIDs(ctx, `SELECT id FROM ci_runs WHERE status IN ('pending','running','queued')`)
390 if err != nil {
391 return err
392 }
393 // By name too: containers from before the label existed carry none.
394 for _, id := range ids {
395 r.removeContainer(ctx, containerName(id))
396 }
397 if err := r.removeLabelledContainers(ctx); err != nil && !errors.Is(err, errNoSocket) {
398 log.Printf("[ci] sweep containers: %v", err)
399 }
400 now := db.NowISO()
401 r.execSQL(ctx, `UPDATE ci_runs SET status = 'cancelled', finished_at = ?
402 WHERE status IN ('pending','running','queued')`, now)
403 r.execSQL(ctx, `UPDATE ci_steps SET status = 'cancelled', finished_at = ?, log = 'Skipped: run was cancelled'
404 WHERE status IN ('pending','running')`, now)
405 return nil
406}
407
408func (r *Runner) artifactDir(runID int64) string {
409 return filepath.Join(r.cfg.CIArtifactsDir(), strconv.FormatInt(runID, 10))
410}
411
412// pruneHistory keeps only the CI_MAX_HISTORY most recently finished runs of
413// a repository. By finish time, not id: a retried old run counts as recent.
414func (r *Runner) pruneHistory(ctx context.Context, repoID int64) {
415 // Active runs never get pruned: their row is what the queue and the
416 // finish path update.
417 ids, err := r.queryIDs(ctx,
418 `SELECT id FROM ci_runs WHERE repo_id = ?
419 AND status NOT IN ('pending','running','queued') ORDER BY finished_at DESC, id DESC`, repoID)
420 if err != nil || len(ids) <= r.cfg.CIMaxHistory {
421 return
422 }
423 toDelete := ids[r.cfg.CIMaxHistory:]
424 args := make([]any, len(toDelete))
425 for i, id := range toDelete {
426 os.RemoveAll(r.artifactDir(id))
427 args[i] = id
428 }
429 placeholders := strings.TrimSuffix(strings.Repeat("?,", len(args)), ",")
430 r.execSQL(ctx, `DELETE FROM ci_runs WHERE id IN (`+placeholders+`)`, args...)
431}
432
433// --- Environment ---
434
435// buildEnvVars assembles the container environment. Later sources win on a
436// name clash: config defaults, then manual overrides, then secrets.
437func buildEnvVars(runID int64, repoName, baseURL, host string, run runRow, cfg *Config,
438 overrides map[string]string, secrets []secret,
439) (envArray, secretValues []string) {
440 order := []string{}
441 vars := map[string]string{}
442 set := func(k, v string) {
443 if _, seen := vars[k]; !seen {
444 order = append(order, k)
445 }
446 vars[k] = v
447 }
448
449 refName := run.CommitTag
450 if refName == "" {
451 refName = run.CommitBranch
452 }
453 shortSha := run.CommitSha
454 if len(shortSha) > 8 {
455 shortSha = shortSha[:8]
456 }
457 set("CI", "true")
458 set("CI_PIPELINE_ID", strconv.FormatInt(runID, 10))
459 set("CI_REPO_NAME", repoName)
460 set("CI_SERVER_URL", baseURL)
461 set("CI_TRIGGER_SOURCE", run.TriggerSource)
462 set("CI_COMMIT_SHA", run.CommitSha)
463 set("CI_COMMIT_SHORT_SHA", shortSha)
464 set("CI_COMMIT_BRANCH", run.CommitBranch)
465 set("CI_COMMIT_TAG", run.CommitTag)
466 set("CI_COMMIT_REF_NAME", refName)
467 // Image names must be lowercase. The registry lowers the repo name too.
468 set("CI_REGISTRY", host+"/"+strings.ToLower(repoName))
469
470 for _, name := range cfg.VariableOrder {
471 if def, ok := cfg.Variables[name]; ok && def.Default != "" {
472 set(name, def.Default)
473 }
474 }
475 names := make([]string, 0, len(overrides))
476 for name := range overrides {
477 names = append(names, name)
478 }
479 sort.Strings(names)
480 for _, name := range names {
481 set(name, overrides[name])
482 }
483 for _, s := range secrets {
484 set(s.name, s.value)
485 secretValues = append(secretValues, s.value)
486 }
487
488 for _, k := range order {
489 envArray = append(envArray, k+"="+vars[k])
490 }
491 return envArray, secretValues
492}
493
494// maskSecrets replaces secret values in a log. It is a plain text match: a
495// secret printed in encoded or split form is not caught.
496func maskSecrets(text string, secrets []string) string {
497 // Longest first: a secret that contains a shorter one would otherwise
498 // be left partly visible.
499 sorted := slices.SortedFunc(slices.Values(secrets), func(a, b string) int { return len(b) - len(a) })
500 for _, s := range sorted {
501 if s != "" {
502 text = strings.ReplaceAll(text, s, "[MASKED]")
503 }
504 }
505 return text
506}
507
508// maskLog masks a step log. A partial log, or one cut at the size limit, may
509// end in the first half of a secret, so a tail that starts one is dropped.
510func maskLog(text string, secrets []string, partial bool) string {
511 body, truncated := strings.CutSuffix(text, logTruncatedNotice)
512 body = maskSecrets(body, secrets)
513 if !partial && !truncated {
514 return body
515 }
516 cut := 0
517 for _, s := range secrets {
518 for k := min(len(s)-1, len(body)); k > cut; k-- {
519 if strings.HasSuffix(body, s[:k]) {
520 cut = k
521 break
522 }
523 }
524 }
525 body = body[:len(body)-cut]
526 if truncated {
527 body += logTruncatedNotice
528 }
529 return body
530}
531
532type secret struct{ name, value string }
533
534type runRow struct {
535 RepoID int64
536 TriggerSource string
537 CommitSha string
538 CommitBranch string
539 CommitTag string
540 Overrides string
541}
542
543func nullString(s string) any {
544 if s == "" {
545 return nil
546 }
547 return s
548}
549
550func (r *Runner) loadRun(ctx context.Context, runID int64) (runRow, error) {
551 var row runRow
552 var sha, branch, tag, overrides sql.NullString
553 err := r.db.QueryRowContext(ctx,
554 `SELECT repo_id, trigger_source, commit_sha, commit_branch, commit_tag, variable_overrides
555 FROM ci_runs WHERE id = ?`, runID).
556 Scan(&row.RepoID, &row.TriggerSource, &sha, &branch, &tag, &overrides)
557 if err != nil {
558 return row, errors.New("run not found")
559 }
560 row.CommitSha, row.CommitBranch = sha.String, branch.String
561 row.CommitTag, row.Overrides = tag.String, overrides.String
562 return row, nil
563}
564
565// --- Checkout upload ---
566
567// uploadCheckout puts the commit's tree into the container at destPath.
568//
569// The checkout deliberately does not happen inside the container: that would
570// force every CI image to carry git. It is not a bind mount either, because
571// the daemon resolves bind sources in the host filesystem.
572//
573// `git archive` writes uid 0 and mode 0644/0755 into the tar headers, and no
574// entry for the archive root, so the destination keeps the ownership the
575// container gave it. No .git reaches the container.
576func (r *Runner) uploadCheckout(ctx context.Context, containerID, repoName, commitSha, destPath string) error {
577 // The archive carries destPath as its prefix and is extracted at /, so
578 // the engine creates the directories and the image needs no mkdir.
579 prefix := strings.TrimPrefix(path.Clean(destPath), "/") + "/"
580 // --end-of-options: commit_sha is unvalidated text from the push.
581 cmd := exec.CommandContext(ctx, "git", "-C", r.repoPath(repoName), "archive",
582 "--format=tar", "--prefix="+prefix, "--end-of-options", commitSha)
583 cmd.Env = gitcmd.Env()
584 var stderr strings.Builder
585 cmd.Stderr = &stderr
586 stdout, err := cmd.StdoutPipe()
587 if err != nil {
588 return err
589 }
590 if err := cmd.Start(); err != nil {
591 return err
592 }
593 resp, body, putErr := r.putArchive(ctx, containerID, "/", stdout)
594 // The PUT stopped reading, so git would block writing into a full pipe.
595 stdout.Close()
596 waitErr := cmd.Wait()
597
598 if putErr != nil {
599 return putErr
600 }
601 if resp.StatusCode >= 300 {
602 return fmt.Errorf("failed to upload the checkout: HTTP %d %s",
603 resp.StatusCode, strings.TrimSpace(string(body)))
604 }
605 if waitErr != nil {
606 return fmt.Errorf("failed to read %s: git archive failed %v %s",
607 commitSha, waitErr, strings.TrimSpace(stderr.String()))
608 }
609 return nil
610}
611
612// --- Artifacts ---
613
614// collectArtifacts copies this step's published files out of the container.
615// Archives are built on the host from the engine's tar stream, so the image
616// needs no archive tools.
617func (r *Runner) collectArtifacts(ctx context.Context, runID int64, containerID string, step Step) {
618 dir := r.artifactDir(runID)
619 if err := os.MkdirAll(dir, 0o755); err != nil {
620 return
621 }
622
623 // store streams the container file straight into the artifact directory.
624 store := func(filename, srcPath string) {
625 dest := filepath.Join(dir, filename)
626 // The run directory is shared by every step. A second file of the
627 // same name would overwrite the first and still add a row.
628 if _, err := os.Stat(dest); err == nil {
629 log.Printf("[ci] run %d: artifact %s already published, skipping %s", runID, filename, srcPath)
630 return
631 }
632 size, err := r.copyFileFromContainer(ctx, containerID, srcPath, dest)
633 if err != nil {
634 log.Printf("[ci] run %d: artifact %s: %v", runID, srcPath, err)
635 return
636 }
637 r.execSQL(ctx, `INSERT INTO ci_artifacts (run_id, filename, size) VALUES (?, ?, ?)`,
638 runID, filename, size)
639 }
640
641 for _, srcPath := range step.PublishFile {
642 store(path.Base(srcPath), srcPath)
643 }
644
645 for _, f := range archiveFormats {
646 for _, srcPath := range f.pick(step) {
647 filename := path.Base(srcPath) + f.ext
648 dest := filepath.Join(dir, filename)
649 if _, err := os.Stat(dest); err == nil {
650 log.Printf("[ci] run %d: artifact %s already published, skipping %s", runID, filename, srcPath)
651 continue
652 }
653 size, err := r.storeArchive(ctx, containerID, srcPath, dest, f.ext, r.cfg.CIMaxArtifactBytes)
654 if err != nil {
655 log.Printf("[ci] run %d: artifact %s: %v", runID, srcPath, err)
656 continue
657 }
658 r.execSQL(ctx, `INSERT INTO ci_artifacts (run_id, filename, size) VALUES (?, ?, ?)`,
659 runID, filename, size)
660 }
661 }
662}
663