internal/ci/sched.go

3797906826d106b7e052e9d1a95958af851ec052
gitbay/internal/ci/sched.go history · blame · raw

133 lines · 3749 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}