ci: reap stale builds on the scheduler tick !159

merged merged by cmc on 2026-09-02 03:29 UTC · krz/gitbay:stack-1-reaper into main

4 files changed, +118 −18

Layout: unified · split

e2e/reap_test.go added +80
@@ -0,0 +1,80 @@
1package e2e
2
3import (
4 "encoding/json"
5 "os"
6 "path/filepath"
7 "strings"
8 "testing"
9 "time"
10)
11
12// A build claimed by a runner that never reports is failed by the
13// scheduler's tick, with no runner alive to trigger it.
14func TestStaleBuildReapedWithoutRunner(t *testing.T) {
15 t.Setenv("GITBAY_SCHED_TICK", "500ms")
16 t.Setenv("GITBAY_STALE_BUILD_DEADLINE", "2s")
17 inst := startInstance(t)
18 aliceKey := inst.newKey(t, "alice")
19 runnerKey := inst.newKey(t, "ci")
20 inst.admin(t, "admin", "user", "create", "alice", "--key", aliceKey+".pub")
21 inst.admin(t, "admin", "user", "create", "ci", "--key", runnerKey+".pub", "--admin")
22 if _, _, code := inst.ssh(t, aliceKey, "", "repo", "create", "alice/app"); code != 0 {
23 t.Fatal("repo create failed")
24 }
25 work := t.TempDir()
26 env := inst.gitEnv(aliceKey)
27 mustGit(t, work, env, "clone", inst.sshURL("alice/app"), "w")
28 dir := filepath.Join(work, "w")
29 os.MkdirAll(filepath.Join(dir, ".gitbay"), 0o755)
30 os.WriteFile(filepath.Join(dir, ".gitbay", "ci.yml"), []byte("jobs:\n ok:\n steps:\n - echo fine\n"), 0o644)
31 mustGit(t, dir, env, "checkout", "-q", "-b", "main")
32 mustGit(t, dir, env, "add", ".")
33 mustGit(t, dir, env, "commit", "-q", "-m", "ci")
34 mustGit(t, dir, env, "push", "-q", "origin", "main")
35
36 // Claim the build the way a runner does, then vanish.
37 out, errOut, code := inst.ssh(t, runnerKey, "", "runner", "next", "--json")
38 if code != 0 || !strings.Contains(out, `"job":"ok"`) {
39 t.Fatalf("runner next: exit %d\n%s%s", code, out, errOut)
40 }
41 var claim struct {
42 Data struct {
43 SHA string `json:"sha"`
44 } `json:"data"`
45 }
46 json.Unmarshal([]byte(out), &claim)
47
48 status := func() (string, string) {
49 out, _, _ := inst.ssh(t, aliceKey, "", "build", "list", "alice/app", "--json")
50 var env struct {
51 Data []struct {
52 Status string `json:"status"`
53 } `json:"data"`
54 }
55 json.Unmarshal([]byte(out), &env)
56 st := ""
57 if len(env.Data) > 0 {
58 st = env.Data[0].Status
59 }
60 cs, _, _ := inst.ssh(t, aliceKey, "", "status", "list", "alice/app", claim.Data.SHA)
61 return st, cs
62 }
63 if st, _ := status(); st != "running" {
64 t.Fatalf("claimed build is %q, want running", st)
65 }
66 deadline := time.Now().Add(20 * time.Second)
67 for {
68 st, cs := status()
69 if st == "failure" && strings.Contains(cs, "build abandoned") {
70 break
71 }
72 if time.Now().After(deadline) {
73 t.Fatalf("never reaped: build %q, statuses:\n%s", st, cs)
74 }
75 time.Sleep(250 * time.Millisecond)
76 }
77 if out, _, _ := inst.ssh(t, aliceKey, "", "build", "log", "alice/app", "1"); !strings.Contains(out, "build abandoned") {
78 t.Fatalf("log lacks the abandonment note:\n%s", out)
79 }
80}
internal/ci/sched.go +25 −2
@@ -57,9 +57,11 @@ func (s *Scheduler) Run(ctx context.Context) {
5757 }
5858}
5959
60// RunDue queues every due scheduled build and advances its next_run. Split
61// from the ticker for tests.
60// RunDue queues every due scheduled build and advances its next_run, and
61// fails any build a runner claimed and never reported. Split from the
62// ticker for tests.
6263func (s *Scheduler) RunDue(now time.Time) {
64 s.reapStale()
6365 due, err := s.St.DueSchedules(isoNow(now))
6466 if err != nil {
6567 slog.Error("scheduler: listing due builds", "err", err)
@@ -108,3 +110,24 @@ func (s *Scheduler) RunDue(now time.Time) {
108110 s.St.SetCommitStatus(repo.ID, sha, "ci/"+job.Name, "pending", "scheduled", url, 0)
109111 }
110112}
113
114// reapStale resolves builds running past the deadline, so a runner that
115// died mid-build leaves neither a build running nor a commit pending
116// forever. It runs on the tick rather than on the next runner claim, so it
117// does not need a runner to be alive.
118func (s *Scheduler) reapStale() {
119 stale, err := s.St.ReapStaleBuilds()
120 if err != nil {
121 slog.Error("scheduler: reaping stale builds", "err", err)
122 return
123 }
124 for _, b := range stale {
125 repo, err := s.St.RepoByID(b.RepoID)
126 if err != nil {
127 continue
128 }
129 url := fmt.Sprintf("%s/%s/builds/%d", s.SiteURL, repo.Path(), b.Number)
130 s.St.SetCommitStatus(repo.ID, b.SHA, "ci/"+b.Job, "failure", "build abandoned", url, 0)
131 slog.Warn("build abandoned", "repo", repo.Path(), "build", b.Number, "job", b.Job)
132 }
133}
internal/control/build.go −15
@@ -314,21 +314,6 @@ func runRunnerNext(c *Ctx, args []string) int {
314314 if code := requireRunner(c); code >= 0 {
315315 return code
316316 }
317 // Resolve anything a previous runner claimed and never reported, so a
318 // killed runner does not leave a build running and a commit pending forever.
319 if stale, err := c.Store.ReapStaleBuilds(); err != nil {
320 return c.fail(protocol.ExitFailure, "%v", err)
321 } else {
322 for _, sb := range stale {
323 repo, err := c.Store.RepoByID(sb.RepoID)
324 if err != nil {
325 continue
326 }
327 url := fmt.Sprintf("%s/%s/builds/%d", c.Cfg.Server.SiteURL, repo.Path(), sb.Number)
328 c.Store.SetCommitStatus(repo.ID, sb.SHA, "ci/"+sb.Job, "failure",
329 "build abandoned", url, c.User.ID)
330 }
331 }
332317 // A runner may limit itself to named repositories. The operator chooses
333318 // what a given runner executes by how they start it; this is scoping the
334319 // runner asks for, not an ACL the server holds over it.
internal/store/builds.go +13 −1
@@ -3,6 +3,7 @@ package store
33import (
44 "database/sql"
55 "errors"
6 "os"
67 "strconv"
78 "strings"
89 "time"
@@ -112,12 +113,23 @@ func (s *Store) ClaimBuild(repoIDs []int64) (Build, bool, error) {
112113// it was killed, restarted, or lost the network mid-build.
113114const StaleBuildDeadline = 90 * time.Minute
114115
116// staleBuildDeadline is StaleBuildDeadline unless GITBAY_STALE_BUILD_DEADLINE
117// shortens it, which tests do.
118func staleBuildDeadline() time.Duration {
119 if v := os.Getenv("GITBAY_STALE_BUILD_DEADLINE"); v != "" {
120 if d, err := time.ParseDuration(v); err == nil && d > 0 {
121 return d
122 }
123 }
124 return StaleBuildDeadline
125}
126
115127// ReapStaleBuilds fails every build that has been running past the deadline and
116128// returns them, so the caller can resolve their commit statuses. A runner that
117129// dies between claiming a build and reporting it otherwise leaves the row
118130// claimed forever, and the commit pending forever with it.
119131func (s *Store) ReapStaleBuilds() ([]Build, error) {
120 cutoff := time.Now().UTC().Add(-StaleBuildDeadline).Format("2006-01-02T15:04:05Z")
132 cutoff := time.Now().UTC().Add(-staleBuildDeadline()).Format("2006-01-02T15:04:05Z")
121133 rows, err := s.DB.Query(buildSelect+
122134 " WHERE status = 'running' AND started_at != '' AND started_at < ?", cutoff)
123135 if err != nil {