internal/ci/sched.go

8246c58f81115a972853af8d0ed7cbaa35bd1339
gitbay/internal/ci/sched.go history · blame · raw

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}