cmd/gitbay-runner/main.go

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

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