internal/ci/sched.go
143 lines · 4230 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 // The push path skips a job whose build for the commit is still
104 // pending or running; so does a tick. Without this a repository
105 // no runner serves gained one row per tick forever (#206). A
106 // finished build does not suppress the tick: a schedule re-runs
107 // an unchanged commit on purpose.
108 if built, err := s.St.BuildsForCommit(repo.ID, sha); err == nil {
109 if b, ok := built[job.Name]; ok && (b.Status == "pending" || b.Status == "running") {
110 continue
111 }
112 }
113 steps, _ := json.Marshal(job.Steps)
114 n, err := s.St.CreateBuild(repo.ID, job.Name, sha, repo.DefaultBranch, string(steps), job.Image, "", true)
115 if err != nil {
116 slog.Error("scheduler: queueing build", "repo", repo.Path(), "job", job.Name, "err", err)
117 continue
118 }
119 url := fmt.Sprintf("%s/%s/builds/%d", s.SiteURL, repo.Path(), n)
120 s.St.SetCommitStatus(repo.ID, sha, "ci/"+job.Name, "pending", "scheduled", url, 0)
121 }
122}
123
124// reapStale resolves builds running past the deadline, so a runner that
125// died mid-build leaves neither a build running nor a commit pending
126// forever. It runs on the tick rather than on the next runner claim, so it
127// does not need a runner to be alive.
128func (s *Scheduler) reapStale() {
129 stale, err := s.St.ReapStaleBuilds()
130 if err != nil {
131 slog.Error("scheduler: reaping stale builds", "err", err)
132 return
133 }
134 for _, b := range stale {
135 repo, err := s.St.RepoByID(b.RepoID)
136 if err != nil {
137 continue
138 }
139 url := fmt.Sprintf("%s/%s/builds/%d", s.SiteURL, repo.Path(), b.Number)
140 s.St.SetCommitStatus(repo.ID, b.SHA, "ci/"+b.Job, "failure", "build abandoned", url, 0)
141 slog.Warn("build abandoned", "repo", repo.Path(), "build", b.Number, "job", b.Job)
142 }
143}