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 periodically enqueues downloads for subscriptions whose next run is
15// due. It reuses the worker Pool and the normal download pipeline; refresh-mode
16// and pruning behaviour 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 cancels the loop and waits for it to return. Waiting matters: main
49// closes the database next, and a checkDue still in flight would query a closed
50// handle or queue a run nobody will service.
51func (s *Scheduler) Stop() {
52 s.cancel()
53 s.wg.Wait()
54}
55
56func (s *Scheduler) loop() {
57 defer s.wg.Done()
58
59 ticker := time.NewTicker(s.interval)
60 defer ticker.Stop()
61
62 // Check once promptly on startup so a subscription that came due during
63 // downtime doesn't wait a full interval.
64 s.checkDue()
65
66 for {
67 select {
68 case <-ticker.C:
69 s.checkDue()
70 case <-s.ctx.Done():
71 return
72 }
73 }
74}
75
76// backfillNextRuns gives any enabled subscription without a next run time one,
77// without running it — so a fresh restart doesn't fire every subscription that
78// happens to have a null next_run_at.
79func (s *Scheduler) backfillNextRuns() {
80 subs, err := s.subscriptionSvc.GetAll()
81 if err != nil {
82 slog.Error("scheduler backfill failed to list subscriptions", "err", err)
83 return
84 }
85 now := time.Now()
86 for _, sub := range subs {
87 if !sub.Enabled || sub.NextRunAt.Valid {
88 continue
89 }
90 next, err := s.subscriptionSvc.ComputeNextRun(sub, now)
91 if err != nil {
92 slog.Warn("scheduler: invalid schedule", "subscription_id", sub.ID, "err", err)
93 continue
94 }
95 sub.NextRunAt = sql.NullTime{Time: next, Valid: true}
96 if err := s.subscriptionSvc.Update(sub); err != nil {
97 slog.Error("scheduler: failed to set next run", "subscription_id", sub.ID, "err", err)
98 }
99 }
100}
101
102func (s *Scheduler) checkDue() {
103 now := time.Now()
104 due, err := s.subscriptionSvc.GetDue(now)
105 if err != nil {
106 slog.Error("scheduler: failed to query due subscriptions", "err", err)
107 return
108 }
109 for _, sub := range due {
110 s.run(sub, now)
111 }
112}
113
114func (s *Scheduler) run(sub *models.Subscription, now time.Time) {
115 next, err := s.subscriptionSvc.ComputeNextRun(sub, now)
116 if err != nil {
117 // Don't keep retrying a broken schedule every tick; push it out a day.
118 slog.Warn("scheduler: invalid schedule, deferring", "subscription_id", sub.ID, "err", err)
119 next = now.Add(24 * time.Hour)
120 }
121
122 // Skip this run if the previous one is still queued or downloading — a run
123 // longer than the interval would otherwise stack duplicate downloads. Advance
124 // next_run_at so we don't re-evaluate it every tick, preserving the last run.
125 if active, aErr := s.downloadSvc.HasActiveForSubscription(sub.ID); aErr != nil {
126 slog.Error("scheduler: active-run check failed", "subscription_id", sub.ID, "err", aErr)
127 } else if active {
128 lastRun := now
129 if sub.LastRunAt.Valid {
130 lastRun = sub.LastRunAt.Time
131 }
132 slog.Info("scheduler: active run in progress, skipping", "subscription_id", sub.ID, "next_run", next.Format(time.RFC3339))
133 if mErr := s.subscriptionSvc.MarkRun(sub.ID, lastRun, next, "skipped"); mErr != nil {
134 slog.Error("scheduler: failed to mark run", "subscription_id", sub.ID, "err", mErr)
135 }
136 return
137 }
138
139 d, err := s.downloadSvc.CreateForSubscription(sub)
140 if err != nil {
141 slog.Error("scheduler: failed to queue subscription", "subscription_id", sub.ID, "err", err)
142 if mErr := s.subscriptionSvc.MarkRun(sub.ID, now, next, "error"); mErr != nil {
143 slog.Error("scheduler: failed to mark run", "subscription_id", sub.ID, "err", mErr)
144 }
145 return
146 }
147
148 s.pool.Submit(d)
149 if err := s.subscriptionSvc.MarkRun(sub.ID, now, next, "queued"); err != nil {
150 slog.Error("scheduler: failed to mark run", "subscription_id", sub.ID, "err", err)
151 }
152 slog.Info("scheduler: queued subscription", "subscription_id", sub.ID, "download_id", d.ID, "next_run", next.Format(time.RFC3339))
153}
154