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}