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 if err := r.CancelRun(ctx, id); err != nil {
352 log.Printf("[ci] stop repo %s: cancel run %d: %v", repoName, id, err)
353 }
354 }
355 // A cancelled run removes its container before it frees its slot, and a
356 // volume cannot be deleted while a container holds it.
357 for deadline := time.Now().Add(30 * time.Second); time.Now().Before(deadline); {
358 if !slices.ContainsFunc(ids, r.Active) {
359 break
360 }
361 time.Sleep(100 * time.Millisecond)
362 }
363 for _, id := range ids {
364 os.RemoveAll(r.artifactDir(id))
365 }
366 if _, err := r.PurgeRepoCaches(ctx, repoName); err != nil && !errors.Is(err, errNoSocket) {
367 log.Printf("[ci] purge caches of %s: %v", repoName, err)
368 }
369}
370
371func (r *Runner) queryIDs(ctx context.Context, q string, args ...any) ([]int64, error) {
372 rows, err := r.db.QueryContext(ctx, q, args...)
373 if err != nil {
374 return nil, err
375 }
376 defer rows.Close()
377 var ids []int64
378 for rows.Next() {
379 var id int64
380 if err := rows.Scan(&id); err != nil {
381 return nil, err
382 }
383 ids = append(ids, id)
384 }
385 return ids, rows.Err()
386}
387
388// CancelStaleRuns cleans up runs a restart interrupted. Their containers are
389// force-removed and the rows read "cancelled", so they can be retried.
390func (r *Runner) CancelStaleRuns(ctx context.Context) error {
391 ids, err := r.queryIDs(ctx, `SELECT id FROM ci_runs WHERE status IN ('pending','running','queued')`)
392 if err != nil {
393 return err
394 }
395 // By name too: containers from before the label existed carry none.
396 for _, id := range ids {
397 r.removeContainer(ctx, containerName(id))
398 }
399 if err := r.removeLabelledContainers(ctx); err != nil && !errors.Is(err, errNoSocket) {
400 log.Printf("[ci] sweep containers: %v", err)
401 }
402 now := db.NowISO()
403 r.execSQL(ctx, `UPDATE ci_runs SET status = 'cancelled', finished_at = ?
404 WHERE status IN ('pending','running','queued')`, now)
405 r.execSQL(ctx, `UPDATE ci_steps SET status = 'cancelled', finished_at = ?, log = 'Skipped: run was cancelled'
406 WHERE status IN ('pending','running')`, now)
407 return nil
408}
409
410func (r *Runner) artifactDir(runID int64) string {
411 return filepath.Join(r.cfg.CIArtifactsDir(), strconv.FormatInt(runID, 10))
412}
413
414// pruneHistory keeps only the CI_MAX_HISTORY most recently finished runs of
415// a repository. By finish time, not id: a retried old run counts as recent.
416func (r *Runner) pruneHistory(ctx context.Context, repoID int64) {
417 // Active runs never get pruned: their row is what the queue and the
418 // finish path update.
419 ids, err := r.queryIDs(ctx,
420 `SELECT id FROM ci_runs WHERE repo_id = ?
421 AND status NOT IN ('pending','running','queued') ORDER BY finished_at DESC, id DESC`, repoID)
422 if err != nil || len(ids) <= r.cfg.CIMaxHistory {
423 return
424 }
425 toDelete := ids[r.cfg.CIMaxHistory:]
426 args := make([]any, len(toDelete))
427 for i, id := range toDelete {
428 os.RemoveAll(r.artifactDir(id))
429 args[i] = id
430 }
431 placeholders := strings.TrimSuffix(strings.Repeat("?,", len(args)), ",")
432 r.execSQL(ctx, `DELETE FROM ci_runs WHERE id IN (`+placeholders+`)`, args...)
433}
434
435// --- Environment ---
436
437// buildEnvVars assembles the container environment. Later sources win on a
438// name clash: config defaults, then manual overrides, then secrets.
439func buildEnvVars(runID int64, repoName, baseURL, host string, run runRow, cfg *Config,
440 overrides map[string]string, secrets []secret,
441) (envArray, secretValues []string) {
442 order := []string{}
443 vars := map[string]string{}
444 set := func(k, v string) {
445 if _, seen := vars[k]; !seen {
446 order = append(order, k)
447 }
448 vars[k] = v
449 }
450
451 refName := run.CommitTag
452 if refName == "" {
453 refName = run.CommitBranch
454 }
455 shortSha := run.CommitSha
456 if len(shortSha) > 8 {
457 shortSha = shortSha[:8]
458 }
459 set("CI", "true")
460 set("CI_PIPELINE_ID", strconv.FormatInt(runID, 10))
461 set("CI_REPO_NAME", repoName)
462 set("CI_SERVER_URL", baseURL)
463 set("CI_TRIGGER_SOURCE", run.TriggerSource)
464 set("CI_COMMIT_SHA", run.CommitSha)
465 set("CI_COMMIT_SHORT_SHA", shortSha)
466 set("CI_COMMIT_BRANCH", run.CommitBranch)
467 set("CI_COMMIT_TAG", run.CommitTag)
468 set("CI_COMMIT_REF_NAME", refName)
469 // Image names must be lowercase. The registry lowers the repo name too.
470 set("CI_REGISTRY", host+"/"+strings.ToLower(repoName))
471
472 for _, name := range cfg.VariableOrder {
473 if def, ok := cfg.Variables[name]; ok && def.Default != "" {
474 set(name, def.Default)
475 }
476 }
477 names := make([]string, 0, len(overrides))
478 for name := range overrides {
479 names = append(names, name)
480 }
481 sort.Strings(names)
482 for _, name := range names {
483 set(name, overrides[name])
484 }
485 for _, s := range secrets {
486 set(s.name, s.value)
487 secretValues = append(secretValues, s.value)
488 }
489
490 for _, k := range order {
491 envArray = append(envArray, k+"="+vars[k])
492 }
493 return envArray, secretValues
494}
495
496// maskSecrets replaces secret values in a log. It is a plain text match: a
497// secret printed in encoded or split form is not caught.
498func maskSecrets(text string, secrets []string) string {
499 // Longest first: a secret that contains a shorter one would otherwise
500 // be left partly visible.
501 sorted := slices.SortedFunc(slices.Values(secrets), func(a, b string) int { return len(b) - len(a) })
502 for _, s := range sorted {
503 if s != "" {
504 text = strings.ReplaceAll(text, s, "[MASKED]")
505 }
506 }
507 return text
508}
509
510// maskLog masks a step log. A partial log, or one cut at the size limit, may
511// end in the first half of a secret, so a tail that starts one is dropped.
512func maskLog(text string, secrets []string, partial bool) string {
513 body, truncated := strings.CutSuffix(text, logTruncatedNotice)
514 body = maskSecrets(body, secrets)
515 if !partial && !truncated {
516 return body
517 }
518 cut := 0
519 for _, s := range secrets {
520 for k := min(len(s)-1, len(body)); k > cut; k-- {
521 if strings.HasSuffix(body, s[:k]) {
522 cut = k
523 break
524 }
525 }
526 }
527 body = body[:len(body)-cut]
528 if truncated {
529 body += logTruncatedNotice
530 }
531 return body
532}
533
534type secret struct{ name, value string }
535
536type runRow struct {
537 RepoID int64
538 TriggerSource string
539 CommitSha string
540 CommitBranch string
541 CommitTag string
542 Overrides string
543}
544
545func nullString(s string) any {
546 if s == "" {
547 return nil
548 }
549 return s
550}
551
552func (r *Runner) loadRun(ctx context.Context, runID int64) (runRow, error) {
553 var row runRow
554 var sha, branch, tag, overrides sql.NullString
555 err := r.db.QueryRowContext(ctx,
556 `SELECT repo_id, trigger_source, commit_sha, commit_branch, commit_tag, variable_overrides
557 FROM ci_runs WHERE id = ?`, runID).
558 Scan(&row.RepoID, &row.TriggerSource, &sha, &branch, &tag, &overrides)
559 if err != nil {
560 return row, errors.New("run not found")
561 }
562 row.CommitSha, row.CommitBranch = sha.String, branch.String
563 row.CommitTag, row.Overrides = tag.String, overrides.String
564 return row, nil
565}
566
567// --- Checkout upload ---
568
569// uploadCheckout puts the commit's tree into the container at destPath.
570//
571// The checkout deliberately does not happen inside the container: that would
572// force every CI image to carry git. It is not a bind mount either, because
573// the daemon resolves bind sources in the host filesystem.
574//
575// `git archive` writes uid 0 and mode 0644/0755 into the tar headers, and no
576// entry for the archive root, so the destination keeps the ownership the
577// container gave it. No .git reaches the container.
578func (r *Runner) uploadCheckout(ctx context.Context, containerID, repoName, commitSha, destPath string) error {
579 // The archive carries destPath as its prefix and is extracted at /, so
580 // the engine creates the directories and the image needs no mkdir.
581 prefix := strings.TrimPrefix(path.Clean(destPath), "/") + "/"
582 // --end-of-options: commit_sha is unvalidated text from the push.
583 cmd := exec.CommandContext(ctx, "git", "-C", r.repoPath(repoName), "archive",
584 "--format=tar", "--prefix="+prefix, "--end-of-options", commitSha)
585 cmd.Env = gitcmd.Env()
586 var stderr strings.Builder
587 cmd.Stderr = &stderr
588 stdout, err := cmd.StdoutPipe()
589 if err != nil {
590 return err
591 }
592 if err := cmd.Start(); err != nil {
593 return err
594 }
595 resp, body, putErr := r.putArchive(ctx, containerID, "/", stdout)
596 // The PUT stopped reading, so git would block writing into a full pipe.
597 stdout.Close()
598 waitErr := cmd.Wait()
599
600 if putErr != nil {
601 return putErr
602 }
603 if resp.StatusCode >= 300 {
604 return fmt.Errorf("failed to upload the checkout: HTTP %d %s",
605 resp.StatusCode, strings.TrimSpace(string(body)))
606 }
607 if waitErr != nil {
608 return fmt.Errorf("failed to read %s: git archive failed %v %s",
609 commitSha, waitErr, strings.TrimSpace(stderr.String()))
610 }
611 return nil
612}
613
614// --- Artifacts ---
615
616// collectArtifacts copies this step's published files out of the container.
617// Archives are built on the host from the engine's tar stream, so the image
618// needs no archive tools.
619func (r *Runner) collectArtifacts(ctx context.Context, runID int64, containerID string, step Step) {
620 dir := r.artifactDir(runID)
621 if err := os.MkdirAll(dir, 0o755); err != nil {
622 return
623 }
624
625 // store streams the container file straight into the artifact directory.
626 store := func(filename, srcPath string) {
627 dest := filepath.Join(dir, filename)
628 // The run directory is shared by every step. A second file of the
629 // same name would overwrite the first and still add a row.
630 if _, err := os.Stat(dest); err == nil {
631 log.Printf("[ci] run %d: artifact %s already published, skipping %s", runID, filename, srcPath)
632 return
633 }
634 size, err := r.copyFileFromContainer(ctx, containerID, srcPath, dest)
635 if err != nil {
636 log.Printf("[ci] run %d: artifact %s: %v", runID, srcPath, err)
637 return
638 }
639 r.execSQL(ctx, `INSERT INTO ci_artifacts (run_id, filename, size) VALUES (?, ?, ?)`,
640 runID, filename, size)
641 }
642
643 for _, srcPath := range step.PublishFile {
644 store(path.Base(srcPath), srcPath)
645 }
646
647 for _, f := range archiveFormats {
648 for _, srcPath := range f.pick(step) {
649 filename := path.Base(srcPath) + f.ext
650 dest := filepath.Join(dir, filename)
651 if _, err := os.Stat(dest); err == nil {
652 log.Printf("[ci] run %d: artifact %s already published, skipping %s", runID, filename, srcPath)
653 continue
654 }
655 size, err := r.storeArchive(ctx, containerID, srcPath, dest, f.ext, r.cfg.CIMaxArtifactBytes)
656 if err != nil {
657 log.Printf("[ci] run %d: artifact %s: %v", runID, srcPath, err)
658 continue
659 }
660 r.execSQL(ctx, `INSERT INTO ci_artifacts (run_id, filename, size) VALUES (?, ?, ?)`,
661 runID, filename, size)
662 }
663 }
664}
665