cmd/gitbay-runner/main.go

329 lines · 10358 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	"time"
 24
 25	"gitbay.org/gitbay/internal/buildinfo"
 26)
 27
 28type job struct {
 29	ID      int64             `json:"id"`
 30	Repo    string            `json:"repo"`
 31	Number  int64             `json:"number"`
 32	Job     string            `json:"job"`
 33	SHA     string            `json:"sha"`
 34	Ref     string            `json:"ref"`
 35	Steps   []string          `json:"steps"`
 36	Secrets map[string]string `json:"secrets"`
 37}
 38
 39type runner struct {
 40	remote    string // ssh destination, e.g. git@gitbay.org
 41	sshOpts   []string
 42	cloneBase string // e.g. ssh://git@gitbay.org
 43	workdir   string
 44	timeout   time.Duration
 45	// repos limits which repositories this runner claims builds for. Empty
 46	// means any, which is what a runner on the server itself wants; a runner
 47	// somewhere that should not execute every repository's steps names them.
 48	repos []string
 49}
 50
 51func main() {
 52	var (
 53		remote    = flag.String("remote", "git@gitbay.org", "ssh destination of the gitbay server")
 54		sshOpts   = flag.String("ssh-opts", "", "extra ssh options, space-separated (also used for git clone)")
 55		cloneBase = flag.String("clone-base", "", "clone URL prefix (default ssh://<remote>)")
 56		workdir   = flag.String("workdir", filepath.Join(os.TempDir(), "gitbay-runner"), "build workspace root")
 57		poll      = flag.Duration("poll", 5*time.Second, "idle poll interval")
 58		timeout   = flag.Duration("timeout", 30*time.Minute, "per-build time limit")
 59		repos     = flag.String("repos", "", "only claim builds for these repositories, comma-separated owner/name (default: any)")
 60		once      = flag.Bool("once", false, "process at most one build, then exit")
 61		jobs      = flag.Int("jobs", 1, "builds to run at once")
 62		version   = flag.Bool("version", false, "print the commit this binary was built from, then exit")
 63	)
 64	flag.Parse()
 65	if *version {
 66		fmt.Println(buildinfo.String())
 67		return
 68	}
 69	// The runner links internal/store, so it goes stale on changes that never
 70	// touch cmd/gitbay-runner. Say which commit is running.
 71	log.Printf("gitbay-runner %s", buildinfo.String())
 72	r := &runner{
 73		remote:    *remote,
 74		cloneBase: *cloneBase,
 75		workdir:   *workdir,
 76		timeout:   *timeout,
 77	}
 78	if *sshOpts != "" {
 79		r.sshOpts = strings.Fields(*sshOpts)
 80	}
 81	for _, name := range strings.Split(*repos, ",") {
 82		if name = strings.TrimSpace(name); name != "" {
 83			r.repos = append(r.repos, name)
 84		}
 85	}
 86	if r.cloneBase == "" {
 87		r.cloneBase = "ssh://" + *remote
 88	}
 89	if err := os.MkdirAll(r.workdir, 0o755); err != nil {
 90		log.Fatal(err)
 91	}
 92	n := *jobs
 93	if n < 1 {
 94		log.Fatal("-jobs must be at least 1")
 95	}
 96	if *once {
 97		// "at most one build" is one build, whatever -jobs says.
 98		n = 1
 99	}
100	// `runner next` claims inside one transaction, so several workers
101	// claiming at once is already safe; the runner just never used that.
102	// Each build works in its own build-<id> directory, so they do not
103	// meet on disk either.
104	var wg sync.WaitGroup
105	for i := 0; i < n; i++ {
106		wg.Add(1)
107		go func(i int) {
108			defer wg.Done()
109			// Spread the idle polls across the interval rather than
110			// having every worker wake together: n workers asking the
111			// same question in the same instant is n times the load for
112			// one answer.
113			if n > 1 {
114				time.Sleep(time.Duration(i) * *poll / time.Duration(n))
115			}
116			for {
117				ran, err := r.step()
118				if err != nil {
119					log.Printf("runner: %v", err)
120				}
121				if *once {
122					return
123				}
124				if !ran {
125					time.Sleep(*poll)
126				}
127			}
128		}(i)
129	}
130	wg.Wait()
131}
132
133// step claims and executes at most one build. ran reports whether there was
134// one, so the caller knows when to idle.
135func (r *runner) step() (bool, error) {
136	out, err := r.ssh(nil, append([]string{"runner", "next"}, append(r.repos, "--json")...)...)
137	if err != nil {
138		return false, fmt.Errorf("claiming build: %w (%s)", err, out)
139	}
140	var env struct {
141		Data job `json:"data"`
142	}
143	if err := json.Unmarshal([]byte(out), &env); err != nil {
144		return false, fmt.Errorf("parsing job: %w", err)
145	}
146	if env.Data.ID == 0 {
147		return false, nil
148	}
149	j := env.Data
150	log.Printf("build %d: %s %s @ %.10s", j.ID, j.Repo, j.Job, j.SHA)
151	status := "failure"
152	if r.run(j) {
153		status = "success"
154	}
155	if out, err := r.ssh(nil, "runner", "done", fmt.Sprint(j.ID), status); err != nil {
156		return true, fmt.Errorf("reporting build %d: %w (%s)", j.ID, err, out)
157	}
158	log.Printf("build %d: %s", j.ID, status)
159	return true, nil
160}
161
162// logSink forwards a build's output to the server and swallows any error
163// doing so. os/exec surfaces a write failure on a step's stdout through
164// cmd.Wait(), so a sink that can fail is a sink that can fail the build it
165// was only recording — a restart or a dropped session used to turn a green
166// suite red, with the explaining line written to the same dead pipe. Losing
167// log lines is the acceptable failure here; losing the build is not.
168type logSink struct {
169	mu sync.Mutex
170	w  io.Writer // nil once a write has failed
171}
172
173func (s *logSink) Write(p []byte) (int, error) {
174	s.mu.Lock()
175	defer s.mu.Unlock()
176	if s.w != nil {
177		if _, err := s.w.Write(p); err != nil {
178			s.w = nil
179		}
180	}
181	return len(p), nil
182}
183
184// broken reports whether the stream was lost, so a build can say its log is
185// incomplete rather than appear to have simply stopped.
186func (s *logSink) broken() bool {
187	s.mu.Lock()
188	defer s.mu.Unlock()
189	return s.w == nil
190}
191
192// run clones, checks out, and executes the steps, streaming output to the
193// server. Returns whether every step succeeded.
194func (r *runner) run(j job) bool {
195	dir := filepath.Join(r.workdir, fmt.Sprintf("build-%d", j.ID))
196	defer os.RemoveAll(dir)
197
198	// One long-lived `runner log` session receives the whole stream.
199	logCmd := exec.Command("ssh", append(r.sshOpts, r.remote, "runner", "log", fmt.Sprint(j.ID))...)
200	pipe, err := logCmd.StdinPipe()
201	if err != nil {
202		log.Printf("build %d: log pipe: %v", j.ID, err)
203		return false
204	}
205	sink := &logSink{w: pipe}
206	logCmd.Stdout, logCmd.Stderr = io.Discard, io.Discard
207	if err := logCmd.Start(); err != nil {
208		log.Printf("build %d: log stream: %v", j.ID, err)
209		return false
210	}
211	// The server ends the log session with exit 3 when the build is
212	// cancelled; any other end is a lost stream, which the sink absorbs.
213	cancelled := make(chan struct{})
214	logExited := make(chan struct{})
215	go func() {
216		defer close(logExited)
217		err := logCmd.Wait()
218		if ee, ok := err.(*exec.ExitError); ok && ee.ExitCode() == 3 {
219			close(cancelled)
220			return
221		}
222		if err != nil {
223			log.Printf("build %d: log session ended: %v", j.ID, err)
224		}
225	}()
226	// runStep starts cmd and waits for it, the cancel signal, or the
227	// deadline. Every phase goes through it, so a cancel during the clone
228	// lands as fast as one during a step.
229	runStep := func(cmd *exec.Cmd, deadline time.Time) (bool, string) {
230		select {
231		case <-cancelled:
232			return false, "cancelled"
233		default:
234		}
235		ownProcessGroup(cmd)
236		if err := cmd.Start(); err != nil {
237			return false, fmt.Sprintf("start: %v", err)
238		}
239		done := make(chan error, 1)
240		go func() { done <- cmd.Wait() }()
241		// After a kill, Wait returns once every holder of the log pipe is
242		// gone; the group kill makes that prompt, and the cap makes sure a
243		// straggler cannot hold the build open.
244		reap := func() {
245			killTree(cmd)
246			select {
247			case <-done:
248			case <-time.After(10 * time.Second):
249			}
250		}
251		select {
252		case err := <-done:
253			if err != nil {
254				return false, fmt.Sprintf("step failed: %v", err)
255			}
256			return true, ""
257		case <-cancelled:
258			reap()
259			return false, "cancelled"
260		case <-time.After(time.Until(deadline)):
261			reap()
262			return false, fmt.Sprintf("build timed out after %s", r.timeout)
263		}
264	}
265	defer func() {
266		select {
267		case <-cancelled:
268			log.Printf("build %d: cancelled", j.ID)
269		default:
270			if sink.broken() {
271				log.Printf("build %d: log stream lost; stored log is incomplete", j.ID)
272			}
273		}
274		pipe.Close()
275		<-logExited
276	}()
277
278	gitSSH := strings.TrimSpace("ssh " + strings.Join(r.sshOpts, " "))
279	cloneURL := r.cloneBase + "/" + j.Repo + ".git"
280	deadline := time.Now().Add(r.timeout)
281	fmt.Fprintf(sink, "$ git clone %s (%.10s)\n", cloneURL, j.SHA)
282	// A merge request head lives under refs/merge-requests/, which a
283	// clone does not fetch; ask for the ref before checking out.
284	steps := [][]string{{"clone", "-q", cloneURL, dir}}
285	if strings.HasPrefix(j.Ref, "refs/") {
286		steps = append(steps, []string{"-C", dir, "fetch", "-q", "origin", j.Ref})
287	}
288	steps = append(steps, []string{"-C", dir, "checkout", "-q", j.SHA})
289	for _, args := range steps {
290		cmd := exec.Command("git", args...)
291		cmd.Env = append(os.Environ(), "GIT_SSH_COMMAND="+gitSSH, "GIT_TERMINAL_PROMPT=0")
292		cmd.Stdout, cmd.Stderr = sink, sink
293		if ok, why := runStep(cmd, deadline); !ok {
294			fmt.Fprintf(sink, "git %s: %s\n", args[0], why)
295			return false
296		}
297	}
298
299	for _, step := range j.Steps {
300		fmt.Fprintf(sink, "$ %s\n", step)
301		cmd := exec.Command("sh", "-c", step)
302		cmd.Dir = dir
303		cmd.Env = append(os.Environ(),
304			"GITBAY_REPO="+j.Repo, "GITBAY_SHA="+j.SHA, "GITBAY_REF="+j.Ref, "GITBAY_JOB="+j.Job, "CI=true")
305		for name, value := range j.Secrets {
306			cmd.Env = append(cmd.Env, name+"="+value)
307		}
308		cmd.Stdout, cmd.Stderr = sink, sink
309		if ok, why := runStep(cmd, deadline); !ok {
310			fmt.Fprintf(sink, "%s\n", why)
311			return false
312		}
313	}
314	return true
315}
316
317// ssh runs one control command against the server and returns stdout.
318func (r *runner) ssh(stdin io.Reader, args ...string) (string, error) {
319	cmd := exec.Command("ssh", append(append(r.sshOpts, r.remote), args...)...)
320	if stdin != nil {
321		cmd.Stdin = stdin
322	}
323	var out, errOut strings.Builder
324	cmd.Stdout, cmd.Stderr = &out, &errOut
325	if err := cmd.Run(); err != nil {
326		return out.String() + errOut.String(), err
327	}
328	return out.String(), nil
329}