package worker import ( "os" "path/filepath" "testing" "time" "vidarchive/internal/config" "vidarchive/internal/database" "vidarchive/internal/models" "vidarchive/internal/repository" "vidarchive/internal/service" ) type schedEnv struct { sched *Scheduler pool *Pool subSvc *service.SubscriptionService subRepo *repository.SubscriptionRepository repo *repository.DownloadRepository } // newTestScheduler builds a scheduler over a real database. The pool is created // but never started, so a queued run stays in the buffer where a test can count // it instead of being executed. func newTestScheduler(t *testing.T) *schedEnv { t.Helper() root := t.TempDir() cfg := &config.Config{ DBPath: ":memory:", DataDir: root, LibraryDir: filepath.Join(root, "library"), TempDir: filepath.Join(root, "temp"), YTDLPPath: "/bin/false", FFmpegPath: "ffmpeg", FFprobePath: "ffprobe", } for _, dir := range []string{cfg.LibraryDir, cfg.TempDir} { if err := os.MkdirAll(dir, 0o755); err != nil { t.Fatal(err) } } db, err := database.New(cfg) if err != nil { t.Fatalf("init db: %v", err) } t.Cleanup(func() { db.Close() }) downloadRepo := repository.NewDownloadRepository(db) subRepo := repository.NewSubscriptionRepository(db) subSvc := service.NewSubscriptionService(subRepo, cfg) downloadSvc := service.NewDownloadService( downloadRepo, service.NewLibraryService(cfg.LibraryDir, cfg.FFmpegPath, cfg.FFprobePath), service.NewPresetService(repository.NewPresetRepository(db)), service.NewSettingsService(repository.NewSettingsRepository(db)), subSvc, cfg, ) pool := New(downloadSvc, 1) return &schedEnv{ sched: NewScheduler(subSvc, downloadSvc, pool, time.Minute), pool: pool, subSvc: subSvc, subRepo: subRepo, repo: downloadRepo, } } func (e *schedEnv) newSub(t *testing.T, name string, sub *models.Subscription) *models.Subscription { t.Helper() sub.Name = name if sub.URL == "" { sub.URL = "https://example.com/" + name } if sub.RefreshMode == "" { sub.RefreshMode = "overwrite" } if sub.ScheduleKind == "" { sub.ScheduleKind = "daily" } if sub.CronExpr == "" { sub.CronExpr = "0 3 * * *" } if sub.OutputDir == "" { sub.OutputDir = name } if err := e.subRepo.Create(sub); err != nil { t.Fatalf("create subscription %s: %v", name, err) } return sub } func (e *schedEnv) reload(t *testing.T, id int64) *models.Subscription { t.Helper() got, err := e.subRepo.GetByID(id) if err != nil { t.Fatalf("reload subscription %d: %v", id, err) } return got } // A restart must not fire every subscription that happens to have no next run // recorded: backfill gives it a time without queueing anything. func TestBackfillNextRunsDoesNotRun(t *testing.T) { e := newTestScheduler(t) sub := e.newSub(t, "channel", &models.Subscription{Enabled: true}) e.sched.backfillNextRuns() got := e.reload(t, sub.ID) if !got.NextRunAt.Valid { t.Fatal("next_run_at was not backfilled") } if !got.NextRunAt.Time.After(time.Now()) { t.Errorf("next_run_at = %v, want a future time", got.NextRunAt.Time) } if got.LastStatus.Valid { t.Errorf("last_status = %q, want nothing: backfill must not run the subscription", got.LastStatus.String) } if len(e.pool.queue) != 0 { t.Errorf("%d downloads queued by a backfill", len(e.pool.queue)) } } // A disabled subscription is not scheduled at all, so it needs no next run. func TestBackfillSkipsDisabled(t *testing.T) { e := newTestScheduler(t) sub := e.newSub(t, "paused", &models.Subscription{Enabled: false}) e.sched.backfillNextRuns() if got := e.reload(t, sub.ID); got.NextRunAt.Valid { t.Errorf("next_run_at = %v, want none for a disabled subscription", got.NextRunAt.Time) } } // An invalid cron expression must not stop the other subscriptions from being // backfilled. func TestBackfillContinuesPastInvalidSchedule(t *testing.T) { e := newTestScheduler(t) bad := e.newSub(t, "bad", &models.Subscription{ Enabled: true, ScheduleKind: "cron", CronExpr: "not a cron", }) good := e.newSub(t, "good", &models.Subscription{Enabled: true}) e.sched.backfillNextRuns() if got := e.reload(t, bad.ID); got.NextRunAt.Valid { t.Error("a subscription with a broken schedule should get no next run") } if got := e.reload(t, good.ID); !got.NextRunAt.Valid { t.Error("the valid subscription was skipped after the broken one") } } // A due subscription is queued, recorded as such, and pushed out to its next // scheduled time so the same tick doesn't pick it up again. func TestCheckDueQueuesAndReschedules(t *testing.T) { e := newTestScheduler(t) sub := e.newSub(t, "channel", &models.Subscription{Enabled: true}) e.sched.checkDue() got := e.reload(t, sub.ID) if got.LastStatus.String != "queued" { t.Errorf("last_status = %q, want queued", got.LastStatus.String) } if !got.NextRunAt.Valid || !got.NextRunAt.Time.After(time.Now()) { t.Errorf("next_run_at = %v, want a future time", got.NextRunAt.Time) } if len(e.pool.queue) != 1 { t.Fatalf("%d downloads queued, want 1", len(e.pool.queue)) } d := <-e.pool.queue if !d.SubscriptionID.Valid || d.SubscriptionID.Int64 != sub.ID { t.Errorf("queued download is not tagged with subscription %d", sub.ID) } if d.URL != sub.URL { t.Errorf("queued URL = %q, want %q", d.URL, sub.URL) } // The subscription is no longer due, so a second tick queues nothing. e.sched.checkDue() if len(e.pool.queue) != 0 { t.Errorf("%d downloads queued on a tick where nothing was due", len(e.pool.queue)) } } // A run longer than the interval must not stack a second one. The subscription // is marked "skipped" and pushed forward, keeping its previous run time. func TestRunSkipsWhileAnEarlierRunIsActive(t *testing.T) { e := newTestScheduler(t) sub := e.newSub(t, "channel", &models.Subscription{Enabled: true}) lastRun := time.Now().Add(-2 * time.Hour).Truncate(time.Second) if err := e.subRepo.MarkRun(sub.ID, lastRun, time.Now().Add(-time.Hour), "queued"); err != nil { t.Fatal(err) } // An in-flight download for this subscription. active := &models.Download{URL: sub.URL, Status: "downloading"} active.SubscriptionID.Int64, active.SubscriptionID.Valid = sub.ID, true if err := e.repo.Create(active); err != nil { t.Fatal(err) } e.sched.checkDue() got := e.reload(t, sub.ID) if got.LastStatus.String != "skipped" { t.Errorf("last_status = %q, want skipped", got.LastStatus.String) } if !got.LastRunAt.Time.Equal(lastRun) { t.Errorf("last_run_at = %v, want the earlier run preserved at %v", got.LastRunAt.Time, lastRun) } if !got.NextRunAt.Time.After(time.Now()) { t.Errorf("next_run_at = %v, want it moved into the future", got.NextRunAt.Time) } if len(e.pool.queue) != 0 { t.Errorf("%d downloads queued while a run was already active", len(e.pool.queue)) } } // A broken schedule must not be retried every tick; the run is deferred instead. func TestRunDefersBrokenSchedule(t *testing.T) { e := newTestScheduler(t) sub := e.newSub(t, "bad", &models.Subscription{ Enabled: true, ScheduleKind: "cron", CronExpr: "not a cron", }) e.sched.checkDue() got := e.reload(t, sub.ID) if !got.NextRunAt.Valid || !got.NextRunAt.Time.After(time.Now().Add(23*time.Hour)) { t.Errorf("next_run_at = %v, want it deferred about a day", got.NextRunAt.Time) } } // Stop must return promptly and end the loop. func TestSchedulerStop(t *testing.T) { e := newTestScheduler(t) e.sched.Start() e.sched.Stop() if err := e.sched.ctx.Err(); err == nil { t.Error("scheduler context still live after Stop") } } // An interval of zero or less would spin the ticker, so the constructor clamps it. func TestNewSchedulerRejectsNonPositiveInterval(t *testing.T) { e := newTestScheduler(t) for _, interval := range []time.Duration{0, -time.Second} { s := NewScheduler(e.subSvc, nil, e.pool, interval) if s.interval != time.Minute { t.Errorf("interval %v became %v, want 1m", interval, s.interval) } } }