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