scheduler.go
⎇
Raw
1package worker
2
3import (
4 "context"
5 "database/sql"
6 "log/slog"
7 "sync"
8 "time"
9
10 "vidarchive/internal/models"
11 "vidarchive/internal/service"
12)
13
14// Scheduler enqueues downloads for subscriptions whose next run is due. It
15// reuses the normal download pipeline, so refresh-mode and pruning behaviour
16// live in DownloadService.ExecuteDownload.
17type Scheduler struct {
18 subscriptionSvc *service.SubscriptionService
19 downloadSvc *service.DownloadService
20 pool *Pool
21 interval time.Duration
22 ctx context.Context
23 cancel context.CancelFunc
24 wg sync.WaitGroup
25}
26
27func NewScheduler(subscriptionSvc *service.SubscriptionService, downloadSvc *service.DownloadService, pool *Pool, interval time.Duration) *Scheduler {
28 if interval <= 0 {
29 interval = time.Minute
30 }
31 ctx, cancel := context.WithCancel(context.Background())
32 return &Scheduler{
33 subscriptionSvc: subscriptionSvc,
34 downloadSvc: downloadSvc,
35 pool: pool,
36 interval: interval,
37 ctx: ctx,
38 cancel: cancel,
39 }
40}
41
42func (s *Scheduler) Start() {
43 s.backfillNextRuns()
44 s.wg.Add(1)
45 go s.loop()
46}
47
48// Stop waits for the loop because main closes the database next: a checkDue
49// still in flight would query a closed handle or queue a run nobody services.
50func (s *Scheduler) Stop() {
51 s.cancel()
52 s.wg.Wait()
53}
54
55func (s *Scheduler) loop() {
56 defer s.wg.Done()
57
58 ticker := time.NewTicker(s.interval)
59 defer ticker.Stop()
60
61 // A subscription that came due during downtime should not wait a full interval.
62 s.checkDue()
63
64 for {
65 select {
66 case <-ticker.C:
67 s.checkDue()
68 case <-s.ctx.Done():
69 return
70 }
71 }
72}
73
74// backfillNextRuns fills in a missing next run time without running anything, so
75// a restart does not fire every subscription with a null next_run_at.
76func (s *Scheduler) backfillNextRuns() {
77 subs, err := s.subscriptionSvc.GetAll()
78 if err != nil {
79 slog.Error("scheduler backfill failed to list subscriptions", "err", err)
80 return
81 }
82 now := time.Now()
83 for _, sub := range subs {
84 if !sub.Enabled || sub.NextRunAt.Valid {
85 continue
86 }
87 next, err := s.subscriptionSvc.ComputeNextRun(sub, now)
88 if err != nil {
89 slog.Warn("scheduler: invalid schedule", "subscription_id", sub.ID, "err", err)
90 continue
91 }
92 sub.NextRunAt = sql.NullTime{Time: next, Valid: true}
93 if err := s.subscriptionSvc.Update(sub); err != nil {
94 slog.Error("scheduler: failed to set next run", "subscription_id", sub.ID, "err", err)
95 }
96 }
97}
98
99func (s *Scheduler) checkDue() {
100 now := time.Now()
101 due, err := s.subscriptionSvc.GetDue(now)
102 if err != nil {
103 slog.Error("scheduler: failed to query due subscriptions", "err", err)
104 return
105 }
106 for _, sub := range due {
107 s.run(sub, now)
108 }
109}
110
111func (s *Scheduler) run(sub *models.Subscription, now time.Time) {
112 next, err := s.subscriptionSvc.ComputeNextRun(sub, now)
113 if err != nil {
114 // Don't keep retrying a broken schedule every tick; push it out a day.
115 slog.Warn("scheduler: invalid schedule, deferring", "subscription_id", sub.ID, "err", err)
116 next = now.Add(24 * time.Hour)
117 }
118
119 // A run longer than the interval would stack duplicate downloads. next_run_at
120 // still advances, so this is not re-evaluated every tick.
121 if active, aErr := s.downloadSvc.HasActiveForSubscription(sub.ID); aErr != nil {
122 slog.Error("scheduler: active-run check failed", "subscription_id", sub.ID, "err", aErr)
123 } else if active {
124 lastRun := now
125 if sub.LastRunAt.Valid {
126 lastRun = sub.LastRunAt.Time
127 }
128 slog.Info("scheduler: active run in progress, skipping", "subscription_id", sub.ID, "next_run", next.Format(time.RFC3339))
129 if mErr := s.subscriptionSvc.MarkRun(sub.ID, lastRun, next, "skipped"); mErr != nil {
130 slog.Error("scheduler: failed to mark run", "subscription_id", sub.ID, "err", mErr)
131 }
132 return
133 }
134
135 d, err := s.downloadSvc.CreateForSubscription(sub)
136 if err != nil {
137 slog.Error("scheduler: failed to queue subscription", "subscription_id", sub.ID, "err", err)
138 if mErr := s.subscriptionSvc.MarkRun(sub.ID, now, next, "error"); mErr != nil {
139 slog.Error("scheduler: failed to mark run", "subscription_id", sub.ID, "err", mErr)
140 }
141 return
142 }
143
144 s.pool.Submit(d)
145 if err := s.subscriptionSvc.MarkRun(sub.ID, now, next, "queued"); err != nil {
146 slog.Error("scheduler: failed to mark run", "subscription_id", sub.ID, "err", err)
147 }
148 slog.Info("scheduler: queued subscription", "subscription_id", sub.ID, "download_id", d.ID, "next_run", next.Format(time.RFC3339))
149}
150