build cancel: end a running build at the runner !168

merged merged by cmc on 2026-09-02 04:06 UTC · krz/gitbay:build-cancel-running into main

4 files changed, +171 −31

Layout: unified · split

cmd/gitbay-runner/main.go +24 −3
@@ -180,12 +180,29 @@ func (r *runner) run(j job) bool {
180180 log.Printf("build %d: log stream: %v", j.ID, err)
181181 return false
182182 }
183 // The server ends the log session with exit 3 when the build is
184 // cancelled; any other end is a lost stream, which the sink absorbs.
185 cancelled := make(chan struct{})
186 logExited := make(chan struct{})
187 go func() {
188 defer close(logExited)
189 if err := logCmd.Wait(); err != nil {
190 if ee, ok := err.(*exec.ExitError); ok && ee.ExitCode() == 3 {
191 close(cancelled)
192 }
193 }
194 }()
183195 defer func() {
184 if sink.broken() {
185 log.Printf("build %d: log stream lost; stored log is incomplete", j.ID)
196 select {
197 case <-cancelled:
198 log.Printf("build %d: cancelled", j.ID)
199 default:
200 if sink.broken() {
201 log.Printf("build %d: log stream lost; stored log is incomplete", j.ID)
202 }
186203 }
187204 pipe.Close()
188 logCmd.Wait()
205 <-logExited
189206 }()
190207
191208 gitSSH := strings.TrimSpace("ssh " + strings.Join(r.sshOpts, " "))
@@ -227,6 +244,10 @@ func (r *runner) run(j job) bool {
227244 fmt.Fprintf(sink, "step failed: %v\n", err)
228245 return false
229246 }
247 case <-cancelled:
248 cmd.Process.Kill()
249 <-done
250 return false
230251 case <-time.After(time.Until(deadline)):
231252 cmd.Process.Kill()
232253 fmt.Fprintf(sink, "build timed out after %s\n", r.timeout)
e2e/build_cancel_test.go +80 −7
@@ -1,10 +1,13 @@
11package e2e
22
33import (
4 "fmt"
45 "os"
6 "os/exec"
57 "path/filepath"
68 "strings"
79 "testing"
10 "time"
811)
912
1013// build cancel withdraws a queued build; a running one is the runner's.
@@ -82,14 +85,84 @@ func TestBuildCancel(t *testing.T) {
8285 if st, _, _ := inst.ssh(t, aliceKey, "", "status", "list", "alice/app", sha); !strings.Contains(st, "success") || !strings.Contains(st, "passed in build 2") {
8386 t.Fatalf("status after cancelling the duplicate:\n%s", st)
8487 }
85 // A running build cannot be cancelled here.
86 if _, _, code := inst.ssh(t, aliceKey, "", "build", "trigger", "alice/app", "unit"); code != 0 {
87 t.Fatal("third trigger failed")
88}
89
90// Cancelling a running build ends it at the runner within seconds: the
91// server closes the log session, the runner kills the step, and its late
92// report lands on a row that already says cancelled.
93func TestBuildCancelRunning(t *testing.T) {
94 inst := startInstance(t)
95 inst.runner = buildRunner(t)
96 aliceKey := inst.newKey(t, "alice")
97 runnerKey := inst.newKey(t, "ci")
98 inst.admin(t, "admin", "user", "create", "alice", "--key", aliceKey+".pub")
99 inst.admin(t, "admin", "user", "create", "ci", "--key", runnerKey+".pub", "--admin")
100 if _, _, code := inst.ssh(t, aliceKey, "", "repo", "create", "alice/slow"); code != 0 {
101 t.Fatal("repo create failed")
88102 }
89 if _, _, code := inst.ssh(t, runnerKey, "", "runner", "next"); code != 0 {
90 t.Fatal("claim failed")
103 work := t.TempDir()
104 env := inst.gitEnv(aliceKey)
105 mustGit(t, work, env, "clone", inst.sshURL("alice/slow"), "w")
106 dir := filepath.Join(work, "w")
107 os.MkdirAll(filepath.Join(dir, ".gitbay"), 0o755)
108 os.WriteFile(filepath.Join(dir, ".gitbay", "ci.yml"), []byte("jobs:\n slow:\n steps:\n - echo starting\n - sleep 120\n - echo never\n"), 0o644)
109 mustGit(t, dir, env, "checkout", "-q", "-b", "main")
110 mustGit(t, dir, env, "add", ".")
111 mustGit(t, dir, env, "commit", "-q", "-m", "slow")
112 mustGit(t, dir, env, "push", "-q", "origin", "main")
113 sha := strings.Fields(strings.TrimSpace(mustGit(t, dir, env, "ls-remote", "origin", "refs/heads/main")))[0]
114 slowList := func() string {
115 out, _, _ := inst.ssh(t, aliceKey, "", "build", "list", "alice/slow")
116 return out
91117 }
92 if _, errOut, code := inst.ssh(t, aliceKey, "", "build", "cancel", "alice/app", "4"); code != 2 || !strings.Contains(errOut, "running") {
93 t.Fatalf("cancelled a running build: exit %d %s", code, errOut)
118
119 // A real runner, in the background, claims the build and sits in sleep.
120 opts := fmt.Sprintf("-p %d -i %s -o IdentitiesOnly=yes -o StrictHostKeyChecking=no -o UserKnownHostsFile=%s -o BatchMode=yes",
121 inst.port, runnerKey, filepath.Join(inst.sshDir, "known_hosts"))
122 runner := exec.Command(inst.runner, "-once", "-remote", "git@127.0.0.1", "-ssh-opts", opts,
123 "-clone-base", fmt.Sprintf("ssh://git@127.0.0.1:%d", inst.port), "-workdir", t.TempDir())
124 runner.Env = append(os.Environ(), "GIT_CONFIG_NOSYSTEM=1", "GIT_CONFIG_GLOBAL=/dev/null")
125 var runnerOut strings.Builder
126 runner.Stdout, runner.Stderr = &runnerOut, &runnerOut
127 if err := runner.Start(); err != nil {
128 t.Fatal(err)
129 }
130 exited := make(chan error, 1)
131 go func() { exited <- runner.Wait() }()
132 t.Cleanup(func() { runner.Process.Kill() })
133 deadline := time.Now().Add(30 * time.Second)
134 for !strings.Contains(slowList(), "slow\trunning") {
135 if time.Now().After(deadline) {
136 t.Fatalf("runner never claimed the build:\n%s\nbuild list:\n%s", runnerOut.String(), slowList())
137 }
138 time.Sleep(200 * time.Millisecond)
139 }
140 started := time.Now()
141 out, errOut, code := inst.ssh(t, aliceKey, "", "build", "cancel", "alice/slow", "1")
142 if code != 0 || !strings.Contains(out, "the runner stops at its next check") {
143 t.Fatalf("cancel running: exit %d %s%s", code, out, errOut)
144 }
145 select {
146 case <-exited:
147 case <-time.After(20 * time.Second):
148 t.Fatalf("runner still running 20s after cancel:\n%s", runnerOut.String())
149 }
150 if took := time.Since(started); took > 15*time.Second {
151 t.Fatalf("runner took %s to stop", took)
152 }
153 if !strings.Contains(runnerOut.String(), "cancelled") {
154 t.Fatalf("runner did not say it was cancelled:\n%s", runnerOut.String())
155 }
156 // The row, log and status say cancelled, and the runner's late report
157 // changed none of them.
158 if list := slowList(); !strings.Contains(list, "slow\tcancelled") {
159 t.Fatalf("build after cancel:\n%s", list)
160 }
161 log, _, _ := inst.ssh(t, aliceKey, "", "build", "log", "alice/slow", "1")
162 if !strings.Contains(log, "starting") || !strings.Contains(log, "cancelled by alice while running") || strings.Contains(log, "never") {
163 t.Fatalf("log after cancel:\n%s", log)
164 }
165 if st, _, _ := inst.ssh(t, aliceKey, "", "status", "list", "alice/slow", sha); !strings.Contains(st, "error") || !strings.Contains(st, "cancelled") {
166 t.Fatalf("status after cancel:\n%s", st)
94167 }
95168}
internal/control/build.go +62 −17
@@ -379,24 +379,53 @@ func runRunnerLog(c *Ctx, args []string) int {
379379 // append that fails drops its chunk and the loop keeps draining: ending
380380 // the session here breaks the runner's pipe, and a broken pipe is how a
381381 // transient SQLITE_BUSY used to fail the build the log belonged to.
382 buf := make([]byte, 64<<10)
382 //
383 // The session is also how a running build is cancelled: while it is
384 // open the build's row is watched, and when the row stops saying
385 // running the session ends with ExitNotFound, which the runner reads as
386 // "stop this build". Any other end of the session is a lost stream.
387 type chunk struct {
388 data []byte
389 err error
390 }
391 chunks := make(chan chunk, 4)
392 go func() {
393 buf := make([]byte, 64<<10)
394 for {
395 n, rerr := c.Stdin.Read(buf)
396 if n > 0 {
397 chunks <- chunk{data: append([]byte(nil), buf[:n]...)}
398 }
399 if rerr != nil {
400 chunks <- chunk{err: rerr}
401 return
402 }
403 }
404 }()
405 watch := time.NewTicker(2 * time.Second)
406 defer watch.Stop()
383407 dropped := 0
384408 for {
385 n, rerr := c.Stdin.Read(buf)
386 if n > 0 {
387 if err := c.Store.AppendBuildLog(id, buf[:n]); err != nil {
388 dropped++
389 slog.Warn("appending build log", "build", id, "err", err)
409 select {
410 case ch := <-chunks:
411 if len(ch.data) > 0 {
412 if err := c.Store.AppendBuildLog(id, ch.data); err != nil {
413 dropped++
414 slog.Warn("appending build log", "build", id, "err", err)
415 }
416 }
417 if ch.err != nil {
418 if dropped > 0 {
419 slog.Warn("build log incomplete", "build", id, "dropped_chunks", dropped)
420 }
421 return c.emit(map[string]string{"log": "ok"}, func(w io.Writer) {})
422 }
423 case <-watch.C:
424 if b, err := c.Store.BuildByID(id); err == nil && b.Status != "running" {
425 return c.fail(protocol.ExitNotFound, "build %d is %s; stop", id, b.Status)
390426 }
391427 }
392 if rerr != nil {
393 break
394 }
395 }
396 if dropped > 0 {
397 slog.Warn("build log incomplete", "build", id, "dropped_chunks", dropped)
398428 }
399 return c.emit(map[string]string{"log": "ok"}, func(w io.Writer) {})
400429}
401430
402431func runRunnerDone(c *Ctx, args []string) int {
@@ -414,6 +443,14 @@ func runRunnerDone(c *Ctx, args []string) int {
414443 if err != nil {
415444 return c.fail(protocol.ExitNotFound, "no build %d", id)
416445 }
446 // Cancelled underneath the runner: its report is late, not wrong.
447 // The row, the status and the log were settled by the cancel.
448 if b.Status == "cancelled" {
449 c.Store.RunnerDone(c.User.ID)
450 return c.emit(map[string]any{"build": b.Number, "status": "cancelled"}, func(w io.Writer) {
451 fmt.Fprintf(w, "build %d was cancelled\n", b.Number)
452 })
453 }
417454 if err := c.Store.FinishBuild(id, args[1]); err != nil {
418455 return c.fail(protocol.ExitFailure, "finishing build %d: %v", id, err)
419456 }
@@ -527,13 +564,17 @@ func runBuildCancel(c *Ctx, args []string) int {
527564 if !policy.CanWrite(c.User, repo, grant) {
528565 return c.fail(protocol.ExitDenied, "cancelling a build needs write access to %s", repo.Path())
529566 }
530 if b.Status != "pending" {
531 return c.fail(protocol.ExitUsage, "build %d is %s; only a queued build can be cancelled", b.Number, b.Status)
567 if b.Status != "pending" && b.Status != "running" {
568 return c.fail(protocol.ExitUsage, "build %d is %s; only a queued or running build can be cancelled", b.Number, b.Status)
532569 }
533570 if err := c.Store.CancelBuild(b.ID); err != nil {
534571 return c.fail(protocol.ExitFailure, "%v", err)
535572 }
536 c.Store.AppendBuildLog(b.ID, []byte(fmt.Sprintf("cancelled by %s before a runner claimed it\n", c.User.Username)))
573 if b.Status == "running" {
574 c.Store.AppendBuildLog(b.ID, []byte(fmt.Sprintf("\ncancelled by %s while running; the runner stops at its next check\n", c.User.Username)))
575 } else {
576 c.Store.AppendBuildLog(b.ID, []byte(fmt.Sprintf("cancelled by %s before a runner claimed it\n", c.User.Username)))
577 }
537578 // The queued status replaced whatever the commit had for this job. If
538579 // the commit passed the job on another ref, that result stands again;
539580 // otherwise the context says it was withdrawn.
@@ -546,7 +587,11 @@ func runBuildCancel(c *Ctx, args []string) int {
546587 c.Store.SetCommitStatus(repo.ID, b.SHA, "ci/"+b.Job, "error", "cancelled", url, c.User.ID)
547588 }
548589 c.Store.RecordEvent(repo.ID, c.User.ID, "build.cancelled", fmt.Sprintf(`{"number":%d,"job":%q}`, b.Number, b.Job))
549 return c.emit(map[string]any{"number": b.Number, "job": b.Job, "status": "cancelled"}, func(w io.Writer) {
590 return c.emit(map[string]any{"number": b.Number, "job": b.Job, "status": "cancelled", "was": b.Status}, func(w io.Writer) {
591 if b.Status == "running" {
592 fmt.Fprintf(w, "cancelled %s build %d (%s); the runner stops at its next check\n", repo.Path(), b.Number, b.Job)
593 return
594 }
550595 fmt.Fprintf(w, "cancelled %s build %d (%s)\n", repo.Path(), b.Number, b.Job)
551596 })
552597}
internal/store/builds.go +5 −4
@@ -291,12 +291,13 @@ func (b Build) Elapsed() time.Duration {
291291 return 0
292292}
293293
294// CancelBuild withdraws a build that no runner has claimed. A running
295// build is the runner's to finish; cancelling it here would leave the
296// runner reporting on a row that says otherwise.
294// CancelBuild withdraws a queued or running build. A running one is
295// ended by the runner, which learns of the cancellation when its log
296// session is closed, and whose later report lands on a row that already
297// says cancelled.
297298func (s *Store) CancelBuild(id int64) error {
298299 res, err := s.DB.Exec(`UPDATE builds SET status = 'cancelled',
299 finished_at = strftime('%Y-%m-%dT%H:%M:%SZ','now') WHERE id = ? AND status = 'pending'`, id)
300 finished_at = strftime('%Y-%m-%dT%H:%M:%SZ','now') WHERE id = ? AND status IN ('pending', 'running')`, id)
300301 if err != nil {
301302 return err
302303 }