A search can now be saved with a schedule (every N hours, daily, or weekly at a given day and time). An internal scheduler runs it and keeps a snapshot of the results on a /data volume (JSON files, no new dependency). The new Watch tab lists the scheduled searches; the review page of each one shows, for any run of its history, the repositories never seen before and the ones that gained the most stars since the previous run. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
164 lines
4.2 KiB
Go
164 lines
4.2 KiB
Go
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)
|
|
}
|