runner: a dropped log stream no longer fails the build it was recording !148

merged merged by cmc on 2026-09-01 07:19 UTC · krz/gitbay:fix-67 into main

5 files changed, +187 −5

Layout: unified · split

cmd/gitbay-runner/main.go +37 −2
@@ -19,6 +19,7 @@ import (
1919 "os/exec"
2020 "path/filepath"
2121 "strings"
22 "sync"
2223 "time"
2324
2425 "gitbay.org/gitbay/internal/buildinfo"
@@ -130,6 +131,36 @@ func (r *runner) step() (bool, error) {
130131 return true, nil
131132}
132133
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.
140type logSink struct {
141 mu sync.Mutex
142 w io.Writer // nil once a write has failed
143}
144
145func (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.
158func (s *logSink) broken() bool {
159 s.mu.Lock()
160 defer s.mu.Unlock()
161 return s.w == nil
162}
163
133164// run clones, checks out, and executes the steps, streaming output to the
134165// server. Returns whether every step succeeded.
135166func (r *runner) run(j job) bool {
@@ -138,18 +169,22 @@ func (r *runner) run(j job) bool {
138169
139170 // One long-lived `runner log` session receives the whole stream.
140171 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()
142173 if err != nil {
143174 log.Printf("build %d: log pipe: %v", j.ID, err)
144175 return false
145176 }
177 sink := &logSink{w: pipe}
146178 logCmd.Stdout, logCmd.Stderr = io.Discard, io.Discard
147179 if err := logCmd.Start(); err != nil {
148180 log.Printf("build %d: log stream: %v", j.ID, err)
149181 return false
150182 }
151183 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()
153188 logCmd.Wait()
154189 }()
155190
cmd/gitbay-runner/main_test.go added +74
@@ -0,0 +1,74 @@
1package main
2
3import (
4 "errors"
5 "os/exec"
6 "strings"
7 "testing"
8)
9
10type failingWriter struct{ n int }
11
12func (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.
21func 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.
35func 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.
46func 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.
60func 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 {
385385 if err != nil {
386386 return c.fail(protocol.ExitUsage, "bad build id %q", args[0])
387387 }
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.
389392 buf := make([]byte, 64<<10)
393 dropped := 0
390394 for {
391395 n, rerr := c.Stdin.Read(buf)
392396 if n > 0 {
393397 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)
395400 }
396401 }
397402 if rerr != nil {
398403 break
399404 }
400405 }
406 if dropped > 0 {
407 slog.Warn("build log incomplete", "build", id, "dropped_chunks", dropped)
408 }
401409 return c.emit(map[string]string{"log": "ok"}, func(w io.Writer) {})
402410}
403411
internal/store/builds.go +20 −1
@@ -3,6 +3,7 @@ package store
33import (
44 "database/sql"
55 "errors"
6 "strconv"
67 "strings"
78 "time"
89)
@@ -25,6 +26,12 @@ type Build struct {
2526// MaxBuildLog caps a build's stored log; appends past it are dropped.
2627const MaxBuildLog = 2 << 20
2728
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.
32var truncNotice = []byte("\n[log truncated: reached the " +
33 strconv.Itoa(MaxBuildLog>>20) + " MiB cap; earlier output is above]\n")
34
2835// CreateBuild allocates the per-repo build number in the same transaction
2936// as the insert, like issue and MR numbers.
3037func (s *Store) CreateBuild(repoID int64, job, sha, ref, stepsJSON string) (int64, error) {
@@ -142,9 +149,21 @@ func (s *Store) ReapStaleBuilds() ([]Build, error) {
142149
143150// AppendBuildLog adds a chunk to the build's log, dropping bytes past the cap.
144151func (s *Store) AppendBuildLog(id int64, chunk []byte) error {
145 _, err := s.DB.Exec(`
152 res, err := s.DB.Exec(`
146153 UPDATE builds SET log = log || ?
147154 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))
148167 return err
149168}
150169
internal/store/builds_test.go +46
@@ -166,3 +166,49 @@ func TestClaimBuildScopedToRepos(t *testing.T) {
166166 t.Fatalf("unscoped claim: err=%v ok=%v repo=%d", err, ok, b.RepoID)
167167 }
168168}
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.
173func 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}