internal/ci/sched.go
110 lines · 2879 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. Split
61// from the ticker for tests.
62func (s *Scheduler) RunDue(now time.Time) {
63 due, err := s.St.DueSchedules(isoNow(now))
64 if err != nil {
65 slog.Error("scheduler: listing due builds", "err", err)
66 return
67 }
68 for _, e := range due {
69 repo, err := s.St.RepoByID(e.RepoID)
70 if err != nil {
71 s.St.RemoveSchedule(e.RepoID, e.Job) // repo gone
72 continue
73 }
74 // Always advance first so a broken repo cannot wedge the loop.
75 s.St.SetScheduleNext(e.RepoID, e.Job, NextRun(e.Cron, now))
76 dir := s.RepoDir(repo.OwnerName, repo.Name)
77 sha, err := gitutil.ResolveRef(dir, "refs/heads/"+repo.DefaultBranch)
78 if err != nil {
79 continue // empty repo
80 }
81 raw, err := gitutil.ReadBlob(dir, sha, ConfigPath, 1<<16)
82 if err != nil {
83 s.St.RemoveSchedule(e.RepoID, e.Job) // config removed
84 continue
85 }
86 jobs, err := Parse(raw)
87 if err != nil {
88 continue
89 }
90 var job *Job
91 for i := range jobs {
92 if jobs[i].Name == e.Job && jobs[i].Schedule != "" {
93 job = &jobs[i]
94 break
95 }
96 }
97 if job == nil {
98 s.St.RemoveSchedule(e.RepoID, e.Job)
99 continue
100 }
101 steps, _ := json.Marshal(job.Steps)
102 n, err := s.St.CreateBuild(repo.ID, job.Name, sha, repo.DefaultBranch, string(steps))
103 if err != nil {
104 slog.Error("scheduler: queueing build", "repo", repo.Path(), "job", job.Name, "err", err)
105 continue
106 }
107 url := fmt.Sprintf("%s/%s/builds/%d", s.SiteURL, repo.Path(), n)
108 s.St.SetCommitStatus(repo.ID, sha, "ci/"+job.Name, "pending", "scheduled", url, 0)
109 }
110}