execute.go
⎇
Raw
1package ci
2
3import (
4 "context"
5 "encoding/json"
6 "errors"
7 "fmt"
8 "strings"
9 "time"
10
11 "hearthforge/internal/db"
12)
13
14// executeRun drives one run from "pending" to a final status. It never
15// returns an error: every outcome is written to the run and step rows.
16func (r *Runner) executeRun(ctx context.Context, runID int64, t *task) {
17 var containerID string
18 // Set once the config is known. The cleanup block cannot see cfg, and
19 // pruning off another branch is wrong.
20 var pruneRepo string
21 var limitsRepo string
22 var cacheEntries []Cache
23
24 // Inserted before any config step, so ordering by id puts it first. It
25 // covers everything up to the first step: config parsing, the image pull,
26 // the container, the copies and the checkout. Not called "setup": a
27 // pipeline may well have its own step by that name.
28 setupStepID := r.insertStep(ctx, runID, "pipeline setup", "running", true)
29
30 // Written after the container is gone, so a retry that sees a terminal
31 // status can never collide with a container name still in use.
32 finalStatus := "success"
33
34 err := func() error {
35 r.execSQL(ctx, `UPDATE ci_runs SET status = 'running', started_at = ? WHERE id = ?`, db.NowISO(), runID)
36
37 run, err := r.loadRun(ctx, runID)
38 if err != nil {
39 return err
40 }
41 var repoName, defaultBranch string
42 var repoID int64
43 err = r.db.QueryRowContext(ctx,
44 `SELECT id, name, default_branch FROM repositories WHERE id = ?`, run.RepoID).
45 Scan(&repoID, &repoName, &defaultBranch)
46 if err != nil {
47 return errors.New("repo not found")
48 }
49
50 cfg, err := r.ConfigAt(ctx, repoName, run.CommitSha)
51 if err != nil {
52 return err
53 }
54 if msg := ValidateCiConfig(cfg); msg != "" {
55 return fmt.Errorf("invalid .hearthforge-ci.toml: %s", msg)
56 }
57
58 // Size caps apply on any branch: an oversized volume is oversized
59 // whoever noticed. Dropping stale volumes is default-branch only,
60 // because the config that names them is read per commit.
61 limitsRepo = repoName
62 cacheEntries = cfg.Cache
63 if run.CommitBranch != "" && run.CommitBranch == defaultBranch {
64 pruneRepo = repoName
65 }
66
67 secrets, err := r.loadSecrets(ctx, repoID)
68 if err != nil {
69 return err
70 }
71 overrides := map[string]string{}
72 if run.Overrides != "" {
73 if err := json.Unmarshal([]byte(run.Overrides), &overrides); err != nil {
74 return fmt.Errorf("invalid overrides: %w", err)
75 }
76 }
77 envArray, secretValues := buildEnvVars(runID, repoName, r.cfg.BaseURL, run, cfg, overrides, secrets)
78
79 // Step names need not be unique, so a name cannot identify a row.
80 // Keep the ids in config order instead.
81 stepIDs := make([]int64, len(cfg.Steps))
82 for i, step := range cfg.Steps {
83 stepIDs[i] = r.insertStep(ctx, runID, step.Name, "pending", false)
84 }
85
86 if err := r.pullImage(ctx, cfg.Image); err != nil {
87 return err
88 }
89 if ctx.Err() != nil {
90 return ctx.Err()
91 }
92
93 containerID, err = r.createContainer(ctx, runID, repoName, cfg, envArray)
94 if err != nil {
95 return err
96 }
97 r.setContainer(t, containerID)
98
99 if err := r.startContainer(ctx, containerID); err != nil {
100 return err
101 }
102 if ctx.Err() != nil {
103 return ctx.Err()
104 }
105
106 if cfg.WorkDir != "" {
107 if _, err := r.exec(ctx, containerID, []string{"mkdir", "-p", cfg.WorkDir}, "", nil, nil); err != nil {
108 return fmt.Errorf("create workdir: %w", err)
109 }
110 }
111
112 // Copies run before the clone, so a step can rely on the tools they
113 // bring in, and so they can supply a shell the image lacks.
114 for _, spec := range cfg.Copy {
115 if err := r.copyFromImage(ctx, containerID, spec); err != nil {
116 return err
117 }
118 if ctx.Err() != nil {
119 return ctx.Err()
120 }
121 }
122
123 if cfg.CloneProjectTo != "" && run.CommitSha != "" {
124 if err := r.uploadCheckout(ctx, containerID, repoName, run.CommitSha, cfg.CloneProjectTo); err != nil {
125 return err
126 }
127 }
128
129 r.execSQL(ctx, `UPDATE ci_steps SET status = 'success', finished_at = ? WHERE id = ?`, db.NowISO(), setupStepID)
130
131 shell := cfg.Shell
132 if len(shell) == 0 {
133 shell = []string{"/bin/sh", "-c"}
134 }
135 runFailed := false
136 sawWarning := false
137
138 for i, step := range cfg.Steps {
139 // Returned, not broken out of: the follow-up SQL below runs on the
140 // cancelled context and would silently do nothing, leaving every
141 // remaining step pending. The caller's cleanup context marks them.
142 if ctx.Err() != nil {
143 return ctx.Err()
144 }
145 // After a failure the skipped steps stay pending and are marked
146 // skipped below.
147 if runFailed && !step.Always {
148 continue
149 }
150 // A step timeout removes the container to kill the command, so
151 // there is nothing left to run an `always` step in.
152 if containerID == "" {
153 continue
154 }
155 stepID := stepIDs[i]
156
157 if step.RunIf != "" {
158 res, err := r.exec(ctx, containerID, append(append([]string{}, shell...), step.RunIf),
159 cfg.WorkDir, envArray, nil)
160 if err != nil {
161 // The engine went away. That is a run failure, not a
162 // condition that did not hold.
163 return err
164 }
165 if res.exitCode != 0 {
166 r.execSQL(ctx,
167 `UPDATE ci_steps SET status = 'skipped', started_at = ?, finished_at = ?, log = ? WHERE id = ?`,
168 db.NowISO(), db.NowISO(), "Skipped: condition not met", stepID)
169 continue
170 }
171 }
172
173 // clear: drop the directory and re-extract. `git clean` would
174 // need git in the image, and this also removes untracked files.
175 if step.Clear && cfg.CloneProjectTo != "" && run.CommitSha != "" {
176 clearErr := func() error {
177 // One exec, not two. `rm -rf` can delete the container's
178 // WorkingDir, and every later exec then fails to chdir.
179 // This exec chdirs first.
180 cmd := append(append([]string{}, shell...),
181 `rm -rf "$1" && mkdir -p "$1"`, "sh", cfg.CloneProjectTo)
182 reset, err := r.exec(ctx, containerID, cmd, "", nil, nil)
183 if err != nil {
184 return err
185 }
186 if reset.exitCode != 0 {
187 return errors.New(strings.TrimSpace(reset.log))
188 }
189 return r.uploadCheckout(ctx, containerID, repoName, run.CommitSha, cfg.CloneProjectTo)
190 }()
191 if clearErr != nil {
192 // Not returned: the outer handler only records a message
193 // while setup still owns the run. The step row survives.
194 r.execSQL(ctx,
195 `UPDATE ci_steps SET status = 'failure', started_at = ?, finished_at = ?, log = ? WHERE id = ?`,
196 db.NowISO(), db.NowISO(),
197 fmt.Sprintf("Failed to reset %s: %s\n", cfg.CloneProjectTo, clearErr), stepID)
198 runFailed = true
199 continue
200 }
201 }
202
203 r.execSQL(ctx, `UPDATE ci_steps SET status = 'running', started_at = ? WHERE id = ?`, db.NowISO(), stepID)
204
205 stepLog := ""
206 stepStatus := "success"
207 // Applies to every way a step can fail, a timeout included.
208 onFail := "failure"
209 if step.WarnOnFail {
210 onFail = "warning"
211 }
212
213 if step.RunSh != "" {
214 command := step.RunSh
215 if cfg.ShellSetup != "" {
216 command = cfg.ShellSetup + "\n" + step.RunSh
217 }
218 timeout := step.Timeout
219 if timeout == 0 {
220 timeout = cfg.Timeout
221 }
222 if timeout == 0 {
223 timeout = r.cfg.CIDefaultTimeout
224 }
225 stepCtx, cancel := context.WithTimeout(ctx, time.Duration(timeout)*time.Second)
226 res, execErr := r.exec(stepCtx, containerID,
227 append(append([]string{}, shell...), command), cfg.WorkDir, envArray,
228 func(partial string) {
229 r.execSQL(ctx, `UPDATE ci_steps SET log = ? WHERE id = ?`,
230 maskSecrets(partial, secretValues), stepID)
231 })
232 timedOut := stepCtx.Err() == context.DeadlineExceeded
233 cancel()
234
235 switch {
236 case ctx.Err() != nil:
237 // A run cancellation is reported as "cancelled" by the
238 // outer handler. Do not relabel it as a step failure.
239 return ctx.Err()
240 case timedOut:
241 // Not subject to warn_on_fail. The container is about to
242 // be destroyed, so every later step is skipped anyway. A
243 // run that cannot continue is a failure.
244 stepStatus = "failure"
245 stepLog = fmt.Sprintf("Step timed out after %ds\n", timeout)
246 // Kill the container now so the timed-out command stops
247 // immediately rather than lingering until cleanup.
248 r.removeContainer(context.WithoutCancel(ctx), containerID)
249 containerID = ""
250 r.setContainer(t, "")
251 case execErr != nil:
252 stepStatus = onFail
253 stepLog = fmt.Sprintf("Step failed: %s\n", execErr)
254 default:
255 stepLog = maskSecrets(res.log, secretValues)
256 if res.exitCode != 0 {
257 stepStatus = onFail
258 }
259 }
260 }
261
262 if stepStatus == "failure" {
263 runFailed = true
264 }
265 if stepStatus == "warning" {
266 sawWarning = true
267 }
268 if stepStatus != "failure" {
269 r.collectArtifacts(ctx, runID, containerID, step, envArray)
270 }
271
272 r.execSQL(ctx, `UPDATE ci_steps SET status = ?, finished_at = ?, log = ? WHERE id = ?`,
273 stepStatus, db.NowISO(), stepLog, stepID)
274 }
275
276 r.execSQL(ctx,
277 `UPDATE ci_steps SET status = 'skipped', started_at = ?, finished_at = ?, log = 'Skipped: previous step failed'
278 WHERE run_id = ? AND status = 'pending'`, db.NowISO(), db.NowISO(), runID)
279
280 if runFailed {
281 finalStatus = "failure"
282 } else if sawWarning {
283 finalStatus = "warning"
284 }
285 return nil
286 }()
287
288 // The run context may be cancelled from here on, so cleanup and the
289 // failure rows use a context that outlives it.
290 cleanupCtx := context.WithoutCancel(ctx)
291
292 if err != nil {
293 cancelled := ctx.Err() != nil
294 dockerUnavailable := errors.Is(err, errNoSocket)
295 status := "failure"
296 stepStatus := "failure"
297 skipLog := "Skipped: run failed"
298 switch {
299 case cancelled:
300 status, stepStatus, skipLog = "cancelled", "cancelled", "Skipped: run was cancelled"
301 case dockerUnavailable:
302 status, stepStatus, skipLog = "skipped", "skipped", "Skipped: Docker/Podman not available"
303 }
304 finalStatus = status
305 setupLog := fmt.Sprintf("Error: %s\n", err)
306 if dockerUnavailable {
307 setupLog = skipLog
308 }
309 // Only while setup still owns the run. Once it succeeds, a failure
310 // belongs to a step, and that step carries its own log.
311 r.execSQL(cleanupCtx,
312 `UPDATE ci_steps SET status = ?, finished_at = ?, log = ? WHERE id = ? AND status = 'running'`,
313 stepStatus, db.NowISO(), setupLog, setupStepID)
314 r.execSQL(cleanupCtx,
315 `UPDATE ci_steps SET status = 'skipped', finished_at = ?, log = ? WHERE run_id = ? AND status = 'pending'`,
316 db.NowISO(), skipLog, runID)
317 }
318
319 if containerID != "" {
320 r.removeContainer(cleanupCtx, containerID)
321 }
322 r.execSQL(cleanupCtx, `UPDATE ci_runs SET status = ?, finished_at = ? WHERE id = ?`,
323 finalStatus, db.NowISO(), runID)
324 r.finishSlot(runID, t)
325
326 var repoID int64
327 if r.db.QueryRowContext(cleanupCtx, `SELECT repo_id FROM ci_runs WHERE id = ?`, runID).Scan(&repoID) == nil {
328 r.pruneHistory(cleanupCtx, repoID)
329 }
330 if pruneRepo != "" {
331 r.pruneStaleCaches(cleanupCtx, pruneRepo, cacheEntries)
332 }
333 // After removeContainer above: the daemon will not delete a volume that
334 // is still mounted.
335 if limitsRepo != "" {
336 if dropped := r.enforceCacheLimits(cleanupCtx, limitsRepo, cacheEntries); len(dropped) > 0 {
337 // Recorded on the run. The only other symptom is a slow next build
338 // with no visible cause.
339 lines := make([]string, len(dropped))
340 for i, d := range dropped {
341 lines[i] = fmt.Sprintf("Dropped cache %s: %s over the %s limit",
342 d.path, formatBytes(d.size), formatBytes(d.maxSize))
343 }
344 r.execSQL(cleanupCtx,
345 `INSERT INTO ci_steps (run_id, name, status, started_at, finished_at, log)
346 VALUES (?, 'cache', 'success', ?, ?, ?)`,
347 runID, db.NowISO(), db.NowISO(), strings.Join(lines, "\n")+"\n")
348 }
349 }
350 // A slot just freed up. Start the next queued run, if any.
351 r.pumpQueue()
352}
353
354func (r *Runner) loadSecrets(ctx context.Context, repoID int64) ([]secret, error) {
355 rows, err := r.db.QueryContext(ctx, `SELECT name, value FROM ci_secrets WHERE repo_id = ?`, repoID)
356 if err != nil {
357 return nil, err
358 }
359 defer rows.Close()
360 var out []secret
361 for rows.Next() {
362 var s secret
363 if err := rows.Scan(&s.name, &s.value); err != nil {
364 return nil, err
365 }
366 out = append(out, s)
367 }
368 return out, rows.Err()
369}
370