internal/ci/sched.go

v1.32.0
gitbay/internal/ci/sched.go history · blame · raw

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}