Commit 99de33111e
Verified · cmc ci/build: success ci/test: failure ci/vuln: success
Layout: unified · split
cmd/gitbay-runner/main.go +24 −3
| @@ -180,12 +180,29 @@ func (r *runner) run(j job) bool { | |||
| 180 | log.Printf("build %d: log stream: %v", j.ID, err) | 180 | log.Printf("build %d: log stream: %v", j.ID, err) |
| 181 | return false | 181 | return false |
| 182 | } | 182 | } |
| 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 | }() | ||
| 183 | defer func() { | 195 | defer func() { |
| 184 | if sink.broken() { | 196 | select { |
| 185 | log.Printf("build %d: log stream lost; stored log is incomplete", j.ID) | 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 | } | ||
| 186 | } | 203 | } |
| 187 | pipe.Close() | 204 | pipe.Close() |
| 188 | logCmd.Wait() | 205 | <-logExited |
| 189 | }() | 206 | }() |
| 190 | 207 | ||
| 191 | gitSSH := strings.TrimSpace("ssh " + strings.Join(r.sshOpts, " ")) | 208 | gitSSH := strings.TrimSpace("ssh " + strings.Join(r.sshOpts, " ")) |
| @@ -227,6 +244,10 @@ func (r *runner) run(j job) bool { | |||
| 227 | fmt.Fprintf(sink, "step failed: %v\n", err) | 244 | fmt.Fprintf(sink, "step failed: %v\n", err) |
| 228 | return false | 245 | return false |
| 229 | } | 246 | } |
| 247 | case <-cancelled: | ||
| 248 | cmd.Process.Kill() | ||
| 249 | <-done | ||
| 250 | return false | ||
| 230 | case <-time.After(time.Until(deadline)): | 251 | case <-time.After(time.Until(deadline)): |
| 231 | cmd.Process.Kill() | 252 | cmd.Process.Kill() |
| 232 | fmt.Fprintf(sink, "build timed out after %s\n", r.timeout) | 253 | fmt.Fprintf(sink, "build timed out after %s\n", r.timeout) |
e2e/build_cancel_test.go +80 −7
| @@ -1,10 +1,13 @@ | |||
| 1 | package e2e | 1 | package e2e |
| 2 | 2 | ||
| 3 | import ( | 3 | import ( |
| 4 | "fmt" | ||
| 4 | "os" | 5 | "os" |
| 6 | "os/exec" | ||
| 5 | "path/filepath" | 7 | "path/filepath" |
| 6 | "strings" | 8 | "strings" |
| 7 | "testing" | 9 | "testing" |
| 10 | "time" | ||
| 8 | ) | 11 | ) |
| 9 | 12 | ||
| 10 | // build cancel withdraws a queued build; a running one is the runner's. | 13 | // build cancel withdraws a queued build; a running one is the runner's. |
| @@ -82,14 +85,84 @@ func TestBuildCancel(t *testing.T) { | |||
| 82 | if st, _, _ := inst.ssh(t, aliceKey, "", "status", "list", "alice/app", sha); !strings.Contains(st, "success") || !strings.Contains(st, "passed in build 2") { | 85 | if st, _, _ := inst.ssh(t, aliceKey, "", "status", "list", "alice/app", sha); !strings.Contains(st, "success") || !strings.Contains(st, "passed in build 2") { |
| 83 | t.Fatalf("status after cancelling the duplicate:\n%s", st) | 86 | t.Fatalf("status after cancelling the duplicate:\n%s", st) |
| 84 | } | 87 | } |
| 85 | // A running build cannot be cancelled here. | 88 | } |
| 86 | if _, _, code := inst.ssh(t, aliceKey, "", "build", "trigger", "alice/app", "unit"); code != 0 { | 89 | |
| 87 | t.Fatal("third trigger failed") | 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. | ||
| 93 | func 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") | ||
| 88 | } | 102 | } |
| 89 | if _, _, code := inst.ssh(t, runnerKey, "", "runner", "next"); code != 0 { | 103 | work := t.TempDir() |
| 90 | t.Fatal("claim failed") | 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 | ||
| 91 | } | 117 | } |
| 92 | if _, errOut, code := inst.ssh(t, aliceKey, "", "build", "cancel", "alice/app", "4"); code != 2 || !strings.Contains(errOut, "running") { | 118 | |
| 93 | t.Fatalf("cancelled a running build: exit %d %s", code, errOut) | 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) | ||
| 94 | } | 167 | } |
| 95 | } | 168 | } |
internal/control/build.go +62 −17
| @@ -379,24 +379,53 @@ func runRunnerLog(c *Ctx, args []string) int { | |||
| 379 | // append that fails drops its chunk and the loop keeps draining: ending | 379 | // append that fails drops its chunk and the loop keeps draining: ending |
| 380 | // the session here breaks the runner's pipe, and a broken pipe is how a | 380 | // the session here breaks the runner's pipe, and a broken pipe is how a |
| 381 | // transient SQLITE_BUSY used to fail the build the log belonged to. | 381 | // 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() | ||
| 383 | dropped := 0 | 407 | dropped := 0 |
| 384 | for { | 408 | for { |
| 385 | n, rerr := c.Stdin.Read(buf) | 409 | select { |
| 386 | if n > 0 { | 410 | case ch := <-chunks: |
| 387 | if err := c.Store.AppendBuildLog(id, buf[:n]); err != nil { | 411 | if len(ch.data) > 0 { |
| 388 | dropped++ | 412 | if err := c.Store.AppendBuildLog(id, ch.data); err != nil { |
| 389 | slog.Warn("appending build log", "build", id, "err", err) | 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) | ||
| 390 | } | 426 | } |
| 391 | } | 427 | } |
| 392 | if rerr != nil { | ||
| 393 | break | ||
| 394 | } | ||
| 395 | } | ||
| 396 | if dropped > 0 { | ||
| 397 | slog.Warn("build log incomplete", "build", id, "dropped_chunks", dropped) | ||
| 398 | } | 428 | } |
| 399 | return c.emit(map[string]string{"log": "ok"}, func(w io.Writer) {}) | ||
| 400 | } | 429 | } |
| 401 | 430 | ||
| 402 | func runRunnerDone(c *Ctx, args []string) int { | 431 | func runRunnerDone(c *Ctx, args []string) int { |
| @@ -414,6 +443,14 @@ func runRunnerDone(c *Ctx, args []string) int { | |||
| 414 | if err != nil { | 443 | if err != nil { |
| 415 | return c.fail(protocol.ExitNotFound, "no build %d", id) | 444 | return c.fail(protocol.ExitNotFound, "no build %d", id) |
| 416 | } | 445 | } |
| 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 | } | ||
| 417 | if err := c.Store.FinishBuild(id, args[1]); err != nil { | 454 | if err := c.Store.FinishBuild(id, args[1]); err != nil { |
| 418 | return c.fail(protocol.ExitFailure, "finishing build %d: %v", id, err) | 455 | return c.fail(protocol.ExitFailure, "finishing build %d: %v", id, err) |
| 419 | } | 456 | } |
| @@ -527,13 +564,17 @@ func runBuildCancel(c *Ctx, args []string) int { | |||
| 527 | if !policy.CanWrite(c.User, repo, grant) { | 564 | if !policy.CanWrite(c.User, repo, grant) { |
| 528 | return c.fail(protocol.ExitDenied, "cancelling a build needs write access to %s", repo.Path()) | 565 | return c.fail(protocol.ExitDenied, "cancelling a build needs write access to %s", repo.Path()) |
| 529 | } | 566 | } |
| 530 | if b.Status != "pending" { | 567 | if b.Status != "pending" && b.Status != "running" { |
| 531 | return c.fail(protocol.ExitUsage, "build %d is %s; only a queued build can be cancelled", b.Number, b.Status) | 568 | return c.fail(protocol.ExitUsage, "build %d is %s; only a queued or running build can be cancelled", b.Number, b.Status) |
| 532 | } | 569 | } |
| 533 | if err := c.Store.CancelBuild(b.ID); err != nil { | 570 | if err := c.Store.CancelBuild(b.ID); err != nil { |
| 534 | return c.fail(protocol.ExitFailure, "%v", err) | 571 | return c.fail(protocol.ExitFailure, "%v", err) |
| 535 | } | 572 | } |
| 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 | } | ||
| 537 | // The queued status replaced whatever the commit had for this job. If | 578 | // The queued status replaced whatever the commit had for this job. If |
| 538 | // the commit passed the job on another ref, that result stands again; | 579 | // the commit passed the job on another ref, that result stands again; |
| 539 | // otherwise the context says it was withdrawn. | 580 | // otherwise the context says it was withdrawn. |
| @@ -546,7 +587,11 @@ func runBuildCancel(c *Ctx, args []string) int { | |||
| 546 | c.Store.SetCommitStatus(repo.ID, b.SHA, "ci/"+b.Job, "error", "cancelled", url, c.User.ID) | 587 | c.Store.SetCommitStatus(repo.ID, b.SHA, "ci/"+b.Job, "error", "cancelled", url, c.User.ID) |
| 547 | } | 588 | } |
| 548 | c.Store.RecordEvent(repo.ID, c.User.ID, "build.cancelled", fmt.Sprintf(`{"number":%d,"job":%q}`, b.Number, b.Job)) | 589 | 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 | } | ||
| 550 | fmt.Fprintf(w, "cancelled %s build %d (%s)\n", repo.Path(), b.Number, b.Job) | 595 | fmt.Fprintf(w, "cancelled %s build %d (%s)\n", repo.Path(), b.Number, b.Job) |
| 551 | }) | 596 | }) |
| 552 | } | 597 | } |
internal/store/builds.go +5 −4
| @@ -291,12 +291,13 @@ func (b Build) Elapsed() time.Duration { | |||
| 291 | return 0 | 291 | return 0 |
| 292 | } | 292 | } |
| 293 | 293 | ||
| 294 | // CancelBuild withdraws a build that no runner has claimed. A running | 294 | // CancelBuild withdraws a queued or running build. A running one is |
| 295 | // build is the runner's to finish; cancelling it here would leave the | 295 | // ended by the runner, which learns of the cancellation when its log |
| 296 | // runner reporting on a row that says otherwise. | 296 | // session is closed, and whose later report lands on a row that already |
| 297 | // says cancelled. | ||
| 297 | func (s *Store) CancelBuild(id int64) error { | 298 | func (s *Store) CancelBuild(id int64) error { |
| 298 | res, err := s.DB.Exec(`UPDATE builds SET status = 'cancelled', | 299 | 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) |
| 300 | if err != nil { | 301 | if err != nil { |
| 301 | return err | 302 | return err |
| 302 | } | 303 | } |