package ci import ( "context" "encoding/json" "errors" "fmt" "strings" "time" "hearthforge/internal/db" ) // setupTimeout bounds everything before the first step: pulls, copies and // the checkout upload. const setupTimeout = 30 * time.Minute // executeRun drives one run from "pending" to a final status. It never // returns an error: every outcome is written to the run and step rows. func (r *Runner) executeRun(ctx context.Context, runID int64, t *task) { var containerID string // Step results must land even when the run is cancelled mid-write. dbCtx := context.WithoutCancel(ctx) // Set once the config is known. The cleanup block cannot see cfg, and // pruning off another branch is wrong. var pruneRepo string var limitsRepo string var cacheEntries []Cache // Inserted before any config step, so ordering by id puts it first. It // covers everything up to the first step: config parsing, the image pull, // the container, the copies and the checkout. Not called "setup": a // pipeline may well have its own step by that name. setupStepID := r.insertStep(ctx, runID, "pipeline setup", "running", true) // Written after the container is gone, so a retry that sees a terminal // status can never collide with a container name still in use. finalStatus := "success" err := func() error { r.execSQL(ctx, `UPDATE ci_runs SET status = 'running', started_at = ? WHERE id = ?`, db.NowISO(), runID) setupCtx, cancelSetup := context.WithTimeout(ctx, setupTimeout) defer cancelSetup() run, err := r.loadRun(ctx, runID) if err != nil { return err } var repoName, defaultBranch string var repoID int64 err = r.db.QueryRowContext(ctx, `SELECT id, name, default_branch FROM repositories WHERE id = ?`, run.RepoID). Scan(&repoID, &repoName, &defaultBranch) if err != nil { return errors.New("repo not found") } cfg, err := r.ConfigAt(ctx, repoName, run.CommitSha) if err != nil { return err } if msg := ValidateCiConfig(cfg); msg != "" { return fmt.Errorf("invalid .hearthforge-ci.toml: %s", msg) } // Size caps apply on any branch: an oversized volume is oversized // whoever noticed. Dropping stale volumes is default-branch only, // because the config that names them is read per commit. limitsRepo = repoName cacheEntries = cfg.Cache if run.CommitBranch != "" && run.CommitBranch == defaultBranch { pruneRepo = repoName } secrets, err := r.loadSecrets(ctx, repoID) if err != nil { return err } overrides := map[string]string{} if run.Overrides != "" { if err := json.Unmarshal([]byte(run.Overrides), &overrides); err != nil { return fmt.Errorf("invalid overrides: %w", err) } } envArray, secretValues := buildEnvVars(runID, repoName, r.cfg.BaseURL, r.cfg.PublicHost, run, cfg, overrides, secrets) // Step names need not be unique, so a name cannot identify a row. // Keep the ids in config order instead. stepIDs := make([]int64, len(cfg.Steps)) for i, step := range cfg.Steps { stepIDs[i] = r.insertStep(ctx, runID, step.Name, "pending", false) } if err := r.pullImage(setupCtx, cfg.Image); err != nil { return err } if ctx.Err() != nil { return ctx.Err() } // Set before the create: a create cut short can still leave it behind. containerID = containerName(runID) id, err := r.createContainer(setupCtx, runID, repoName, cfg, envArray) if err != nil { return err } containerID = id r.setContainer(t, containerID) if err := r.startContainer(setupCtx, containerID); err != nil { return err } if ctx.Err() != nil { return ctx.Err() } if cfg.WorkDir != "" { if err := r.mkdirInContainer(setupCtx, containerID, cfg.WorkDir); err != nil { return fmt.Errorf("create workdir: %w", err) } } // Copies run before the clone, so a step can rely on the tools they // bring in, and so they can supply a shell the image lacks. for _, spec := range cfg.Copy { if err := r.copyFromImage(setupCtx, containerID, spec); err != nil { return err } if ctx.Err() != nil { return ctx.Err() } } if cfg.CloneProjectTo != "" && run.CommitSha != "" { if err := r.uploadCheckout(setupCtx, containerID, repoName, run.CommitSha, cfg.CloneProjectTo); err != nil { return err } } cancelSetup() r.execSQL(ctx, `UPDATE ci_steps SET status = 'success', finished_at = ? WHERE id = ?`, db.NowISO(), setupStepID) shell := cfg.Shell if len(shell) == 0 { shell = []string{"/bin/sh", "-c"} } runFailed := false sawWarning := false for i, step := range cfg.Steps { // Returned, not broken out of: the caller marks the remaining // steps cancelled, not skipped. if ctx.Err() != nil { return ctx.Err() } // After a failure the skipped steps stay pending and are marked // skipped below. if runFailed && !step.Always { continue } // A step timeout removes the container to kill the command, so // there is nothing left to run an `always` step in. if containerID == "" { continue } stepID := stepIDs[i] timeout := step.Timeout if timeout == 0 { timeout = cfg.Timeout } if timeout == 0 { timeout = r.cfg.CIDefaultTimeout } if step.RunIf != "" { condCtx, cancel := context.WithTimeout(ctx, time.Duration(timeout)*time.Second) res, err := r.exec(condCtx, containerID, append(append([]string{}, shell...), step.RunIf), cfg.WorkDir, envArray, nil) cancel() if ctx.Err() != nil { return ctx.Err() } if err != nil { // A hung condition or a lost engine is a failure, not a // condition that did not hold. r.execSQL(dbCtx, `UPDATE ci_steps SET status = 'failure', started_at = ?, finished_at = ?, log = ? WHERE id = ?`, db.NowISO(), db.NowISO(), fmt.Sprintf("run_if failed: %s\n", err), stepID) runFailed = true continue } if res.exitCode != 0 { r.execSQL(dbCtx, `UPDATE ci_steps SET status = 'skipped', started_at = ?, finished_at = ?, log = ? WHERE id = ?`, db.NowISO(), db.NowISO(), "Skipped: condition not met", stepID) continue } } // clear: drop the directory and re-extract. `git clean` would // need git in the image, and this also removes untracked files. if step.Clear && cfg.CloneProjectTo != "" && run.CommitSha != "" { clearCtx, cancel := context.WithTimeout(ctx, time.Duration(timeout)*time.Second) clearErr := func() error { // One exec, not two. `rm -rf` can delete the container's // WorkingDir, and every later exec then fails to chdir. // This exec chdirs first. cmd := append(append([]string{}, shell...), `rm -rf "$1" && mkdir -p "$1"`, "sh", cfg.CloneProjectTo) reset, err := r.exec(clearCtx, containerID, cmd, "", nil, nil) if err != nil { return err } if reset.exitCode != 0 { return errors.New(strings.TrimSpace(reset.log)) } return r.uploadCheckout(clearCtx, containerID, repoName, run.CommitSha, cfg.CloneProjectTo) }() cancel() if ctx.Err() != nil { return ctx.Err() } if clearErr != nil { // Not returned: the outer handler only records a message // while setup still owns the run. The step row survives. r.execSQL(dbCtx, `UPDATE ci_steps SET status = 'failure', started_at = ?, finished_at = ?, log = ? WHERE id = ?`, db.NowISO(), db.NowISO(), fmt.Sprintf("Failed to reset %s: %s\n", cfg.CloneProjectTo, clearErr), stepID) runFailed = true continue } } r.execSQL(ctx, `UPDATE ci_steps SET status = 'running', started_at = ? WHERE id = ?`, db.NowISO(), stepID) stepLog := "" stepStatus := "success" // Applies to every way a step can fail, a timeout included. onFail := "failure" if step.WarnOnFail { onFail = "warning" } if step.RunSh != "" { command := step.RunSh if cfg.ShellSetup != "" { command = cfg.ShellSetup + "\n" + step.RunSh } stepCtx, cancel := context.WithTimeout(ctx, time.Duration(timeout)*time.Second) res, execErr := r.exec(stepCtx, containerID, append(append([]string{}, shell...), command), cfg.WorkDir, envArray, func(partial string) { r.execSQL(ctx, `UPDATE ci_steps SET log = ? WHERE id = ?`, maskLog(partial, secretValues, true), stepID) }) timedOut := stepCtx.Err() == context.DeadlineExceeded cancel() switch { case ctx.Err() != nil: // A run cancellation is reported as "cancelled" by the // outer handler. Do not relabel it as a step failure. return ctx.Err() case timedOut: // Not subject to warn_on_fail. The container is about to // be destroyed, so every later step is skipped anyway. A // run that cannot continue is a failure. stepStatus = "failure" stepLog = fmt.Sprintf("Step timed out after %ds\n", timeout) // Kill the container now so the timed-out command stops // immediately rather than lingering until cleanup. r.removeContainer(ctx, containerID) containerID = "" r.setContainer(t, "") case execErr != nil: stepStatus = onFail stepLog = fmt.Sprintf("Step failed: %s\n", execErr) default: stepLog = maskLog(res.log, secretValues, false) if res.exitCode != 0 { stepStatus = onFail } } } if step.BuildImage != nil && stepStatus == "success" { prefix := stepLog buildCtx, cancel := context.WithTimeout(ctx, time.Duration(timeout)*time.Second) buildLog, buildErr := r.buildImage(buildCtx, runID, containerID, repoName, cfg, step.BuildImage, buildVars(envArray, secrets), secrets, func(partial string) { r.execSQL(ctx, `UPDATE ci_steps SET log = ? WHERE id = ?`, maskLog(prefix+partial, secretValues, true), stepID) }) timedOut := buildCtx.Err() == context.DeadlineExceeded cancel() stepLog = maskLog(prefix+buildLog, secretValues, false) switch { case ctx.Err() != nil: return ctx.Err() case timedOut: stepStatus = "failure" stepLog += fmt.Sprintf("Step timed out after %ds\n", timeout) case buildErr != nil: stepStatus = onFail stepLog += fmt.Sprintf("Image build failed: %s\n", buildErr) } } if stepStatus == "failure" { runFailed = true } if stepStatus == "warning" { sawWarning = true } if stepStatus != "failure" { r.collectArtifacts(ctx, runID, containerID, step) } r.execSQL(dbCtx, `UPDATE ci_steps SET status = ?, finished_at = ?, log = ? WHERE id = ?`, stepStatus, db.NowISO(), stepLog, stepID) } if ctx.Err() != nil { return ctx.Err() } r.execSQL(dbCtx, `UPDATE ci_steps SET status = 'skipped', started_at = ?, finished_at = ?, log = 'Skipped: previous step failed' WHERE run_id = ? AND status = 'pending'`, db.NowISO(), db.NowISO(), runID) if runFailed { finalStatus = "failure" } else if sawWarning { finalStatus = "warning" } return nil }() if err != nil { cancelled := ctx.Err() != nil dockerUnavailable := errors.Is(err, errNoSocket) status := "failure" stepStatus := "failure" skipLog := "Skipped: run failed" switch { case cancelled: status, stepStatus, skipLog = "cancelled", "cancelled", "Skipped: run was cancelled" case dockerUnavailable: status, stepStatus, skipLog = "skipped", "skipped", "Skipped: Docker/Podman not available" } finalStatus = status setupLog := fmt.Sprintf("Error: %s\n", err) switch { case dockerUnavailable: setupLog = skipLog case !cancelled && errors.Is(err, context.DeadlineExceeded): setupLog = fmt.Sprintf("Error: pipeline setup timed out after %s\n", setupTimeout) } // Only while setup still owns the run. Once it succeeds, a failure // belongs to a step, and that step carries its own log. r.execSQL(dbCtx, `UPDATE ci_steps SET status = ?, finished_at = ?, log = ? WHERE id = ? AND status = 'running'`, stepStatus, db.NowISO(), setupLog, setupStepID) // The step a cancel interrupted keeps its partial log. r.execSQL(dbCtx, `UPDATE ci_steps SET status = ?, finished_at = ? WHERE run_id = ? AND status = 'running'`, stepStatus, db.NowISO(), runID) r.execSQL(dbCtx, `UPDATE ci_steps SET status = 'skipped', finished_at = ?, log = ? WHERE run_id = ? AND status = 'pending'`, db.NowISO(), skipLog, runID) } if containerID != "" { r.removeContainer(dbCtx, containerID) } // CancelRun may have written 'cancelled' after the last ctx check. r.execSQL(dbCtx, `UPDATE ci_runs SET status = ?, finished_at = ? WHERE id = ? AND status <> 'cancelled'`, finalStatus, db.NowISO(), runID) var repoID int64 if r.db.QueryRowContext(dbCtx, `SELECT repo_id FROM ci_runs WHERE id = ?`, runID).Scan(&repoID) == nil { r.pruneHistory(dbCtx, repoID) } engineCtx, cancelEngine := context.WithTimeout(dbCtx, cleanupTimeout) defer cancelEngine() if pruneRepo != "" { r.pruneStaleCaches(engineCtx, pruneRepo, cacheEntries) } // After removeContainer above: the daemon will not delete a volume that // is still mounted. if limitsRepo != "" { if dropped := r.enforceCacheLimits(engineCtx, limitsRepo, cacheEntries); len(dropped) > 0 { // Recorded on the run. The only other symptom is a slow next build // with no visible cause. lines := make([]string, len(dropped)) for i, d := range dropped { lines[i] = fmt.Sprintf("Dropped cache %s: %s over the %s limit", d.path, formatBytes(d.size), formatBytes(d.maxSize)) } r.execSQL(dbCtx, `INSERT INTO ci_steps (run_id, name, status, started_at, finished_at, log) VALUES (?, 'cache', 'success', ?, ?, ?)`, runID, db.NowISO(), db.NowISO(), strings.Join(lines, "\n")+"\n") } } // Last, so StopRepo's wait also covers the cleanup writes above. r.finishSlot(runID, t) r.pumpQueue() } func (r *Runner) loadSecrets(ctx context.Context, repoID int64) ([]secret, error) { rows, err := r.db.QueryContext(ctx, `SELECT name, value FROM ci_secrets WHERE repo_id = ?`, repoID) if err != nil { return nil, err } defer rows.Close() var out []secret for rows.Next() { var s secret if err := rows.Scan(&s.name, &s.value); err != nil { return nil, err } out = append(out, s) } return out, rows.Err() }