scheduler_test.go
⎇
Raw
1package worker
2
3import (
4 "os"
5 "path/filepath"
6 "testing"
7 "time"
8
9 "vidarchive/internal/config"
10 "vidarchive/internal/database"
11 "vidarchive/internal/models"
12 "vidarchive/internal/repository"
13 "vidarchive/internal/service"
14 "vidarchive/internal/tools"
15)
16
17type schedEnv struct {
18 sched *Scheduler
19 pool *Pool
20 subSvc *service.SubscriptionService
21 settingsSvc *service.SettingsService
22 subRepo *repository.SubscriptionRepository
23 repo *repository.DownloadRepository
24}
25
26// newTestScheduler builds a scheduler over a real database. The pool is created
27// but never started, so a queued run stays in the buffer where a test can count
28// it instead of being executed.
29func newTestScheduler(t *testing.T) *schedEnv {
30 t.Helper()
31
32 root := t.TempDir()
33 cfg := &config.Config{
34 DBPath: ":memory:",
35 DataDir: root,
36 LibraryDir: filepath.Join(root, "library"),
37 TempDir: filepath.Join(root, "temp"),
38 YTDLPPath: "/bin/false",
39 FFmpegPath: "ffmpeg",
40 FFprobePath: "ffprobe",
41 }
42 for _, dir := range []string{cfg.LibraryDir, cfg.TempDir} {
43 if err := os.MkdirAll(dir, 0o755); err != nil {
44 t.Fatal(err)
45 }
46 }
47
48 db, err := database.New(cfg)
49 if err != nil {
50 t.Fatalf("init db: %v", err)
51 }
52 t.Cleanup(func() { db.Close() })
53
54 downloadRepo := repository.NewDownloadRepository(db)
55 subRepo := repository.NewSubscriptionRepository(db)
56 subSvc := service.NewSubscriptionService(subRepo, cfg)
57 settingsSvc := service.NewSettingsService(repository.NewSettingsRepository(db))
58 downloadSvc := service.NewDownloadService(
59 downloadRepo,
60 service.NewLibraryService(cfg.LibraryDir, cfg.FFmpegPath, cfg.FFprobePath),
61 service.NewPresetService(repository.NewPresetRepository(db)),
62 settingsSvc,
63 subSvc,
64 cfg,
65 )
66
67 pool := New(downloadSvc, 1)
68 return &schedEnv{
69 sched: NewScheduler(subSvc, downloadSvc, settingsSvc, tools.New(cfg), pool, time.Minute),
70 pool: pool,
71 subSvc: subSvc,
72 settingsSvc: settingsSvc,
73 subRepo: subRepo,
74 repo: downloadRepo,
75 }
76}
77
78func (e *schedEnv) newSub(t *testing.T, name string, sub *models.Subscription) *models.Subscription {
79 t.Helper()
80 sub.Name = name
81 if sub.URL == "" {
82 sub.URL = "https://example.com/" + name
83 }
84 if sub.RefreshMode == "" {
85 sub.RefreshMode = "overwrite"
86 }
87 if sub.ScheduleKind == "" {
88 sub.ScheduleKind = "daily"
89 }
90 if sub.CronExpr == "" {
91 sub.CronExpr = "0 3 * * *"
92 }
93 if sub.OutputDir == "" {
94 sub.OutputDir = name
95 }
96 if err := e.subRepo.Create(sub); err != nil {
97 t.Fatalf("create subscription %s: %v", name, err)
98 }
99 return sub
100}
101
102func (e *schedEnv) reload(t *testing.T, id int64) *models.Subscription {
103 t.Helper()
104 got, err := e.subRepo.GetByID(id)
105 if err != nil {
106 t.Fatalf("reload subscription %d: %v", id, err)
107 }
108 return got
109}
110
111// A restart must not fire every subscription that happens to have no next run
112// recorded: backfill gives it a time without queueing anything.
113func TestBackfillNextRunsDoesNotRun(t *testing.T) {
114 e := newTestScheduler(t)
115 sub := e.newSub(t, "channel", &models.Subscription{Enabled: true})
116
117 e.sched.backfillNextRuns()
118
119 got := e.reload(t, sub.ID)
120 if !got.NextRunAt.Valid {
121 t.Fatal("next_run_at was not backfilled")
122 }
123 if !got.NextRunAt.Time.After(time.Now()) {
124 t.Errorf("next_run_at = %v, want a future time", got.NextRunAt.Time)
125 }
126 if got.LastStatus.Valid {
127 t.Errorf("last_status = %q, want nothing: backfill must not run the subscription", got.LastStatus.String)
128 }
129 if len(e.pool.queue) != 0 {
130 t.Errorf("%d downloads queued by a backfill", len(e.pool.queue))
131 }
132}
133
134// A disabled subscription is not scheduled at all, so it needs no next run.
135func TestBackfillSkipsDisabled(t *testing.T) {
136 e := newTestScheduler(t)
137 sub := e.newSub(t, "paused", &models.Subscription{Enabled: false})
138
139 e.sched.backfillNextRuns()
140
141 if got := e.reload(t, sub.ID); got.NextRunAt.Valid {
142 t.Errorf("next_run_at = %v, want none for a disabled subscription", got.NextRunAt.Time)
143 }
144}
145
146// An invalid cron expression must not stop the other subscriptions from being
147// backfilled.
148func TestBackfillContinuesPastInvalidSchedule(t *testing.T) {
149 e := newTestScheduler(t)
150 bad := e.newSub(t, "bad", &models.Subscription{
151 Enabled: true, ScheduleKind: "cron", CronExpr: "not a cron",
152 })
153 good := e.newSub(t, "good", &models.Subscription{Enabled: true})
154
155 e.sched.backfillNextRuns()
156
157 if got := e.reload(t, bad.ID); got.NextRunAt.Valid {
158 t.Error("a subscription with a broken schedule should get no next run")
159 }
160 if got := e.reload(t, good.ID); !got.NextRunAt.Valid {
161 t.Error("the valid subscription was skipped after the broken one")
162 }
163}
164
165// A due subscription is queued, recorded as such, and pushed out to its next
166// scheduled time so the same tick doesn't pick it up again.
167func TestCheckDueQueuesAndReschedules(t *testing.T) {
168 e := newTestScheduler(t)
169 sub := e.newSub(t, "channel", &models.Subscription{Enabled: true})
170
171 e.sched.checkDue()
172
173 got := e.reload(t, sub.ID)
174 if got.LastStatus.String != "queued" {
175 t.Errorf("last_status = %q, want queued", got.LastStatus.String)
176 }
177 if !got.NextRunAt.Valid || !got.NextRunAt.Time.After(time.Now()) {
178 t.Errorf("next_run_at = %v, want a future time", got.NextRunAt.Time)
179 }
180 if len(e.pool.queue) != 1 {
181 t.Fatalf("%d downloads queued, want 1", len(e.pool.queue))
182 }
183 d := <-e.pool.queue
184 if !d.SubscriptionID.Valid || d.SubscriptionID.Int64 != sub.ID {
185 t.Errorf("queued download is not tagged with subscription %d", sub.ID)
186 }
187 if d.URL != sub.URL {
188 t.Errorf("queued URL = %q, want %q", d.URL, sub.URL)
189 }
190
191 // The subscription is no longer due, so a second tick queues nothing.
192 e.sched.checkDue()
193 if len(e.pool.queue) != 0 {
194 t.Errorf("%d downloads queued on a tick where nothing was due", len(e.pool.queue))
195 }
196}
197
198// A run longer than the interval must not stack a second one. The subscription
199// is marked "skipped" and pushed forward, keeping its previous run time.
200func TestRunSkipsWhileAnEarlierRunIsActive(t *testing.T) {
201 e := newTestScheduler(t)
202 sub := e.newSub(t, "channel", &models.Subscription{Enabled: true})
203
204 lastRun := time.Now().Add(-2 * time.Hour).Truncate(time.Second)
205 if err := e.subRepo.MarkRun(sub.ID, lastRun, time.Now().Add(-time.Hour), "queued"); err != nil {
206 t.Fatal(err)
207 }
208 active := &models.Download{URL: sub.URL, Status: "downloading"}
209 active.SubscriptionID.Int64, active.SubscriptionID.Valid = sub.ID, true
210 if err := e.repo.Create(active); err != nil {
211 t.Fatal(err)
212 }
213
214 e.sched.checkDue()
215
216 got := e.reload(t, sub.ID)
217 if got.LastStatus.String != "skipped" {
218 t.Errorf("last_status = %q, want skipped", got.LastStatus.String)
219 }
220 if !got.LastRunAt.Time.Equal(lastRun) {
221 t.Errorf("last_run_at = %v, want the earlier run preserved at %v", got.LastRunAt.Time, lastRun)
222 }
223 if !got.NextRunAt.Time.After(time.Now()) {
224 t.Errorf("next_run_at = %v, want it moved into the future", got.NextRunAt.Time)
225 }
226 if len(e.pool.queue) != 0 {
227 t.Errorf("%d downloads queued while a run was already active", len(e.pool.queue))
228 }
229}
230
231// A broken schedule must not be retried every tick; the run is deferred instead.
232func TestRunDefersBrokenSchedule(t *testing.T) {
233 e := newTestScheduler(t)
234 sub := e.newSub(t, "bad", &models.Subscription{
235 Enabled: true, ScheduleKind: "cron", CronExpr: "not a cron",
236 })
237
238 e.sched.checkDue()
239
240 got := e.reload(t, sub.ID)
241 if !got.NextRunAt.Valid || !got.NextRunAt.Time.After(time.Now().Add(23*time.Hour)) {
242 t.Errorf("next_run_at = %v, want it deferred about a day", got.NextRunAt.Time)
243 }
244}
245
246// Stop must return promptly and end the loop. It waits for the loop goroutine,
247// because main closes the database right after it returns.
248func TestSchedulerStop(t *testing.T) {
249 e := newTestScheduler(t)
250 e.sched.Start()
251 stopWithin(t, e.sched, 5*time.Second)
252
253 if err := e.sched.ctx.Err(); err == nil {
254 t.Error("scheduler context still live after Stop")
255 }
256}
257
258// Stop on a scheduler that was never started must return rather than block
259// forever waiting for a loop goroutine that does not exist.
260func TestSchedulerStopWithoutStart(t *testing.T) {
261 e := newTestScheduler(t)
262 stopWithin(t, e.sched, 5*time.Second)
263}
264
265// stopWithin fails the test if Stop has not returned within d, instead of
266// hanging until the whole package times out.
267func stopWithin(t *testing.T, s *Scheduler, d time.Duration) {
268 t.Helper()
269 done := make(chan struct{})
270 go func() {
271 s.Stop()
272 close(done)
273 }()
274 select {
275 case <-done:
276 case <-time.After(d):
277 t.Fatalf("Stop did not return within %v", d)
278 }
279}
280
281// An interval of zero or less would spin the ticker, so the constructor clamps it.
282func TestNewSchedulerRejectsNonPositiveInterval(t *testing.T) {
283 e := newTestScheduler(t)
284 for _, interval := range []time.Duration{0, -time.Second} {
285 s := NewScheduler(e.subSvc, nil, nil, nil, e.pool, interval)
286 if s.interval != time.Minute {
287 t.Errorf("interval %v became %v, want 1m", interval, s.interval)
288 }
289 }
290}
291
292// A comparison flipped the wrong way turns the daily check into one network
293// round trip per tick, with nothing in the UI to show for it.
294//
295// The scheduler's config leaves both tools external and not opted in, so Update
296// finds nothing to do and the test never reaches the network.
297func TestCheckToolUpdateDailyGate(t *testing.T) {
298 now := time.Date(2026, 3, 1, 12, 0, 0, 0, time.UTC)
299
300 cases := []struct {
301 name string
302 autoUpdate bool
303 lastCheck time.Time
304 wantStamp bool
305 }{
306 {"opted out stays quiet", false, now.Add(-25 * time.Hour), false},
307 {"checked an hour ago waits", true, now.Add(-time.Hour), false},
308 {"never checked runs", true, time.Time{}, true},
309 {"checked 25 hours ago runs", true, now.Add(-25 * time.Hour), true},
310 }
311
312 for _, tc := range cases {
313 t.Run(tc.name, func(t *testing.T) {
314 e := newTestScheduler(t)
315 if err := e.settingsSvc.SetToolAutoUpdate(tc.autoUpdate); err != nil {
316 t.Fatal(err)
317 }
318 if !tc.lastCheck.IsZero() {
319 if err := e.settingsSvc.SetToolsCheckedAt(tc.lastCheck); err != nil {
320 t.Fatal(err)
321 }
322 }
323
324 e.sched.checkToolUpdate(now)
325
326 settings, err := e.settingsSvc.GetAll()
327 if err != nil {
328 t.Fatal(err)
329 }
330 // The stamp is written before the run, so it is what says whether the
331 // gate opened.
332 if got := settings.ToolsCheckedAt.Equal(now); got != tc.wantStamp {
333 t.Errorf("stamped = %v, want %v (last check %v)", got, tc.wantStamp, settings.ToolsCheckedAt)
334 }
335 })
336 }
337}
338