package ci import ( "context" "database/sql" "encoding/json" "errors" "fmt" "io" "log" "net/http" "os" "os/exec" "path" "path/filepath" "slices" "sort" "strconv" "strings" "sync" "time" "hearthforge/internal/config" "hearthforge/internal/db" "hearthforge/internal/gitcmd" ) // ciMaxLogBytes caps the log stored per step. const ciMaxLogBytes = 2 * 1024 * 1024 // TriggerOpts describes why a run starts and which commit it builds. type TriggerOpts struct { // TriggerSource is "push", "tag" or "manual". TriggerSource string CommitSha string CommitBranch string CommitTag string TriggeredBy int64 VariableOverrides map[string]string } // task is one run that holds a concurrency slot. type task struct { cancel context.CancelFunc containerID string } // Runner owns the CI queue and talks to the container engine. type Runner struct { cfg *config.Config db *db.DB // ImportImage stores an image from a build VM in the registry and // returns its manifest digest. The web server provides it. ImportImage func(ctx context.Context, repoName, image string, tags []string, src io.Reader) (string, error) // vmImageSlot serialises builds of the default build VM image. vmImageSlot chan struct{} socketMu sync.Mutex socket string client *http.Client mu sync.Mutex running map[int64]*task // queue holds run ids whose DB row is "queued". pumpQueue is the sole // writer of `running`, so the slot check and the reservation happen // together under one lock. queue []int64 } func New(cfg *config.Config, database *db.DB) *Runner { return &Runner{cfg: cfg, db: database, running: map[int64]*task{}, vmImageSlot: make(chan struct{}, 1)} } func (r *Runner) repoPath(name string) string { return filepath.Join(r.cfg.ReposDir(), name+".git") } // ErrNoConfig means the commit carries no .hearthforge-ci.toml. Callers tell // it apart from a parse error to explain why CI cannot run. var ErrNoConfig = errors.New(".hearthforge-ci.toml not found at commit") // ConfigAt reads and parses .hearthforge-ci.toml at a commit or ref. func (r *Runner) ConfigAt(ctx context.Context, repoName, ref string) (*Config, error) { cmd := exec.CommandContext(ctx, "git", "-C", r.repoPath(repoName), "show", "--end-of-options", ref+":.hearthforge-ci.toml") cmd.Env = gitcmd.Env() out, err := cmd.Output() if err != nil || len(out) == 0 { return nil, ErrNoConfig } return ParseCiConfig(string(out)) } // --- Queue --- // QueuePosition returns the 1-based place of a queued run, or 0. func (r *Runner) QueuePosition(runID int64) int { r.mu.Lock() defer r.mu.Unlock() for i, id := range r.queue { if id == runID { return i + 1 } } return 0 } // enqueue puts a run in the queue and starts it when a slot is free. // The row is flipped to "queued" under the same lock that decides there is no // slot, so two triggers that arrive together cannot both stay "pending". func (r *Runner) enqueue(ctx context.Context, runID int64) { r.mu.Lock() r.queue = append(r.queue, runID) if len(r.running)+len(r.queue) > r.cfg.CIMaxConcurrent { // pumpQueue is the only starter and it needs the same lock, so this // update cannot clobber a row that just went running. r.execSQL(ctx, `UPDATE ci_runs SET status = 'queued' WHERE id = ?`, runID) } r.mu.Unlock() r.pumpQueue() } func (r *Runner) pumpQueue() { for { r.mu.Lock() if len(r.queue) == 0 || len(r.running) >= r.cfg.CIMaxConcurrent { r.mu.Unlock() return } runID := r.queue[0] r.queue = r.queue[1:] // The run outlives the request that created it, so it gets its own // context. CancelRun cancels it. ctx, cancel := context.WithCancel(context.Background()) t := &task{cancel: cancel} r.running[runID] = t r.mu.Unlock() go r.spawnRun(ctx, runID, t) } } // spawnRun promotes the row to pending, runs it, and reconciles the row when // executeRun itself fails. func (r *Runner) spawnRun(ctx context.Context, runID int64, t *task) { defer func() { if rec := recover(); rec != nil { log.Printf("[ci] run %d panicked: %v", runID, rec) r.finishSlot(runID, t) r.execSQL(context.Background(), `UPDATE ci_runs SET status = 'failure', finished_at = ? WHERE id = ? AND status IN ('pending','running','queued')`, db.NowISO(), runID) r.pumpQueue() } }() // Promote queued→pending before executing, so executeRun's failure path // (which only matches pending/running/queued) can mark it failed. r.execSQL(ctx, `UPDATE ci_runs SET status = 'pending' WHERE id = ?`, runID) r.executeRun(ctx, runID, t) } // finishSlot releases the run's concurrency slot. It only drops the slot this // execution owns: a retry that started the same run again holds a different // task, and deleting that one would leak its slot forever. func (r *Runner) finishSlot(runID int64, t *task) { r.mu.Lock() if r.running[runID] == t { delete(r.running, runID) } r.mu.Unlock() } func (r *Runner) setContainer(t *task, containerID string) { r.mu.Lock() t.containerID = containerID r.mu.Unlock() } // --- Small DB helpers --- func (r *Runner) execSQL(ctx context.Context, query string, args ...any) { if _, err := r.db.ExecContext(ctx, query, args...); err != nil { log.Printf("[ci] sql failed: %v", err) } } func (r *Runner) insertStep(ctx context.Context, runID int64, name, status string, started bool) int64 { var startedAt any if started { startedAt = db.NowISO() } res, err := r.db.ExecContext(ctx, `INSERT INTO ci_steps (run_id, name, status, started_at) VALUES (?, ?, ?, ?)`, runID, name, status, startedAt) if err != nil { log.Printf("[ci] insert step failed: %v", err) return 0 } id, _ := res.LastInsertId() return id } // --- Public API --- // TriggerRun records a run and starts it, or queues it when every slot is // taken. It returns the run id. func (r *Runner) TriggerRun(ctx context.Context, repoName string, opts TriggerOpts) (int64, error) { var repoID int64 err := r.db.QueryRowContext(ctx, `SELECT id FROM repositories WHERE name = ?`, repoName).Scan(&repoID) if err != nil { return 0, errors.New("repository not found") } var overrides any if len(opts.VariableOverrides) > 0 { buf, err := json.Marshal(opts.VariableOverrides) if err != nil { return 0, err } overrides = string(buf) } var triggeredBy any if opts.TriggeredBy > 0 { triggeredBy = opts.TriggeredBy } res, err := r.db.ExecContext(ctx, `INSERT INTO ci_runs (repo_id, triggered_by, trigger_source, commit_sha, commit_branch, commit_tag, status, variable_overrides) VALUES (?, ?, ?, ?, ?, ?, 'pending', ?)`, repoID, triggeredBy, opts.TriggerSource, opts.CommitSha, nullString(opts.CommitBranch), nullString(opts.CommitTag), overrides) if err != nil { return 0, err } runID, err := res.LastInsertId() if err != nil { return 0, err } // The human-facing run number comes from a per-repo counter. One // upsert-and-increment cannot collide under concurrent triggers, and it // never reuses a number after pruneHistory shrinks the table. var repoRunID int64 err = r.db.QueryRowContext(ctx, `INSERT INTO ci_run_counters (repo_id, last_run_id) VALUES (?, 1) ON CONFLICT(repo_id) DO UPDATE SET last_run_id = last_run_id + 1 RETURNING last_run_id`, repoID).Scan(&repoRunID) if err != nil { return 0, err } r.execSQL(ctx, `UPDATE ci_runs SET repo_run_id = ? WHERE id = ?`, repoRunID, runID) r.enqueue(ctx, runID) return runID, nil } // ErrRunNotFinished is returned when a retry targets a run that is still // queued or executing. var ErrRunNotFinished = errors.New("run is not finished") // Active reports whether a run still holds its slot. A run frees it only // after its cleanup, so a finished status alone does not mean idle. func (r *Runner) Active(runID int64) bool { r.mu.Lock() defer r.mu.Unlock() return r.running[runID] != nil } // RetryRun clears a finished run's steps and artifacts and runs it again. // The reset only applies to a run in a terminal status, so a retry cannot // hijack a run that is still executing. func (r *Runner) RetryRun(ctx context.Context, runID, retriedBy int64) error { var repoID int64 if err := r.db.QueryRowContext(ctx, `SELECT repo_id FROM ci_runs WHERE id = ?`, runID).Scan(&repoID); err != nil { return errors.New("run not found") } // CancelRun marks the row before the execution has cleaned up. Its // cleanup would clobber the new attempt's rows and container. if r.Active(runID) { return ErrRunNotFinished } res, err := r.db.ExecContext(ctx, `UPDATE ci_runs SET status = 'pending', triggered_by = ?, started_at = NULL, finished_at = NULL WHERE id = ? AND status IN ('success','failure','warning','cancelled','skipped')`, retriedBy, runID) if err != nil { return err } if n, err := res.RowsAffected(); err != nil { return err } else if n == 0 { return ErrRunNotFinished } r.execSQL(ctx, `DELETE FROM ci_steps WHERE run_id = ?`, runID) os.RemoveAll(r.artifactDir(runID)) r.execSQL(ctx, `DELETE FROM ci_artifacts WHERE run_id = ?`, runID) r.enqueue(ctx, runID) return nil } // CancelRun stops a running or queued run. func (r *Runner) CancelRun(ctx context.Context, runID int64) error { r.mu.Lock() t := r.running[runID] var containerID string if t != nil { containerID = t.containerID } for i, id := range r.queue { if id == runID { r.queue = append(r.queue[:i], r.queue[i+1:]...) break } } r.mu.Unlock() if t != nil { t.cancel() if containerID != "" { // Docker has no per-exec kill, so removing the container is what // stops the command running inside it. r.removeContainer(ctx, containerID) } } r.execSQL(ctx, `UPDATE ci_runs SET status = 'cancelled', finished_at = ? WHERE id = ? AND status IN ('pending','running','queued')`, db.NowISO(), runID) return nil } // PurgeRepoCaches deletes this repo's cache volumes and returns how many // went. It fails when the engine is unreachable. func (r *Runner) PurgeRepoCaches(ctx context.Context, repoName string) (int, error) { names, err := r.listRepoVolumes(ctx, repoName) if err != nil { return 0, err } removed := 0 for _, name := range names { if r.removeVolume(ctx, name) { removed++ } } return removed, nil } // StopRepo cancels the repo's active runs, waits for them to clean up, and // deletes their artifacts and cache volumes. Call it before the repo row goes: // the run ids come from the database. func (r *Runner) StopRepo(ctx context.Context, repoID int64, repoName string) { ids, err := r.queryIDs(ctx, `SELECT id FROM ci_runs WHERE repo_id = ?`, repoID) if err != nil { log.Printf("[ci] stop repo %s: %v", repoName, err) return } for _, id := range ids { if err := r.CancelRun(ctx, id); err != nil { log.Printf("[ci] stop repo %s: cancel run %d: %v", repoName, id, err) } } // A cancelled run removes its container before it frees its slot, and a // volume cannot be deleted while a container holds it. for deadline := time.Now().Add(30 * time.Second); time.Now().Before(deadline); { if !slices.ContainsFunc(ids, r.Active) { break } time.Sleep(100 * time.Millisecond) } for _, id := range ids { os.RemoveAll(r.artifactDir(id)) } if _, err := r.PurgeRepoCaches(ctx, repoName); err != nil && !errors.Is(err, errNoSocket) { log.Printf("[ci] purge caches of %s: %v", repoName, err) } } func (r *Runner) queryIDs(ctx context.Context, q string, args ...any) ([]int64, error) { rows, err := r.db.QueryContext(ctx, q, args...) if err != nil { return nil, err } defer rows.Close() var ids []int64 for rows.Next() { var id int64 if err := rows.Scan(&id); err != nil { return nil, err } ids = append(ids, id) } return ids, rows.Err() } // CancelStaleRuns cleans up runs a restart interrupted. Their containers are // force-removed and the rows read "cancelled", so they can be retried. func (r *Runner) CancelStaleRuns(ctx context.Context) error { ids, err := r.queryIDs(ctx, `SELECT id FROM ci_runs WHERE status IN ('pending','running','queued')`) if err != nil { return err } // By name too: containers from before the label existed carry none. for _, id := range ids { r.removeContainer(ctx, containerName(id)) } if err := r.removeLabelledContainers(ctx); err != nil && !errors.Is(err, errNoSocket) { log.Printf("[ci] sweep containers: %v", err) } now := db.NowISO() r.execSQL(ctx, `UPDATE ci_runs SET status = 'cancelled', finished_at = ? WHERE status IN ('pending','running','queued')`, now) r.execSQL(ctx, `UPDATE ci_steps SET status = 'cancelled', finished_at = ?, log = 'Skipped: run was cancelled' WHERE status IN ('pending','running')`, now) return nil } func (r *Runner) artifactDir(runID int64) string { return filepath.Join(r.cfg.CIArtifactsDir(), strconv.FormatInt(runID, 10)) } // pruneHistory keeps only the CI_MAX_HISTORY most recently finished runs of // a repository. By finish time, not id: a retried old run counts as recent. func (r *Runner) pruneHistory(ctx context.Context, repoID int64) { // Active runs never get pruned: their row is what the queue and the // finish path update. ids, err := r.queryIDs(ctx, `SELECT id FROM ci_runs WHERE repo_id = ? AND status NOT IN ('pending','running','queued') ORDER BY finished_at DESC, id DESC`, repoID) if err != nil || len(ids) <= r.cfg.CIMaxHistory { return } toDelete := ids[r.cfg.CIMaxHistory:] args := make([]any, len(toDelete)) for i, id := range toDelete { os.RemoveAll(r.artifactDir(id)) args[i] = id } placeholders := strings.TrimSuffix(strings.Repeat("?,", len(args)), ",") r.execSQL(ctx, `DELETE FROM ci_runs WHERE id IN (`+placeholders+`)`, args...) } // --- Environment --- // buildEnvVars assembles the container environment. Later sources win on a // name clash: config defaults, then manual overrides, then secrets. func buildEnvVars(runID int64, repoName, baseURL, host string, run runRow, cfg *Config, overrides map[string]string, secrets []secret, ) (envArray, secretValues []string) { order := []string{} vars := map[string]string{} set := func(k, v string) { if _, seen := vars[k]; !seen { order = append(order, k) } vars[k] = v } refName := run.CommitTag if refName == "" { refName = run.CommitBranch } shortSha := run.CommitSha if len(shortSha) > 8 { shortSha = shortSha[:8] } set("CI", "true") set("CI_PIPELINE_ID", strconv.FormatInt(runID, 10)) set("CI_REPO_NAME", repoName) set("CI_SERVER_URL", baseURL) set("CI_TRIGGER_SOURCE", run.TriggerSource) set("CI_COMMIT_SHA", run.CommitSha) set("CI_COMMIT_SHORT_SHA", shortSha) set("CI_COMMIT_BRANCH", run.CommitBranch) set("CI_COMMIT_TAG", run.CommitTag) set("CI_COMMIT_REF_NAME", refName) // Image names must be lowercase. The registry lowers the repo name too. set("CI_REGISTRY", host+"/"+strings.ToLower(repoName)) for _, name := range cfg.VariableOrder { if def, ok := cfg.Variables[name]; ok && def.Default != "" { set(name, def.Default) } } names := make([]string, 0, len(overrides)) for name := range overrides { names = append(names, name) } sort.Strings(names) for _, name := range names { set(name, overrides[name]) } for _, s := range secrets { set(s.name, s.value) secretValues = append(secretValues, s.value) } for _, k := range order { envArray = append(envArray, k+"="+vars[k]) } return envArray, secretValues } // buildVars maps the run environment without the secrets, for expanding // build_image tags and args. func buildVars(envArray []string, secrets []secret) map[string]string { vars := map[string]string{} for _, kv := range envArray { k, v, _ := strings.Cut(kv, "=") vars[k] = v } for _, s := range secrets { delete(vars, s.name) } return vars } // maskSecrets replaces secret values in a log. It is a plain text match: a // secret printed in encoded or split form is not caught. func maskSecrets(text string, secrets []string) string { // Longest first: a secret that contains a shorter one would otherwise // be left partly visible. sorted := slices.SortedFunc(slices.Values(secrets), func(a, b string) int { return len(b) - len(a) }) for _, s := range sorted { if s != "" { text = strings.ReplaceAll(text, s, "[MASKED]") } } return text } // maskLog masks a step log. A partial log, or one cut at the size limit, may // end in the first half of a secret, so a tail that starts one is dropped. func maskLog(text string, secrets []string, partial bool) string { body, truncated := strings.CutSuffix(text, logTruncatedNotice) body = maskSecrets(body, secrets) if !partial && !truncated { return body } cut := 0 for _, s := range secrets { for k := min(len(s)-1, len(body)); k > cut; k-- { if strings.HasSuffix(body, s[:k]) { cut = k break } } } body = body[:len(body)-cut] if truncated { body += logTruncatedNotice } return body } type secret struct{ name, value string } type runRow struct { RepoID int64 TriggerSource string CommitSha string CommitBranch string CommitTag string Overrides string } func nullString(s string) any { if s == "" { return nil } return s } func (r *Runner) loadRun(ctx context.Context, runID int64) (runRow, error) { var row runRow var sha, branch, tag, overrides sql.NullString err := r.db.QueryRowContext(ctx, `SELECT repo_id, trigger_source, commit_sha, commit_branch, commit_tag, variable_overrides FROM ci_runs WHERE id = ?`, runID). Scan(&row.RepoID, &row.TriggerSource, &sha, &branch, &tag, &overrides) if err != nil { return row, errors.New("run not found") } row.CommitSha, row.CommitBranch = sha.String, branch.String row.CommitTag, row.Overrides = tag.String, overrides.String return row, nil } // --- Checkout upload --- // uploadCheckout puts the commit's tree into the container at destPath. // // The checkout deliberately does not happen inside the container: that would // force every CI image to carry git. It is not a bind mount either, because // the daemon resolves bind sources in the host filesystem. // // `git archive` writes uid 0 and mode 0644/0755 into the tar headers, and no // entry for the archive root, so the destination keeps the ownership the // container gave it. No .git reaches the container. func (r *Runner) uploadCheckout(ctx context.Context, containerID, repoName, commitSha, destPath string) error { // The archive carries destPath as its prefix and is extracted at /, so // the engine creates the directories and the image needs no mkdir. prefix := strings.TrimPrefix(path.Clean(destPath), "/") + "/" // --end-of-options: commit_sha is unvalidated text from the push. cmd := exec.CommandContext(ctx, "git", "-C", r.repoPath(repoName), "archive", "--format=tar", "--prefix="+prefix, "--end-of-options", commitSha) cmd.Env = gitcmd.Env() var stderr strings.Builder cmd.Stderr = &stderr stdout, err := cmd.StdoutPipe() if err != nil { return err } if err := cmd.Start(); err != nil { return err } resp, body, putErr := r.putArchive(ctx, containerID, "/", stdout) // The PUT stopped reading, so git would block writing into a full pipe. stdout.Close() waitErr := cmd.Wait() if putErr != nil { return putErr } if resp.StatusCode >= 300 { return fmt.Errorf("failed to upload the checkout: HTTP %d %s", resp.StatusCode, strings.TrimSpace(string(body))) } if waitErr != nil { return fmt.Errorf("failed to read %s: git archive failed %v %s", commitSha, waitErr, strings.TrimSpace(stderr.String())) } return nil } // --- Artifacts --- // collectArtifacts copies this step's published files out of the container. // Archives are built on the host from the engine's tar stream, so the image // needs no archive tools. func (r *Runner) collectArtifacts(ctx context.Context, runID int64, containerID string, step Step) { dir := r.artifactDir(runID) if err := os.MkdirAll(dir, 0o755); err != nil { return } // store streams the container file straight into the artifact directory. store := func(filename, srcPath string) { dest := filepath.Join(dir, filename) // The run directory is shared by every step. A second file of the // same name would overwrite the first and still add a row. if _, err := os.Stat(dest); err == nil { log.Printf("[ci] run %d: artifact %s already published, skipping %s", runID, filename, srcPath) return } size, err := r.copyFileFromContainer(ctx, containerID, srcPath, dest) if err != nil { log.Printf("[ci] run %d: artifact %s: %v", runID, srcPath, err) return } r.execSQL(ctx, `INSERT INTO ci_artifacts (run_id, filename, size) VALUES (?, ?, ?)`, runID, filename, size) } for _, srcPath := range step.PublishFile { store(path.Base(srcPath), srcPath) } for _, f := range archiveFormats { for _, srcPath := range f.pick(step) { filename := path.Base(srcPath) + f.ext dest := filepath.Join(dir, filename) if _, err := os.Stat(dest); err == nil { log.Printf("[ci] run %d: artifact %s already published, skipping %s", runID, filename, srcPath) continue } size, err := r.storeArchive(ctx, containerID, srcPath, dest, f.ext, r.cfg.CIMaxArtifactBytes) if err != nil { log.Printf("[ci] run %d: artifact %s: %v", runID, srcPath, err) continue } r.execSQL(ctx, `INSERT INTO ci_artifacts (run_id, filename, size) VALUES (?, ?, ?)`, runID, filename, size) } } }