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