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, host 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 set("CI_REGISTRY", host+"/"+repoName)
428
429 for _, name := range cfg.VariableOrder {
430 if def, ok := cfg.Variables[name]; ok && def.Default != "" {
431 set(name, def.Default)
432 }
433 }
434 names := make([]string, 0, len(overrides))
435 for name := range overrides {
436 names = append(names, name)
437 }
438 sort.Strings(names)
439 for _, name := range names {
440 set(name, overrides[name])
441 }
442 for _, s := range secrets {
443 set(s.name, s.value)
444 secretValues = append(secretValues, s.value)
445 }
446
447 for _, k := range order {
448 envArray = append(envArray, k+"="+vars[k])
449 }
450 return envArray, secretValues
451}
452
453// maskSecrets replaces secret values in a log. It is a plain text match: a
454// secret printed in encoded or split form is not caught.
455func maskSecrets(text string, secrets []string) string {
456 for _, s := range secrets {
457 if s != "" {
458 text = strings.ReplaceAll(text, s, "[MASKED]")
459 }
460 }
461 return text
462}
463
464type secret struct{ name, value string }
465
466type runRow struct {
467 RepoID int64
468 TriggerSource string
469 CommitSha string
470 CommitBranch string
471 CommitTag string
472 Overrides string
473}
474
475func nullString(s string) any {
476 if s == "" {
477 return nil
478 }
479 return s
480}
481
482func (r *Runner) loadRun(ctx context.Context, runID int64) (runRow, error) {
483 var row runRow
484 var sha, branch, tag, overrides sql.NullString
485 err := r.db.QueryRowContext(ctx,
486 `SELECT repo_id, trigger_source, commit_sha, commit_branch, commit_tag, variable_overrides
487 FROM ci_runs WHERE id = ?`, runID).
488 Scan(&row.RepoID, &row.TriggerSource, &sha, &branch, &tag, &overrides)
489 if err != nil {
490 return row, errors.New("run not found")
491 }
492 row.CommitSha, row.CommitBranch = sha.String, branch.String
493 row.CommitTag, row.Overrides = tag.String, overrides.String
494 return row, nil
495}
496
497// --- Checkout upload ---
498
499// uploadCheckout puts the commit's tree into the container at destPath.
500//
501// The checkout deliberately does not happen inside the container: that would
502// force every CI image to carry git. It is not a bind mount either, because
503// the daemon resolves bind sources in the host filesystem.
504//
505// `git archive` writes uid 0 and mode 0644/0755 into the tar headers, and no
506// entry for the archive root, so the destination keeps the ownership the
507// container gave it. No .git reaches the container.
508func (r *Runner) uploadCheckout(ctx context.Context, containerID, repoName, commitSha, destPath string) error {
509 // The archive carries destPath as its prefix and is extracted at /, so
510 // the engine creates the directories and the image needs no mkdir.
511 prefix := strings.TrimPrefix(path.Clean(destPath), "/") + "/"
512 // --end-of-options: commit_sha is unvalidated text from the push.
513 cmd := exec.CommandContext(ctx, "git", "-C", r.repoPath(repoName), "archive",
514 "--format=tar", "--prefix="+prefix, "--end-of-options", commitSha)
515 cmd.Env = gitcmd.Env()
516 var stderr strings.Builder
517 cmd.Stderr = &stderr
518 stdout, err := cmd.StdoutPipe()
519 if err != nil {
520 return err
521 }
522 if err := cmd.Start(); err != nil {
523 return err
524 }
525 resp, body, putErr := r.putArchive(ctx, containerID, "/", stdout)
526 // The PUT stopped reading, so git would block writing into a full pipe.
527 stdout.Close()
528 waitErr := cmd.Wait()
529
530 if putErr != nil {
531 return putErr
532 }
533 if resp.StatusCode >= 300 {
534 return fmt.Errorf("failed to upload the checkout: HTTP %d %s",
535 resp.StatusCode, strings.TrimSpace(string(body)))
536 }
537 if waitErr != nil {
538 return fmt.Errorf("failed to read %s: git archive failed %v %s",
539 commitSha, waitErr, strings.TrimSpace(stderr.String()))
540 }
541 return nil
542}
543
544// --- Artifacts ---
545
546// collectArtifacts copies this step's published files out of the container.
547// Archives are built on the host from the engine's tar stream, so the image
548// needs no archive tools.
549func (r *Runner) collectArtifacts(ctx context.Context, runID int64, containerID string, step Step) {
550 dir := r.artifactDir(runID)
551 if err := os.MkdirAll(dir, 0o755); err != nil {
552 return
553 }
554
555 // store streams the container file straight into the artifact directory.
556 store := func(filename, srcPath string) {
557 dest := filepath.Join(dir, filename)
558 // The run directory is shared by every step. A second file of the
559 // same name would overwrite the first and still add a row.
560 if _, err := os.Stat(dest); err == nil {
561 log.Printf("[ci] run %d: artifact %s already published, skipping %s", runID, filename, srcPath)
562 return
563 }
564 size, err := r.copyFileFromContainer(ctx, containerID, srcPath, dest)
565 if err != nil {
566 log.Printf("[ci] run %d: artifact %s: %v", runID, srcPath, err)
567 return
568 }
569 r.execSQL(ctx, `INSERT INTO ci_artifacts (run_id, filename, size) VALUES (?, ?, ?)`,
570 runID, filename, size)
571 }
572
573 for _, srcPath := range step.PublishFile {
574 store(path.Base(srcPath), srcPath)
575 }
576
577 for _, f := range archiveFormats {
578 for _, srcPath := range f.pick(step) {
579 filename := path.Base(srcPath) + f.ext
580 dest := filepath.Join(dir, filename)
581 if _, err := os.Stat(dest); err == nil {
582 log.Printf("[ci] run %d: artifact %s already published, skipping %s", runID, filename, srcPath)
583 continue
584 }
585 size, err := r.storeArchive(ctx, containerID, srcPath, dest, f.ext, r.cfg.CIMaxArtifactBytes)
586 if err != nil {
587 log.Printf("[ci] run %d: artifact %s: %v", runID, srcPath, err)
588 continue
589 }
590 r.execSQL(ctx, `INSERT INTO ci_artifacts (run_id, filename, size) VALUES (?, ?, ?)`,
591 runID, filename, size)
592 }
593 }
594}
595