internal/ci/sched.go
133 lines · 3745 bytes
1package ci
2
3import (
4 "context"
5 "encoding/json"
6 "fmt"
7 "log/slog"
8 "os"
9 "time"
10
11 "gitbay.org/gitbay/internal/gitutil"
12 "gitbay.org/gitbay/internal/store"
13)
14
15// isoNow is the timestamp format schedules are compared in.
16func isoNow(t time.Time) string { return t.UTC().Format("2006-01-02T15:04:05Z") }
17
18// NextRun returns the next firing time for a cron expression as a stored
19// timestamp. The caller has already validated the expression.
20func NextRun(expr string, after time.Time) string {
21 c, err := ParseCron(expr)
22 if err != nil {
23 return isoNow(after.Add(24 * time.Hour)) // unreachable after Parse validation
24 }
25 n := c.Next(after)
26 if n.IsZero() {
27 return isoNow(after.AddDate(1, 0, 0))
28 }
29 return isoNow(n)
30}
31
32// Scheduler fires scheduled builds. repoDir maps a repo to its bare path.
33type Scheduler struct {
34 St *store.Store
35 RepoDir func(owner, name string) string
36 SiteURL string
37}
38
39// Run ticks until the context ends. GITBAY_SCHED_TICK overrides the
40// interval for tests.
41func (s *Scheduler) Run(ctx context.Context) {
42 tick := time.Minute
43 if v := os.Getenv("GITBAY_SCHED_TICK"); v != "" {
44 if d, err := time.ParseDuration(v); err == nil {
45 tick = d
46 }
47 }
48 t := time.NewTicker(tick)
49 defer t.Stop()
50 for {
51 select {
52 case <-ctx.Done():
53 return
54 case <-t.C:
55 s.RunDue(time.Now())
56 }
57 }
58}
59
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.
63func (s *Scheduler) RunDue(now time.Time) {
64 s.reapStale()
65 due, err := s.St.DueSchedules(isoNow(now))
66 if err != nil {
67 slog.Error("scheduler: listing due builds", "err", err)
68 return
69 }
70 for _, e := range due {
71 repo, err := s.St.RepoByID(e.RepoID)
72 if err != nil {
73 s.St.RemoveSchedule(e.RepoID, e.Job) // repo gone
74 continue
75 }
76 // Always advance first so a broken repo cannot wedge the loop.
77 s.St.SetScheduleNext(e.RepoID, e.Job, NextRun(e.Cron, now))
78 dir := s.RepoDir(repo.OwnerName, repo.Name)
79 sha, err := gitutil.ResolveRef(dir, "refs/heads/"+repo.DefaultBranch)
80 if err != nil {
81 continue // empty repo
82 }
83 raw, err := gitutil.ReadBlob(dir, sha, ConfigPath, 1<<16)
84 if err != nil {
85 s.St.RemoveSchedule(e.RepoID, e.Job) // config removed
86 continue
87 }
88 jobs, err := Parse(raw)
89 if err != nil {
90 continue
91 }
92 var job *Job
93 for i := range jobs {
94 if jobs[i].Name == e.Job && jobs[i].Schedule != "" {
95 job = &jobs[i]
96 break
97 }
98 }
99 if job == nil {
100 s.St.RemoveSchedule(e.RepoID, e.Job)
101 continue
102 }
103 steps, _ := json.Marshal(job.Steps)
104 n, err := s.St.CreateBuild(repo.ID, job.Name, sha, repo.DefaultBranch, string(steps), job.Image, true)
105 if err != nil {
106 slog.Error("scheduler: queueing build", "repo", repo.Path(), "job", job.Name, "err", err)
107 continue
108 }
109 url := fmt.Sprintf("%s/%s/builds/%d", s.SiteURL, repo.Path(), n)
110 s.St.SetCommitStatus(repo.ID, sha, "ci/"+job.Name, "pending", "scheduled", url, 0)
111 }
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.
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}