cmd/gitbay-runner/main.go

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

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