package tasks import ( "context" "fmt" "strconv" "strings" "time" ) // Due is a schedule row that has come due: its id and its schedule string. The // id identifies the row for claiming and tells the runner which task to // enqueue — for system schedules it is the schedule type. type Due struct { ID string Schedule string } // ScheduleStore reads and claims due schedules in one table. It traffics only // in the coordination columns every schedule table shares — an id, a schedule // string, and next/last run times — so one Runner drives any such table // regardless of the domain columns alongside them. type ScheduleStore interface { // Due returns every schedule whose next run time has passed. Due(ctx context.Context) ([]Due, error) // Claim advances a schedule's run times, setting its next run to next. It // reports whether this caller won the claim; a lost claim (no rows) means // another tick or instance already fired this occurrence. Claim(ctx context.Context, id string, next time.Time) (bool, error) } // Schedule is a recurring schedule. Next reports the first fire time strictly // after `after`; all times are UTC. String round-trips through ParseSchedule. type Schedule interface { Next(after time.Time) time.Time String() string } // Every fires on a fixed period anchored to the Unix epoch (UTC). Because the // fire times form one arithmetic sequence from a fixed anchor — never reset per // hour or day — any whole-minute period is exact with no wrap-gap (unlike // cron's */N, which resets each hour). type Every struct { Period time.Duration } func (e Every) Next(after time.Time) time.Time { epoch := time.Unix(0, 0).UTC() elapsed := after.Sub(epoch) periods := elapsed / e.Period return epoch.Add((periods + 1) * e.Period) } func (e Every) String() string { m := e.Period / time.Minute if m%60 == 0 { return fmt.Sprintf("every %dh", m/60) } return fmt.Sprintf("every %dm", m) } // DailyAt fires once a day at a wall-clock time in UTC. type DailyAt struct { Hour int Minute int } func (d DailyAt) Next(after time.Time) time.Time { after = after.UTC() t := time.Date(after.Year(), after.Month(), after.Day(), d.Hour, d.Minute, 0, 0, time.UTC) if !t.After(after) { t = t.AddDate(0, 0, 1) } return t } func (d DailyAt) String() string { return fmt.Sprintf("daily %02d:%02d", d.Hour, d.Minute) } // WeeklyAt fires once a week on a weekday at a wall-clock time in UTC. type WeeklyAt struct { Weekday time.Weekday Hour int Minute int } func (w WeeklyAt) Next(after time.Time) time.Time { after = after.UTC() t := time.Date(after.Year(), after.Month(), after.Day(), w.Hour, w.Minute, 0, 0, time.UTC) daysAhead := (int(w.Weekday) - int(t.Weekday()) + 7) % 7 t = t.AddDate(0, 0, daysAhead) if !t.After(after) { t = t.AddDate(0, 0, 7) } return t } func (w WeeklyAt) String() string { return fmt.Sprintf("weekly %s %02d:%02d", w.Weekday.String()[:3], w.Hour, w.Minute) } // ParseSchedule parses one of: // // every e.g. "every 5m", "every 2h" (whole minutes, 1m..366d) // daily HH:MM e.g. "daily 03:00" (UTC) // weekly WD HH:MM e.g. "weekly Sun 03:00" (UTC, WD = Sun..Sat) func ParseSchedule(s string) (Schedule, error) { // maxPeriod bounds Every to catch typos and keep epoch arithmetic sane. const maxPeriod = 366 * 24 * time.Hour fields := strings.Fields(s) if len(fields) == 0 { return nil, fmt.Errorf("empty schedule") } switch fields[0] { case "every": if len(fields) != 2 { return nil, fmt.Errorf("schedule %q: 'every' needs a duration, e.g. 'every 5m'", s) } d, err := time.ParseDuration(fields[1]) if err != nil { return nil, fmt.Errorf("schedule %q: %w", s, err) } if d < time.Minute { return nil, fmt.Errorf("schedule %q: period must be at least 1m", s) } if d%time.Minute != 0 { return nil, fmt.Errorf("schedule %q: period must be a whole number of minutes", s) } if d > maxPeriod { return nil, fmt.Errorf("schedule %q: period must be at most %s", s, maxPeriod) } return Every{Period: d}, nil case "daily": if len(fields) != 2 { return nil, fmt.Errorf("schedule %q: 'daily' needs a time, e.g. 'daily 03:00'", s) } h, m, err := parseHHMM(fields[1]) if err != nil { return nil, fmt.Errorf("schedule %q: %w", s, err) } return DailyAt{Hour: h, Minute: m}, nil case "weekly": if len(fields) != 3 { return nil, fmt.Errorf("schedule %q: 'weekly' needs a weekday and time, e.g. 'weekly Sun 03:00'", s) } wd, err := parseWeekday(fields[1]) if err != nil { return nil, fmt.Errorf("schedule %q: %w", s, err) } h, m, err := parseHHMM(fields[2]) if err != nil { return nil, fmt.Errorf("schedule %q: %w", s, err) } return WeeklyAt{Weekday: wd, Hour: h, Minute: m}, nil default: return nil, fmt.Errorf("schedule %q: unknown kind %q (want every/daily/weekly)", s, fields[0]) } } func parseHHMM(s string) (int, int, error) { hh, mm, ok := strings.Cut(s, ":") if !ok { return 0, 0, fmt.Errorf("time %q must be HH:MM", s) } h, err := strconv.Atoi(hh) if err != nil || h < 0 || h > 23 { return 0, 0, fmt.Errorf("time %q: hour must be 00-23", s) } m, err := strconv.Atoi(mm) if err != nil || m < 0 || m > 59 { return 0, 0, fmt.Errorf("time %q: minute must be 00-59", s) } return h, m, nil } var weekdays = map[string]time.Weekday{ "sun": time.Sunday, "mon": time.Monday, "tue": time.Tuesday, "wed": time.Wednesday, "thu": time.Thursday, "fri": time.Friday, "sat": time.Saturday, } func parseWeekday(s string) (time.Weekday, error) { wd, ok := weekdays[strings.ToLower(s)] if !ok { return 0, fmt.Errorf("weekday %q must be Sun/Mon/Tue/Wed/Thu/Fri/Sat", s) } return wd, nil }