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