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