package main import ( "context" "errors" "log" "net/url" "strconv" "sync" "time" ) // Schedule says when a saved search runs, in the server time zone (TZ). type Schedule struct { Every string `json:"every"` // "hours", "day" or "week" Hours int `json:"hours,omitempty"` // every: hours (1..168) Weekday int `json:"weekday,omitempty"` // every: week, 0 = Sunday Hour int `json:"hour"` // every: day / week Minute int `json:"minute"` } func (sc Schedule) Valid() bool { switch sc.Every { case "hours": return sc.Hours >= 1 && sc.Hours <= 168 case "day", "week": return sc.Hour >= 0 && sc.Hour < 24 && sc.Minute >= 0 && sc.Minute < 60 && sc.Weekday >= 0 && sc.Weekday < 7 } return false } // Next returns the first run time strictly after "after". last is the time of // the last run, used by the "hours" schedule. func (sc Schedule) Next(after, last time.Time) time.Time { switch sc.Every { case "hours": every := time.Duration(sc.Hours) * time.Hour if last.IsZero() || last.Add(every).Before(after) { return after.Add(time.Minute) } return last.Add(every) case "day", "week": t := time.Date(after.Year(), after.Month(), after.Day(), sc.Hour, sc.Minute, 0, 0, after.Location()) for !t.After(after) || (sc.Every == "week" && int(t.Weekday()) != sc.Weekday) { t = t.AddDate(0, 0, 1) } return t } return time.Time{} } // Scheduler runs the saved searches when they are due. type Scheduler struct { store *Store gh *GitHub runMu sync.Mutex // one run at a time, to spare the GitHub quota } func NewScheduler(store *Store, gh *GitHub) *Scheduler { return &Scheduler{store: store, gh: gh} } // Start checks every 30 seconds for due searches until ctx is done. func (s *Scheduler) Start(ctx context.Context) { now := time.Now() for _, x := range s.store.List() { // plan searches that have no next run yet if x.Enabled && x.NextRunAt.IsZero() { _, _ = s.store.Update(x.ID, func(v *SavedSearch) { v.NextRunAt = v.Schedule.Next(now, v.LastRunAt) }) } } go func() { tick := time.NewTicker(30 * time.Second) defer tick.Stop() for { s.runDue(ctx) select { case <-ctx.Done(): return case <-tick.C: } } }() } func (s *Scheduler) runDue(ctx context.Context) { for _, x := range s.store.List() { if ctx.Err() != nil { return } if x.Enabled && !x.NextRunAt.IsZero() && !time.Now().Before(x.NextRunAt) { if _, err := s.Run(ctx, x.ID, "schedule"); err != nil { log.Printf("saved search %q: %v", x.Name, err) } } } } // Run executes a saved search now and stores the snapshot. func (s *Scheduler) Run(ctx context.Context, id, trigger string) (Run, error) { s.runMu.Lock() defer s.runMu.Unlock() x, err := s.store.Get(id) if err != nil { return Run{}, err } start := time.Now() info := RunInfo{ID: start.UTC().Format("20060102T150405Z"), At: start, Trigger: trigger} var repos []Repo c, err := criteriaFor(x) if err == nil { info.Query = c.Query(start) var res SearchResult if res, err = s.gh.Search(ctx, c); err == nil { info.Total = res.Total repos = res.Items } } info.TookMs = time.Since(start).Milliseconds() if err != nil { info.Error = err.Error() } run, serr := s.store.AddRun(id, info, repos) if serr != nil { return Run{}, serr } _, uerr := s.store.Update(id, func(v *SavedSearch) { v.LastRunAt = start v.LastError = info.Error ri := run.RunInfo v.LastRun = &ri next := v.Schedule.Next(time.Now(), start) if err != nil { var ce *codedError if errors.As(err, &ce) && ce.code == "rate_limited" { // Quota exhausted: try again in a few minutes rather than next week. next = time.Now().Add(3 * time.Minute) } } v.NextRunAt = next }) if uerr != nil { return run, uerr } if err != nil { return run, err } return run, nil } // criteriaFor parses the stored parameters; a saved search always reads the // first page, sorted as saved. func criteriaFor(x SavedSearch) (Criteria, error) { v, err := url.ParseQuery(x.Params) if err != nil { return Criteria{}, &codedError{code: "bad_param", msg: "invalid parameters", detail: err.Error()} } v.Del("page") v.Set("per_page", strconv.Itoa(x.MaxResults)) return ParseCriteria(v) }