ci: reap stale builds on the scheduler tick !159
4 files changed, +118 −18
Layout: unified · split
e2e/reap_test.go added +80
| @@ -0,0 +1,80 @@ | ||
| 1 | package e2e | |
| 2 | ||
| 3 | import ( | |
| 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. | |
| 14 | func 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) { | ||
| 57 | 57 | } |
| 58 | 58 | } |
| 59 | 59 | |
| 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. | |
| 62 | 63 | func (s *Scheduler) RunDue(now time.Time) { |
| 64 | s.reapStale() | |
| 63 | 65 | due, err := s.St.DueSchedules(isoNow(now)) |
| 64 | 66 | if err != nil { |
| 65 | 67 | slog.Error("scheduler: listing due builds", "err", err) |
| @@ -108,3 +110,24 @@ func (s *Scheduler) RunDue(now time.Time) { | ||
| 108 | 110 | s.St.SetCommitStatus(repo.ID, sha, "ci/"+job.Name, "pending", "scheduled", url, 0) |
| 109 | 111 | } |
| 110 | 112 | } |
| 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. | |
| 118 | func (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 { | ||
| 314 | 314 | if code := requireRunner(c); code >= 0 { |
| 315 | 315 | return code |
| 316 | 316 | } |
| 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 | } | |
| 332 | 317 | // A runner may limit itself to named repositories. The operator chooses |
| 333 | 318 | // what a given runner executes by how they start it; this is scoping the |
| 334 | 319 | // runner asks for, not an ACL the server holds over it. |
internal/store/builds.go +13 −1
| @@ -3,6 +3,7 @@ package store | ||
| 3 | 3 | import ( |
| 4 | 4 | "database/sql" |
| 5 | 5 | "errors" |
| 6 | "os" | |
| 6 | 7 | "strconv" |
| 7 | 8 | "strings" |
| 8 | 9 | "time" |
| @@ -112,12 +113,23 @@ func (s *Store) ClaimBuild(repoIDs []int64) (Build, bool, error) { | ||
| 112 | 113 | // it was killed, restarted, or lost the network mid-build. |
| 113 | 114 | const StaleBuildDeadline = 90 * time.Minute |
| 114 | 115 | |
| 116 | // staleBuildDeadline is StaleBuildDeadline unless GITBAY_STALE_BUILD_DEADLINE | |
| 117 | // shortens it, which tests do. | |
| 118 | func 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 | ||
| 115 | 127 | // ReapStaleBuilds fails every build that has been running past the deadline and |
| 116 | 128 | // returns them, so the caller can resolve their commit statuses. A runner that |
| 117 | 129 | // dies between claiming a build and reporting it otherwise leaves the row |
| 118 | 130 | // claimed forever, and the commit pending forever with it. |
| 119 | 131 | func (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") | |
| 121 | 133 | rows, err := s.DB.Query(buildSelect+ |
| 122 | 134 | " WHERE status = 'running' AND started_at != '' AND started_at < ?", cutoff) |
| 123 | 135 | if err != nil { |