cmd/gitbay-runner/main.go

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