Commit 3a243db45e
Verified · cmc ci/build: success ci/test: success ci/vuln: success
Layout: unified · split
cmd/gitbay-runner/main.go +37 −2
| @@ -19,6 +19,7 @@ import ( | |||
| 19 | "os/exec" | 19 | "os/exec" |
| 20 | "path/filepath" | 20 | "path/filepath" |
| 21 | "strings" | 21 | "strings" |
| 22 | "sync" | ||
| 22 | "time" | 23 | "time" |
| 23 | 24 | ||
| 24 | "gitbay.org/gitbay/internal/buildinfo" | 25 | "gitbay.org/gitbay/internal/buildinfo" |
| @@ -130,6 +131,36 @@ func (r *runner) step() (bool, error) { | |||
| 130 | return true, nil | 131 | return true, nil |
| 131 | } | 132 | } |
| 132 | 133 | ||
| 134 | // logSink forwards a build's output to the server and swallows any error | ||
| 135 | // doing so. os/exec surfaces a write failure on a step's stdout through | ||
| 136 | // cmd.Wait(), so a sink that can fail is a sink that can fail the build it | ||
| 137 | // was only recording — a restart or a dropped session used to turn a green | ||
| 138 | // suite red, with the explaining line written to the same dead pipe. Losing | ||
| 139 | // log lines is the acceptable failure here; losing the build is not. | ||
| 140 | type logSink struct { | ||
| 141 | mu sync.Mutex | ||
| 142 | w io.Writer // nil once a write has failed | ||
| 143 | } | ||
| 144 | |||
| 145 | func (s *logSink) Write(p []byte) (int, error) { | ||
| 146 | s.mu.Lock() | ||
| 147 | defer s.mu.Unlock() | ||
| 148 | if s.w != nil { | ||
| 149 | if _, err := s.w.Write(p); err != nil { | ||
| 150 | s.w = nil | ||
| 151 | } | ||
| 152 | } | ||
| 153 | return len(p), nil | ||
| 154 | } | ||
| 155 | |||
| 156 | // broken reports whether the stream was lost, so a build can say its log is | ||
| 157 | // incomplete rather than appear to have simply stopped. | ||
| 158 | func (s *logSink) broken() bool { | ||
| 159 | s.mu.Lock() | ||
| 160 | defer s.mu.Unlock() | ||
| 161 | return s.w == nil | ||
| 162 | } | ||
| 163 | |||
| 133 | // run clones, checks out, and executes the steps, streaming output to the | 164 | // run clones, checks out, and executes the steps, streaming output to the |
| 134 | // server. Returns whether every step succeeded. | 165 | // server. Returns whether every step succeeded. |
| 135 | func (r *runner) run(j job) bool { | 166 | func (r *runner) run(j job) bool { |
| @@ -138,18 +169,22 @@ func (r *runner) run(j job) bool { | |||
| 138 | 169 | ||
| 139 | // One long-lived `runner log` session receives the whole stream. | 170 | // One long-lived `runner log` session receives the whole stream. |
| 140 | logCmd := exec.Command("ssh", append(r.sshOpts, r.remote, "runner", "log", fmt.Sprint(j.ID))...) | 171 | logCmd := exec.Command("ssh", append(r.sshOpts, r.remote, "runner", "log", fmt.Sprint(j.ID))...) |
| 141 | sink, err := logCmd.StdinPipe() | 172 | pipe, err := logCmd.StdinPipe() |
| 142 | if err != nil { | 173 | if err != nil { |
| 143 | log.Printf("build %d: log pipe: %v", j.ID, err) | 174 | log.Printf("build %d: log pipe: %v", j.ID, err) |
| 144 | return false | 175 | return false |
| 145 | } | 176 | } |
| 177 | sink := &logSink{w: pipe} | ||
| 146 | logCmd.Stdout, logCmd.Stderr = io.Discard, io.Discard | 178 | logCmd.Stdout, logCmd.Stderr = io.Discard, io.Discard |
| 147 | if err := logCmd.Start(); err != nil { | 179 | if err := logCmd.Start(); err != nil { |
| 148 | log.Printf("build %d: log stream: %v", j.ID, err) | 180 | log.Printf("build %d: log stream: %v", j.ID, err) |
| 149 | return false | 181 | return false |
| 150 | } | 182 | } |
| 151 | defer func() { | 183 | defer func() { |
| 152 | sink.Close() | 184 | if sink.broken() { |
| 185 | log.Printf("build %d: log stream lost; stored log is incomplete", j.ID) | ||
| 186 | } | ||
| 187 | pipe.Close() | ||
| 153 | logCmd.Wait() | 188 | logCmd.Wait() |
| 154 | }() | 189 | }() |
| 155 | 190 | ||
cmd/gitbay-runner/main_test.go added +74
| @@ -0,0 +1,74 @@ | |||
| 1 | package main | ||
| 2 | |||
| 3 | import ( | ||
| 4 | "errors" | ||
| 5 | "os/exec" | ||
| 6 | "strings" | ||
| 7 | "testing" | ||
| 8 | ) | ||
| 9 | |||
| 10 | type failingWriter struct{ n int } | ||
| 11 | |||
| 12 | func (f *failingWriter) Write(p []byte) (int, error) { | ||
| 13 | f.n++ | ||
| 14 | return 0, errors.New("broken pipe") | ||
| 15 | } | ||
| 16 | |||
| 17 | // A build's log stream is a recording device. When it breaks, os/exec | ||
| 18 | // surfaces the write error through cmd.Wait(), which used to mark a step | ||
| 19 | // that had exited 0 as failed — a green suite reported red, with the line | ||
| 20 | // explaining it written to the same dead pipe. | ||
| 21 | func TestBrokenLogSinkDoesNotFailTheStep(t *testing.T) { | ||
| 22 | sink := &logSink{w: &failingWriter{}} | ||
| 23 | |||
| 24 | cmd := exec.Command("sh", "-c", "echo out; echo err >&2; exit 0") | ||
| 25 | cmd.Stdout, cmd.Stderr = sink, sink | ||
| 26 | if err := cmd.Run(); err != nil { | ||
| 27 | t.Fatalf("step reported failed because its log sink broke: %v", err) | ||
| 28 | } | ||
| 29 | if !sink.broken() { | ||
| 30 | t.Error("sink did not record that the stream was lost") | ||
| 31 | } | ||
| 32 | } | ||
| 33 | |||
| 34 | // The converse: a step that genuinely fails still fails. | ||
| 35 | func TestGenuineStepFailureStillFails(t *testing.T) { | ||
| 36 | sink := &logSink{w: &failingWriter{}} | ||
| 37 | cmd := exec.Command("sh", "-c", "exit 3") | ||
| 38 | cmd.Stdout, cmd.Stderr = sink, sink | ||
| 39 | if err := cmd.Run(); err == nil { | ||
| 40 | t.Fatal("a step exiting 3 was reported as success") | ||
| 41 | } | ||
| 42 | } | ||
| 43 | |||
| 44 | // Once the stream is gone the sink stops touching it, rather than calling a | ||
| 45 | // broken pipe once per write for the rest of a long build. | ||
| 46 | func TestLogSinkStopsWritingAfterFailure(t *testing.T) { | ||
| 47 | w := &failingWriter{} | ||
| 48 | sink := &logSink{w: w} | ||
| 49 | for i := 0; i < 5; i++ { | ||
| 50 | if _, err := sink.Write([]byte("x")); err != nil { | ||
| 51 | t.Fatalf("sink returned an error: %v", err) | ||
| 52 | } | ||
| 53 | } | ||
| 54 | if w.n != 1 { | ||
| 55 | t.Errorf("underlying writer called %d times, want 1", w.n) | ||
| 56 | } | ||
| 57 | } | ||
| 58 | |||
| 59 | // A healthy sink still forwards everything. | ||
| 60 | func TestHealthyLogSinkForwards(t *testing.T) { | ||
| 61 | var b strings.Builder | ||
| 62 | sink := &logSink{w: &b} | ||
| 63 | cmd := exec.Command("sh", "-c", "echo hello") | ||
| 64 | cmd.Stdout, cmd.Stderr = sink, sink | ||
| 65 | if err := cmd.Run(); err != nil { | ||
| 66 | t.Fatal(err) | ||
| 67 | } | ||
| 68 | if got := b.String(); !strings.Contains(got, "hello") { | ||
| 69 | t.Errorf("sink dropped output: %q", got) | ||
| 70 | } | ||
| 71 | if sink.broken() { | ||
| 72 | t.Error("healthy sink reported broken") | ||
| 73 | } | ||
| 74 | } | ||
internal/control/build.go +10 −2
| @@ -385,19 +385,27 @@ func runRunnerLog(c *Ctx, args []string) int { | |||
| 385 | if err != nil { | 385 | if err != nil { |
| 386 | return c.fail(protocol.ExitUsage, "bad build id %q", args[0]) | 386 | return c.fail(protocol.ExitUsage, "bad build id %q", args[0]) |
| 387 | } | 387 | } |
| 388 | // Stream stdin into the log in chunks so long builds appear live. | 388 | // Stream stdin into the log in chunks so long builds appear live. An |
| 389 | // append that fails drops its chunk and the loop keeps draining: ending | ||
| 390 | // the session here breaks the runner's pipe, and a broken pipe is how a | ||
| 391 | // transient SQLITE_BUSY used to fail the build the log belonged to. | ||
| 389 | buf := make([]byte, 64<<10) | 392 | buf := make([]byte, 64<<10) |
| 393 | dropped := 0 | ||
| 390 | for { | 394 | for { |
| 391 | n, rerr := c.Stdin.Read(buf) | 395 | n, rerr := c.Stdin.Read(buf) |
| 392 | if n > 0 { | 396 | if n > 0 { |
| 393 | if err := c.Store.AppendBuildLog(id, buf[:n]); err != nil { | 397 | if err := c.Store.AppendBuildLog(id, buf[:n]); err != nil { |
| 394 | return c.fail(protocol.ExitFailure, "%v", err) | 398 | dropped++ |
| 399 | slog.Warn("appending build log", "build", id, "err", err) | ||
| 395 | } | 400 | } |
| 396 | } | 401 | } |
| 397 | if rerr != nil { | 402 | if rerr != nil { |
| 398 | break | 403 | break |
| 399 | } | 404 | } |
| 400 | } | 405 | } |
| 406 | if dropped > 0 { | ||
| 407 | slog.Warn("build log incomplete", "build", id, "dropped_chunks", dropped) | ||
| 408 | } | ||
| 401 | return c.emit(map[string]string{"log": "ok"}, func(w io.Writer) {}) | 409 | return c.emit(map[string]string{"log": "ok"}, func(w io.Writer) {}) |
| 402 | } | 410 | } |
| 403 | 411 | ||
internal/store/builds.go +20 −1
| @@ -3,6 +3,7 @@ package store | |||
| 3 | import ( | 3 | import ( |
| 4 | "database/sql" | 4 | "database/sql" |
| 5 | "errors" | 5 | "errors" |
| 6 | "strconv" | ||
| 6 | "strings" | 7 | "strings" |
| 7 | "time" | 8 | "time" |
| 8 | ) | 9 | ) |
| @@ -25,6 +26,12 @@ type Build struct { | |||
| 25 | // MaxBuildLog caps a build's stored log; appends past it are dropped. | 26 | // MaxBuildLog caps a build's stored log; appends past it are dropped. |
| 26 | const MaxBuildLog = 2 << 20 | 27 | const MaxBuildLog = 2 << 20 |
| 27 | 28 | ||
| 29 | // truncNotice is appended once when a log first hits the cap. A log that | ||
| 30 | // simply stops is indistinguishable from a build that died mid-step, which | ||
| 31 | // is the reading that sent people hunting for a nonexistent test failure. | ||
| 32 | var truncNotice = []byte("\n[log truncated: reached the " + | ||
| 33 | strconv.Itoa(MaxBuildLog>>20) + " MiB cap; earlier output is above]\n") | ||
| 34 | |||
| 28 | // CreateBuild allocates the per-repo build number in the same transaction | 35 | // CreateBuild allocates the per-repo build number in the same transaction |
| 29 | // as the insert, like issue and MR numbers. | 36 | // as the insert, like issue and MR numbers. |
| 30 | func (s *Store) CreateBuild(repoID int64, job, sha, ref, stepsJSON string) (int64, error) { | 37 | func (s *Store) CreateBuild(repoID int64, job, sha, ref, stepsJSON string) (int64, error) { |
| @@ -142,9 +149,21 @@ func (s *Store) ReapStaleBuilds() ([]Build, error) { | |||
| 142 | 149 | ||
| 143 | // AppendBuildLog adds a chunk to the build's log, dropping bytes past the cap. | 150 | // AppendBuildLog adds a chunk to the build's log, dropping bytes past the cap. |
| 144 | func (s *Store) AppendBuildLog(id int64, chunk []byte) error { | 151 | func (s *Store) AppendBuildLog(id int64, chunk []byte) error { |
| 145 | _, err := s.DB.Exec(` | 152 | res, err := s.DB.Exec(` |
| 146 | UPDATE builds SET log = log || ? | 153 | UPDATE builds SET log = log || ? |
| 147 | WHERE id = ? AND length(log) < ?`, chunk, id, MaxBuildLog) | 154 | WHERE id = ? AND length(log) < ?`, chunk, id, MaxBuildLog) |
| 155 | if err != nil { | ||
| 156 | return err | ||
| 157 | } | ||
| 158 | if n, _ := res.RowsAffected(); n > 0 { | ||
| 159 | return nil | ||
| 160 | } | ||
| 161 | // Over the cap. The bounds match exactly once: appending the notice puts | ||
| 162 | // the log past the upper bound, so later chunks fall through silently. | ||
| 163 | _, err = s.DB.Exec(` | ||
| 164 | UPDATE builds SET log = log || ? | ||
| 165 | WHERE id = ? AND length(log) >= ? AND length(log) < ?`, | ||
| 166 | truncNotice, id, MaxBuildLog, MaxBuildLog+len(truncNotice)) | ||
| 148 | return err | 167 | return err |
| 149 | } | 168 | } |
| 150 | 169 | ||
internal/store/builds_test.go +46
| @@ -166,3 +166,49 @@ func TestClaimBuildScopedToRepos(t *testing.T) { | |||
| 166 | t.Fatalf("unscoped claim: err=%v ok=%v repo=%d", err, ok, b.RepoID) | 166 | t.Fatalf("unscoped claim: err=%v ok=%v repo=%d", err, ok, b.RepoID) |
| 167 | } | 167 | } |
| 168 | } | 168 | } |
| 169 | |||
| 170 | // A log that stops at the cap reads exactly like a build that died mid-step, | ||
| 171 | // which is what sent people hunting for a test failure that was never there. | ||
| 172 | // It says so instead, once. | ||
| 173 | func TestBuildLogSaysWhenItTruncates(t *testing.T) { | ||
| 174 | s := open(t) | ||
| 175 | if err := s.MigrateUp(); err != nil { | ||
| 176 | t.Fatal(err) | ||
| 177 | } | ||
| 178 | uid, err := s.CreateUser("cmc", true) | ||
| 179 | if err != nil { | ||
| 180 | t.Fatal(err) | ||
| 181 | } | ||
| 182 | if _, err := s.CreateRepo("user", uid, "orgo", "public"); err != nil { | ||
| 183 | t.Fatal(err) | ||
| 184 | } | ||
| 185 | id, err := s.CreateBuild(1, "test", "abc123", "main", `["true"]`) | ||
| 186 | if err != nil { | ||
| 187 | t.Fatal(err) | ||
| 188 | } | ||
| 189 | |||
| 190 | chunk := make([]byte, 256<<10) | ||
| 191 | for i := range chunk { | ||
| 192 | chunk[i] = 'x' | ||
| 193 | } | ||
| 194 | // Well past the cap, so plenty of appends land after it. | ||
| 195 | for written := 0; written < MaxBuildLog+(4*len(chunk)); written += len(chunk) { | ||
| 196 | if err := s.AppendBuildLog(id, chunk); err != nil { | ||
| 197 | t.Fatal(err) | ||
| 198 | } | ||
| 199 | } | ||
| 200 | |||
| 201 | log, err := s.BuildLog(id) | ||
| 202 | if err != nil { | ||
| 203 | t.Fatal(err) | ||
| 204 | } | ||
| 205 | if n := strings.Count(string(log), "log truncated"); n != 1 { | ||
| 206 | t.Errorf("truncation notice appears %d times, want exactly 1", n) | ||
| 207 | } | ||
| 208 | if !strings.HasSuffix(string(log), string(truncNotice)) { | ||
| 209 | t.Error("notice is not at the end of the log") | ||
| 210 | } | ||
| 211 | if len(log) > MaxBuildLog+len(truncNotice)+len(chunk) { | ||
| 212 | t.Errorf("log grew to %d, past the cap plus one chunk", len(log)) | ||
| 213 | } | ||
| 214 | } | ||