cmd/gitbay-runner/main.go

88fc476788ab9c640649427d815217360569b4d9
gitbay/cmd/gitbay-runner/main.go history · blame · raw

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