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