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