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