package ci import ( "context" "database/sql" "encoding/json" "errors" "fmt" "log" "net/http" "os" "os/exec" "path" "path/filepath" "sort" "strconv" "strings" "sync" "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 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{}} } 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") // 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") } 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 } // 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 { rows, err := r.db.QueryContext(ctx, `SELECT id FROM ci_runs WHERE status IN ('pending','running','queued')`) if err != nil { return err } var ids []int64 for rows.Next() { var id int64 if err := rows.Scan(&id); err != nil { rows.Close() return err } ids = append(ids, id) } rows.Close() if err := rows.Err(); err != nil { return err } for _, id := range ids { r.removeContainer(ctx, containerName(id)) } 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 newest CI_MAX_HISTORY runs of a repository. 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. rows, err := r.db.QueryContext(ctx, `SELECT id FROM ci_runs WHERE repo_id = ? AND status NOT IN ('pending','running','queued') ORDER BY id DESC`, repoID) if err != nil { return } var ids []int64 for rows.Next() { var id int64 if rows.Scan(&id) == nil { ids = append(ids, id) } } rows.Close() if 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) set("CI_REGISTRY", host+"/"+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 } // 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 { for _, s := range secrets { if s != "" { text = strings.ReplaceAll(text, s, "[MASKED]") } } return text } 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) } } }