cmd/gitbay-runner/main.go
553 lines · 19385 bytes
1// gitbay-runner executes CI builds queued by a gitbay server. It polls over
2// SSH — the same authenticated channel everything else uses — claims one
3// build at a time, clones the repo, runs each step with `sh -c`, streams the
4// combined output back, and reports success or failure.
5//
6// The account behind the runner's key must be an instance admin: a runner
7// executes arbitrary repo code, so handing out jobs is the operator's call.
8// v1 runs steps directly on the host under this process's user; run it as a
9// dedicated unprivileged user.
10package main
11
12import (
13 "encoding/json"
14 "flag"
15 "fmt"
16 "io"
17 "log"
18 "os"
19 "os/exec"
20 "os/signal"
21 "path/filepath"
22 "strings"
23 "sync"
24 "syscall"
25 "time"
26
27 "gitbay.org/gitbay/internal/buildinfo"
28 "gitbay.org/gitbay/internal/toolpath"
29)
30
31type job struct {
32 ID int64 `json:"id"`
33 Repo string `json:"repo"`
34 Number int64 `json:"number"`
35 Job string `json:"job"`
36 SHA string `json:"sha"`
37 Ref string `json:"ref"`
38 Steps []string `json:"steps"`
39 Image string `json:"image"`
40 Secrets map[string]string `json:"secrets"`
41}
42
43type runner struct {
44 remote string // ssh destination, e.g. git@gitbay.org
45 sshOpts []string
46 cloneBase string // e.g. ssh://git@gitbay.org
47 workdir string
48 timeout time.Duration
49 // image is the container image for a job that names none, and
50 // isolation selects how steps run: "podman" or "none".
51 image string
52 isolation string
53 // memory and cpus cap one build's cgroup; empty means no cap.
54 memory string
55 cpus string
56 // cgroups is the runner's build cgroup subtree, nil where the unit
57 // is not delegated and builds run unconfined in the service cgroup.
58 cgroups *buildCgroups
59 // stepFn is step, replaceable by tests.
60 stepFn func() (bool, error)
61 // repos limits which repositories this runner claims builds for. Empty
62 // means any, which is what a runner on the server itself wants; a runner
63 // somewhere that should not execute every repository's steps names them.
64 repos []string
65 // untrusted also claims merge request heads from forks.
66 untrusted bool
67}
68
69func main() {
70 if len(os.Args) > 1 && os.Args[1] == "init" {
71 os.Exit(runInit(os.Args[2:]))
72 }
73 var (
74 configPath = flag.String("config", defaultConfigPath(), "config file; keys are these flag names, flags override it")
75 identity = flag.String("identity", "", "ssh private key to poll and clone with (default: the key gitbay-runner init generated, if present)")
76 untrusted = flag.Bool("untrusted", false, "also claim untrusted builds: merge request heads from forks (needs -isolation podman to be safe)")
77 remote = flag.String("remote", "git@gitbay.org", "ssh destination of the gitbay server")
78 sshOpts = flag.String("ssh-opts", "", "extra ssh options, space-separated (also used for git clone)")
79 cloneBase = flag.String("clone-base", "", "clone URL prefix (default ssh://<remote>)")
80 workdir = flag.String("workdir", defaultWorkdir(), "build workspace root")
81 poll = flag.Duration("poll", 5*time.Second, "idle poll interval")
82 timeout = flag.Duration("timeout", 30*time.Minute, "per-build time limit")
83 repos = flag.String("repos", "", "only claim builds for these repositories, comma-separated owner/name (default: any)")
84 once = flag.Bool("once", false, "process at most one build, then exit")
85 jobs = flag.Int("jobs", 1, "builds to run at once")
86 image = flag.String("image", "", "default container image for jobs that name none")
87 isolation = flag.String("isolation", "podman", "how steps run: podman, or none for no container")
88 memory = flag.String("memory", "", "memory limit per build, e.g. 4g (podman only, needs a delegated cgroup; default unlimited)")
89 cpus = flag.String("cpus", "", "CPU limit per build, e.g. 2 (podman only, needs a delegated cgroup; default unlimited)")
90 version = flag.Bool("version", false, "print the commit this binary was built from, then exit")
91 )
92 path := configPathFromArgs(os.Args[1:], *configPath)
93 if values, found, err := loadConfig(path); err != nil {
94 log.Fatal(err)
95 } else if found {
96 if err := applyConfig(flag.CommandLine, values); err != nil {
97 log.Fatal(err)
98 }
99 log.Printf("config: %s", path)
100 }
101 flag.Parse()
102 if *version {
103 fmt.Println(buildinfo.String())
104 return
105 }
106 // The runner links internal/store, so it goes stale on changes that never
107 // touch cmd/gitbay-runner. Say which commit is running.
108 log.Printf("gitbay-runner %s", buildinfo.String())
109 r := &runner{
110 remote: *remote,
111 cloneBase: *cloneBase,
112 workdir: *workdir,
113 timeout: *timeout,
114 image: *image,
115 isolation: *isolation,
116 memory: *memory,
117 cpus: *cpus,
118 }
119 if r.isolation == isolationPodman {
120 // Before podman runs anything: its pause process lands in the
121 // cgroup of the first invocation, and that must be the runner's
122 // leaf, not a build's.
123 cg, err := prepareBuildCgroups()
124 switch {
125 case err == nil:
126 r.cgroups = cg
127 case r.memory != "" || r.cpus != "":
128 // Limits that cannot be applied are refused, not dropped:
129 // a runner that accepted -memory and ran uncapped is what
130 // #188 was.
131 log.Fatalf("-memory/-cpus: build cgroups unavailable: %v", err)
132 default:
133 log.Printf("build cgroups unavailable (%v); builds run unconfined in the service cgroup", err)
134 }
135 }
136 if err := r.checkIsolation(); err != nil {
137 // Refusing to start is the point. A runner that quietly fell back
138 // to running repository code on the host would drop isolation
139 // with nothing to surface it, which is worse than a stopped
140 // runner: the operator sees a failed unit either way, but only
141 // one of them is honest about why (#144).
142 log.Fatalf("isolation: %v", err)
143 }
144 if *sshOpts != "" {
145 r.sshOpts = strings.Fields(*sshOpts)
146 }
147 if *identity == "" {
148 if p := filepath.Join(configDir(), "id_ed25519"); fileExists(p) {
149 *identity = p
150 }
151 }
152 r.sshOpts = append(identityOpts(*identity), r.sshOpts...)
153 r.untrusted = *untrusted
154 for _, name := range strings.Split(*repos, ",") {
155 if name = strings.TrimSpace(name); name != "" {
156 r.repos = append(r.repos, name)
157 }
158 }
159 if r.cloneBase == "" {
160 r.cloneBase = "ssh://" + *remote
161 }
162 // 0o700, not 0o755: a build's checkout and its secrets-bearing
163 // environment are this user's business alone, and the default sits
164 // beside other users' data on a shared host.
165 if err := os.MkdirAll(r.workdir, 0o700); err != nil {
166 log.Fatal(err)
167 }
168 if err := checkWorkdir(r.workdir); err != nil {
169 log.Fatal(err)
170 }
171 n := *jobs
172 if n < 1 {
173 log.Fatal("-jobs must be at least 1")
174 }
175 if *once {
176 // "at most one build" is one build, whatever -jobs says.
177 n = 1
178 }
179 // `runner next` claims inside one transaction, so several workers
180 // claiming at once is already safe; the runner just never used that.
181 // Each build works in its own build-<id> directory, so they do not
182 // meet on disk either.
183 // A stop signal drains: no build is claimed after it, and each build
184 // already in flight runs to completion and is reported. The old
185 // behaviour was to die mid-build, which left the build "running" on
186 // the server with nothing executing (#179). The unit's
187 // TimeoutStopSec bounds the drain; a second signal ends it now.
188 stop := make(chan struct{})
189 go func() {
190 sigs := make(chan os.Signal, 2)
191 signal.Notify(sigs, syscall.SIGTERM, syscall.SIGINT)
192 <-sigs
193 log.Printf("draining: finishing builds in flight, claiming no more")
194 close(stop)
195 <-sigs
196 log.Printf("second signal: exiting without draining")
197 os.Exit(1)
198 }()
199 r.serve(n, *once, *poll, stop)
200}
201
202// serve runs n workers until stop closes. A worker checks stop only
203// between builds, so closing it never interrupts one.
204func (r *runner) serve(n int, once bool, poll time.Duration, stop <-chan struct{}) {
205 if r.stepFn == nil {
206 r.stepFn = r.step
207 }
208 var wg sync.WaitGroup
209 for i := 0; i < n; i++ {
210 wg.Add(1)
211 go func(i int) {
212 defer wg.Done()
213 if n > 1 {
214 time.Sleep(time.Duration(i) * poll / time.Duration(n))
215 }
216 for {
217 select {
218 case <-stop:
219 return
220 default:
221 }
222 ran, err := r.stepFn()
223 if err != nil {
224 log.Printf("runner: %v", err)
225 }
226 if once {
227 return
228 }
229 if !ran {
230 select {
231 case <-stop:
232 return
233 case <-time.After(poll):
234 }
235 }
236 }
237 }(i)
238 }
239 wg.Wait()
240}
241
242// step claims and executes at most one build. ran reports whether there was
243// one, so the caller knows when to idle.
244func (r *runner) step() (bool, error) {
245 args := []string{"runner", "next"}
246 if r.untrusted {
247 args = append(args, "--untrusted")
248 }
249 args = append(append(args, r.repos...), "--json")
250 out, err := r.ssh(nil, args...)
251 if err != nil {
252 return false, fmt.Errorf("claiming build: %w (%s)", err, out)
253 }
254 var env struct {
255 Data job `json:"data"`
256 }
257 if err := json.Unmarshal([]byte(out), &env); err != nil {
258 return false, fmt.Errorf("parsing job: %w", err)
259 }
260 if env.Data.ID == 0 {
261 return false, nil
262 }
263 j := env.Data
264 log.Printf("build %d: %s %s @ %.10s", j.ID, j.Repo, j.Job, j.SHA)
265 status := "failure"
266 if r.run(j) {
267 status = "success"
268 }
269 if err := r.reportDone(j.ID, status); err != nil {
270 return true, err
271 }
272 log.Printf("build %d: %s", j.ID, status)
273 return true, nil
274}
275
276// logSink forwards a build's output to the server and swallows any error
277// doing so. os/exec surfaces a write failure on a step's stdout through
278// cmd.Wait(), so a sink that can fail is a sink that can fail the build it
279// was only recording — a restart or a dropped session used to turn a green
280// suite red, with the explaining line written to the same dead pipe. Losing
281// log lines is the acceptable failure here; losing the build is not.
282type logSink struct {
283 mu sync.Mutex
284 w io.Writer // nil once a write has failed
285}
286
287func (s *logSink) Write(p []byte) (int, error) {
288 s.mu.Lock()
289 defer s.mu.Unlock()
290 if s.w != nil {
291 if _, err := s.w.Write(p); err != nil {
292 s.w = nil
293 }
294 }
295 return len(p), nil
296}
297
298// broken reports whether the stream was lost, so a build can say its log is
299// incomplete rather than appear to have simply stopped.
300func (s *logSink) broken() bool {
301 s.mu.Lock()
302 defer s.mu.Unlock()
303 return s.w == nil
304}
305
306// run clones, checks out, and executes the steps, streaming output to the
307// server. Returns whether every step succeeded.
308func (r *runner) run(j job) bool {
309 dir := filepath.Join(r.workdir, fmt.Sprintf("build-%d", j.ID))
310 defer os.RemoveAll(dir)
311
312 // A build's HOME. Not the workspace, which is removed after every
313 // build: the Go module cache, the sonar scanner and every other tool
314 // cache live under HOME, so a per-build one re-downloads the world
315 // each time. Not the runner's own home either, where its SSH key and
316 // credential dotfiles are. A directory beside the workspaces is
317 // neither.
318 //
319 // One per repository: shared across repositories, a step could poison
320 // the module cache or plant a .gitconfig that another repository's
321 // build would honour, and the container mounts the home read-write
322 // (#184).
323 buildHome, err := buildHomeFor(r.workdir, j.Repo)
324 if err != nil {
325 log.Printf("build %d: build home: %v", j.ID, err)
326 return false
327 }
328
329 // One long-lived `runner log` session receives the whole stream.
330 logCmd := exec.Command(toolpath.Look("ssh"), append(r.sshOpts, r.remote, "runner", "log", fmt.Sprint(j.ID))...)
331 pipe, err := logCmd.StdinPipe()
332 if err != nil {
333 log.Printf("build %d: log pipe: %v", j.ID, err)
334 return false
335 }
336 sink := &logSink{w: pipe}
337 logCmd.Stdout, logCmd.Stderr = io.Discard, io.Discard
338 if err := logCmd.Start(); err != nil {
339 log.Printf("build %d: log stream: %v", j.ID, err)
340 return false
341 }
342 // The server ends the log session with exit 3 when the build is
343 // cancelled; any other end is a lost stream, which the sink absorbs.
344 cancelled := make(chan struct{})
345 logExited := make(chan struct{})
346 go func() {
347 defer close(logExited)
348 err := logCmd.Wait()
349 if ee, ok := err.(*exec.ExitError); ok && ee.ExitCode() == 3 {
350 close(cancelled)
351 return
352 }
353 if err != nil {
354 log.Printf("build %d: log session ended: %v", j.ID, err)
355 }
356 }()
357 // runStep starts cmd and waits for it, the cancel signal, or the
358 // deadline. Every phase goes through it, so a cancel during the clone
359 // lands as fast as one during a step.
360 runStep := func(cmd *exec.Cmd, deadline time.Time) (bool, string) {
361 select {
362 case <-cancelled:
363 return false, "cancelled"
364 default:
365 }
366 ownProcessGroup(cmd)
367 if err := cmd.Start(); err != nil {
368 return false, fmt.Sprintf("start: %v", err)
369 }
370 done := make(chan error, 1)
371 go func() { done <- cmd.Wait() }()
372 // After a kill, Wait returns once every holder of the log pipe is
373 // gone; the group kill makes that prompt, and the cap makes sure a
374 // straggler cannot hold the build open.
375 reap := func() {
376 killTree(cmd)
377 select {
378 case <-done:
379 case <-time.After(10 * time.Second):
380 }
381 }
382 select {
383 case err := <-done:
384 if err != nil {
385 return false, fmt.Sprintf("step failed: %v", err)
386 }
387 return true, ""
388 case <-cancelled:
389 reap()
390 return false, "cancelled"
391 case <-time.After(time.Until(deadline)):
392 reap()
393 return false, fmt.Sprintf("build timed out after %s", r.timeout)
394 }
395 }
396 defer func() {
397 select {
398 case <-cancelled:
399 log.Printf("build %d: cancelled", j.ID)
400 default:
401 if sink.broken() {
402 log.Printf("build %d: log stream lost; stored log is incomplete", j.ID)
403 }
404 }
405 pipe.Close()
406 <-logExited
407 }()
408
409 gitSSH := strings.TrimSpace("ssh " + strings.Join(r.sshOpts, " "))
410 cloneURL := r.cloneBase + "/" + j.Repo + ".git"
411 deadline := time.Now().Add(r.timeout)
412 fmt.Fprintf(sink, "$ git clone %s (%.10s)\n", cloneURL, j.SHA)
413 // A merge request head lives under refs/merge-requests/, which a
414 // clone does not fetch; ask for the ref before checking out.
415 steps := [][]string{{"clone", "-q", cloneURL, dir}}
416 if strings.HasPrefix(j.Ref, "refs/") {
417 steps = append(steps, []string{"-C", dir, "fetch", "-q", "origin", j.Ref})
418 }
419 steps = append(steps, []string{"-C", dir, "checkout", "-q", j.SHA})
420 for _, args := range steps {
421 cmd := exec.Command(toolpath.Look("git"), args...)
422 cmd.Env = append(os.Environ(), "GIT_SSH_COMMAND="+gitSSH, "GIT_TERMINAL_PROMPT=0")
423 cmd.Stdout, cmd.Stderr = sink, sink
424 if ok, why := runStep(cmd, deadline); !ok {
425 fmt.Fprintf(sink, "git %s: %s\n", args[0], why)
426 return false
427 }
428 }
429
430 env := stepEnv(j, buildHome)
431 return r.runSteps(j, dir, env, sink, deadline, runStep)
432}
433
434// stepEnv builds the environment a build step runs with. It is
435// constructed, not inherited: os.Environ() would hand repository content
436// the runner's entire environment, including anything an operator set on
437// the service (#144).
438//
439// HOME is the repository's build home, not the runner's own: tools read
440// credentials out of dotfiles — .netrc, .npmrc, .gitconfig — and a build
441// has no business finding the runner's. It is not the workspace either,
442// because the workspace is deleted after every build and every tool
443// cache lives under HOME.
444//
445// PATH is the one thing carried over: without it a step cannot find the
446// tools the host was provisioned with.
447// buildHomeFor is the build home for one repository: <workdir>/home/<owner>/<name>,
448// created on first use. The repository path comes from the server, but a
449// home must still never resolve outside the home root.
450func buildHomeFor(workdir, repo string) (string, error) {
451 root := filepath.Join(workdir, "home")
452 dir := filepath.Join(root, filepath.FromSlash(repo))
453 if rel, err := filepath.Rel(root, dir); err != nil || rel == "." || strings.HasPrefix(rel, "..") {
454 return "", fmt.Errorf("repository path %q escapes the build home root", repo)
455 }
456 if err := os.MkdirAll(dir, 0o700); err != nil {
457 return "", err
458 }
459 return dir, nil
460}
461
462func stepEnv(j job, home string) []string {
463 path := os.Getenv("PATH")
464 if path == "" {
465 path = "/usr/local/sbin:/usr/local/bin:/usr/sbin:/usr/bin:/sbin:/bin"
466 }
467 env := []string{
468 "PATH=" + path,
469 "HOME=" + home,
470 "LANG=C.UTF-8",
471 "CI=true",
472 "GITBAY_REPO=" + j.Repo,
473 "GITBAY_SHA=" + j.SHA,
474 "GITBAY_REF=" + j.Ref,
475 "GITBAY_JOB=" + j.Job,
476 }
477 // The server sends secrets only for a trusted build — a merge request
478 // head from a fork arrives with none — so this loop is empty exactly
479 // when it should be.
480 for name, value := range j.Secrets {
481 env = append(env, name+"="+value)
482 }
483 return env
484}
485
486// ssh runs one control command against the server and returns stdout.// ssh runs one control command against the server and returns stdout.
487func (r *runner) ssh(stdin io.Reader, args ...string) (string, error) {
488 cmd := exec.Command(toolpath.Look("ssh"), append(append(r.sshOpts, r.remote), args...)...)
489 if stdin != nil {
490 cmd.Stdin = stdin
491 }
492 var out, errOut strings.Builder
493 cmd.Stdout, cmd.Stderr = &out, &errOut
494 if err := cmd.Run(); err != nil {
495 return out.String() + errOut.String(), err
496 }
497 return out.String(), nil
498}
499
500func fileExists(p string) bool { _, err := os.Stat(p); return err == nil }
501
502// defaultWorkdir picks a build workspace that another local user cannot
503// have created first.
504//
505// The default used to be <tmp>/gitbay-runner: a fixed name inside a
506// world-writable directory, created with MkdirAll, which succeeds against
507// an existing directory whoever owns it. On a shared host another user
508// could have made it — or symlinked it — before the runner started, and
509// this is the process that clones repositories and exports build secrets
510// into step environments (go:S5445, #153).
511//
512// The user's cache directory is not world-writable and is per-user by
513// construction. Falling back to tmp keeps a runner working where HOME is
514// unset, and checkWorkdir refuses the unsafe cases there.
515func defaultWorkdir() string {
516 if cache, err := os.UserCacheDir(); err == nil && cache != "" {
517 return filepath.Join(cache, "gitbay-runner")
518 }
519 return filepath.Join(os.TempDir(), "gitbay-runner")
520}
521
522// checkWorkdir makes sure the workspace is a directory this user owns
523// privately. MkdirAll is happy with one that already exists, so being
524// able to create it proves nothing about who made it.
525//
526// A directory we own that is merely too permissive is tightened rather
527// than refused: every runner before this one created its workspace 0755,
528// so refusing would take the runner down on upgrade to fix a permission
529// we are entitled to change. What cannot be repaired — a symlink, or
530// something owned by someone else — is refused, because those are what an
531// attacker leaves behind and neither is ours to correct.
532func checkWorkdir(dir string) error {
533 fi, err := os.Lstat(dir)
534 if err != nil {
535 return err
536 }
537 if fi.Mode()&os.ModeSymlink != 0 {
538 return fmt.Errorf("workdir %s is a symlink; point -workdir at a real directory", dir)
539 }
540 if !fi.IsDir() {
541 return fmt.Errorf("workdir %s is not a directory", dir)
542 }
543 if st, ok := fi.Sys().(*syscall.Stat_t); ok && int(st.Uid) != os.Getuid() {
544 return fmt.Errorf("workdir %s is owned by uid %d, not this process's %d", dir, st.Uid, os.Getuid())
545 }
546 if perm := fi.Mode().Perm(); perm&0o077 != 0 {
547 log.Printf("workdir %s was mode %04o; tightening to 0700 (builds and their secrets are this user's alone)", dir, perm)
548 if err := os.Chmod(dir, 0o700); err != nil {
549 return fmt.Errorf("tightening workdir %s: %w", dir, err)
550 }
551 }
552 return nil
553}