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 "vidarchive/internal/tools"
13)
14
15// Scheduler enqueues downloads for subscriptions whose next run is due. It
16// reuses the normal download pipeline, so refresh-mode and pruning behaviour
17// live in DownloadService.ExecuteDownload.
18type Scheduler struct {
19 subscriptionSvc *service.SubscriptionService
20 downloadSvc *service.DownloadService
21 settingsSvc *service.SettingsService
22 tools *tools.Manager
23 pool *Pool
24 interval time.Duration
25 ctx context.Context
26 cancel context.CancelFunc
27 wg sync.WaitGroup
28}
29
30func NewScheduler(subscriptionSvc *service.SubscriptionService, downloadSvc *service.DownloadService, settingsSvc *service.SettingsService, toolMgr *tools.Manager, pool *Pool, interval time.Duration) *Scheduler {
31 if interval <= 0 {
32 interval = time.Minute
33 }
34 ctx, cancel := context.WithCancel(context.Background())
35 return &Scheduler{
36 subscriptionSvc: subscriptionSvc,
37 downloadSvc: downloadSvc,
38 settingsSvc: settingsSvc,
39 tools: toolMgr,
40 pool: pool,
41 interval: interval,
42 ctx: ctx,
43 cancel: cancel,
44 }
45}
46
47func (s *Scheduler) Start() {
48 s.backfillNextRuns()
49 s.wg.Add(1)
50 go s.loop()
51}
52
53// Stop waits for the loop because main closes the database next: a checkDue
54// still in flight would query a closed handle or queue a run nobody services.
55func (s *Scheduler) Stop() {
56 s.cancel()
57 s.wg.Wait()
58}
59
60func (s *Scheduler) loop() {
61 defer s.wg.Done()
62
63 ticker := time.NewTicker(s.interval)
64 defer ticker.Stop()
65
66 // A subscription that came due during downtime should not wait a full interval.
67 s.checkDue()
68
69 for {
70 select {
71 case <-ticker.C:
72 s.checkDue()
73 s.checkToolUpdate(time.Now())
74 case <-s.ctx.Done():
75 return
76 }
77 }
78}
79
80// backfillNextRuns fills in a missing next run time without running anything, so
81// a restart does not fire every subscription with a null next_run_at.
82func (s *Scheduler) backfillNextRuns() {
83 subs, err := s.subscriptionSvc.GetAll()
84 if err != nil {
85 slog.Error("scheduler backfill failed to list subscriptions", "err", err)
86 return
87 }
88 now := time.Now()
89 for _, sub := range subs {
90 if !sub.Enabled || sub.NextRunAt.Valid {
91 continue
92 }
93 next, err := s.subscriptionSvc.ComputeNextRun(sub, now)
94 if err != nil {
95 slog.Warn("scheduler: invalid schedule", "subscription_id", sub.ID, "err", err)
96 continue
97 }
98 sub.NextRunAt = sql.NullTime{Time: next, Valid: true}
99 if err := s.subscriptionSvc.Update(sub); err != nil {
100 slog.Error("scheduler: failed to set next run", "subscription_id", sub.ID, "err", err)
101 }
102 }
103}
104
105func (s *Scheduler) checkDue() {
106 now := time.Now()
107 due, err := s.subscriptionSvc.GetDue(now)
108 if err != nil {
109 slog.Error("scheduler: failed to query due subscriptions", "err", err)
110 return
111 }
112 for _, sub := range due {
113 s.run(sub, now)
114 }
115}
116
117// toolUpdateInterval: yt-dlp releases roughly weekly, and its own updater is a
118// network round trip, so a daily check is plenty.
119const toolUpdateInterval = 24 * time.Hour
120
121func (s *Scheduler) checkToolUpdate(now time.Time) {
122 if s.tools == nil || s.settingsSvc == nil {
123 return
124 }
125 settings, err := s.settingsSvc.GetAll()
126 if err != nil {
127 slog.Error("scheduler: failed to read settings for tool update", "err", err)
128 return
129 }
130 if !settings.ToolAutoUpdate || now.Sub(settings.ToolsCheckedAt) < toolUpdateInterval {
131 return
132 }
133
134 // Stamped before the run, so a failure waits a day instead of retrying every
135 // tick. It also keeps the goroutine below off the database, which matters
136 // because shutdown does not wait for it.
137 if err := s.settingsSvc.SetToolsCheckedAt(now); err != nil {
138 slog.Error("scheduler: failed to record tool check time", "err", err)
139 return
140 }
141
142 // An update takes minutes. Blocking the tick would delay every subscription,
143 // so it runs detached. Shutdown cancels ctx and kills the child. Both
144 // updaters replace the binary by rename, so a killed run leaves the old one
145 // intact.
146 go func() {
147 summary, err := s.tools.Update(s.ctx)
148 if err != nil {
149 slog.Error("scheduled tool update failed", "err", err)
150 return
151 }
152 slog.Info("scheduled tool update finished", "result", summary)
153 }()
154}
155
156func (s *Scheduler) run(sub *models.Subscription, now time.Time) {
157 next, err := s.subscriptionSvc.ComputeNextRun(sub, now)
158 if err != nil {
159 // Don't keep retrying a broken schedule every tick; push it out a day.
160 slog.Warn("scheduler: invalid schedule, deferring", "subscription_id", sub.ID, "err", err)
161 next = now.Add(24 * time.Hour)
162 }
163
164 // A run longer than the interval would stack duplicate downloads. next_run_at
165 // still advances, so this is not re-evaluated every tick.
166 if active, aErr := s.downloadSvc.HasActiveForSubscription(sub.ID); aErr != nil {
167 slog.Error("scheduler: active-run check failed", "subscription_id", sub.ID, "err", aErr)
168 } else if active {
169 lastRun := now
170 if sub.LastRunAt.Valid {
171 lastRun = sub.LastRunAt.Time
172 }
173 slog.Info("scheduler: active run in progress, skipping", "subscription_id", sub.ID, "next_run", next.Format(time.RFC3339))
174 if mErr := s.subscriptionSvc.MarkRun(sub.ID, lastRun, next, "skipped"); mErr != nil {
175 slog.Error("scheduler: failed to mark run", "subscription_id", sub.ID, "err", mErr)
176 }
177 return
178 }
179
180 d, err := s.downloadSvc.CreateForSubscription(sub)
181 if err != nil {
182 slog.Error("scheduler: failed to queue subscription", "subscription_id", sub.ID, "err", err)
183 if mErr := s.subscriptionSvc.MarkRun(sub.ID, now, next, "error"); mErr != nil {
184 slog.Error("scheduler: failed to mark run", "subscription_id", sub.ID, "err", mErr)
185 }
186 return
187 }
188
189 s.pool.Submit(d)
190 if err := s.subscriptionSvc.MarkRun(sub.ID, now, next, "queued"); err != nil {
191 slog.Error("scheduler: failed to mark run", "subscription_id", sub.ID, "err", err)
192 }
193 slog.Info("scheduler: queued subscription", "subscription_id", sub.ID, "download_id", d.ID, "next_run", next.Format(time.RFC3339))
194}
195