cmd/gitbay-runner/main.go
390 lines · 13018 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 for _, step := range j.Steps {
308 fmt.Fprintf(sink, "$ %s\n", step)
309 cmd := exec.Command(toolpath.Look("sh"), "-c", step)
310 cmd.Dir = dir
311 cmd.Env = append(os.Environ(),
312 "GITBAY_REPO="+j.Repo, "GITBAY_SHA="+j.SHA, "GITBAY_REF="+j.Ref, "GITBAY_JOB="+j.Job, "CI=true")
313 for name, value := range j.Secrets {
314 cmd.Env = append(cmd.Env, name+"="+value)
315 }
316 cmd.Stdout, cmd.Stderr = sink, sink
317 if ok, why := runStep(cmd, deadline); !ok {
318 fmt.Fprintf(sink, "%s\n", why)
319 return false
320 }
321 }
322 return true
323}
324
325// ssh runs one control command against the server and returns stdout.
326func (r *runner) ssh(stdin io.Reader, args ...string) (string, error) {
327 cmd := exec.Command(toolpath.Look("ssh"), append(append(r.sshOpts, r.remote), args...)...)
328 if stdin != nil {
329 cmd.Stdin = stdin
330 }
331 var out, errOut strings.Builder
332 cmd.Stdout, cmd.Stderr = &out, &errOut
333 if err := cmd.Run(); err != nil {
334 return out.String() + errOut.String(), err
335 }
336 return out.String(), nil
337}
338
339// defaultWorkdir picks a build workspace that another local user cannot
340// have created first.
341//
342// The default used to be <tmp>/gitbay-runner: a fixed name inside a
343// world-writable directory, created with MkdirAll, which succeeds against
344// an existing directory whoever owns it. On a shared host another user
345// could have made it — or symlinked it — before the runner started, and
346// this is the process that clones repositories and exports build secrets
347// into step environments (go:S5445, #153).
348//
349// The user's cache directory is not world-writable and is per-user by
350// construction. Falling back to tmp keeps a runner working where HOME is
351// unset, and checkWorkdir refuses the unsafe cases there.
352func defaultWorkdir() string {
353 if cache, err := os.UserCacheDir(); err == nil && cache != "" {
354 return filepath.Join(cache, "gitbay-runner")
355 }
356 return filepath.Join(os.TempDir(), "gitbay-runner")
357}
358
359// checkWorkdir makes sure the workspace is a directory this user owns
360// privately. MkdirAll is happy with one that already exists, so being
361// able to create it proves nothing about who made it.
362//
363// A directory we own that is merely too permissive is tightened rather
364// than refused: every runner before this one created its workspace 0755,
365// so refusing would take the runner down on upgrade to fix a permission
366// we are entitled to change. What cannot be repaired — a symlink, or
367// something owned by someone else — is refused, because those are what an
368// attacker leaves behind and neither is ours to correct.
369func checkWorkdir(dir string) error {
370 fi, err := os.Lstat(dir)
371 if err != nil {
372 return err
373 }
374 if fi.Mode()&os.ModeSymlink != 0 {
375 return fmt.Errorf("workdir %s is a symlink; point -workdir at a real directory", dir)
376 }
377 if !fi.IsDir() {
378 return fmt.Errorf("workdir %s is not a directory", dir)
379 }
380 if st, ok := fi.Sys().(*syscall.Stat_t); ok && int(st.Uid) != os.Getuid() {
381 return fmt.Errorf("workdir %s is owned by uid %d, not this process's %d", dir, st.Uid, os.Getuid())
382 }
383 if perm := fi.Mode().Perm(); perm&0o077 != 0 {
384 log.Printf("workdir %s was mode %04o; tightening to 0700 (builds and their secrets are this user's alone)", dir, perm)
385 if err := os.Chmod(dir, 0o700); err != nil {
386 return fmt.Errorf("tightening workdir %s: %w", dir, err)
387 }
388 }
389 return nil
390}