package worker import ( "context" "database/sql" "log/slog" "sync" "time" "vidarchive/internal/models" "vidarchive/internal/service" "vidarchive/internal/tools" ) // Scheduler enqueues downloads for subscriptions whose next run is due. It // reuses the normal download pipeline, so refresh-mode and pruning behaviour // live in DownloadService.ExecuteDownload. type Scheduler struct { subscriptionSvc *service.SubscriptionService downloadSvc *service.DownloadService settingsSvc *service.SettingsService tools *tools.Manager pool *Pool interval time.Duration ctx context.Context cancel context.CancelFunc wg sync.WaitGroup } func NewScheduler(subscriptionSvc *service.SubscriptionService, downloadSvc *service.DownloadService, settingsSvc *service.SettingsService, toolMgr *tools.Manager, pool *Pool, interval time.Duration) *Scheduler { if interval <= 0 { interval = time.Minute } ctx, cancel := context.WithCancel(context.Background()) return &Scheduler{ subscriptionSvc: subscriptionSvc, downloadSvc: downloadSvc, settingsSvc: settingsSvc, tools: toolMgr, pool: pool, interval: interval, ctx: ctx, cancel: cancel, } } func (s *Scheduler) Start() { s.backfillNextRuns() s.wg.Add(1) go s.loop() } // Stop waits for the loop because main closes the database next: a checkDue // still in flight would query a closed handle or queue a run nobody services. func (s *Scheduler) Stop() { s.cancel() s.wg.Wait() } func (s *Scheduler) loop() { defer s.wg.Done() ticker := time.NewTicker(s.interval) defer ticker.Stop() // A subscription that came due during downtime should not wait a full interval. s.checkDue() for { select { case <-ticker.C: s.checkDue() s.checkToolUpdate(time.Now()) case <-s.ctx.Done(): return } } } // backfillNextRuns fills in a missing next run time without running anything, so // a restart does not fire every subscription with a null next_run_at. func (s *Scheduler) backfillNextRuns() { subs, err := s.subscriptionSvc.GetAll() if err != nil { slog.Error("scheduler backfill failed to list subscriptions", "err", err) return } now := time.Now() for _, sub := range subs { if !sub.Enabled || sub.NextRunAt.Valid { continue } next, err := s.subscriptionSvc.ComputeNextRun(sub, now) if err != nil { slog.Warn("scheduler: invalid schedule", "subscription_id", sub.ID, "err", err) continue } sub.NextRunAt = sql.NullTime{Time: next, Valid: true} if err := s.subscriptionSvc.Update(sub); err != nil { slog.Error("scheduler: failed to set next run", "subscription_id", sub.ID, "err", err) } } } func (s *Scheduler) checkDue() { now := time.Now() due, err := s.subscriptionSvc.GetDue(now) if err != nil { slog.Error("scheduler: failed to query due subscriptions", "err", err) return } for _, sub := range due { s.run(sub, now) } } // toolUpdateInterval: yt-dlp releases roughly weekly, and its own updater is a // network round trip, so a daily check is plenty. const toolUpdateInterval = 24 * time.Hour func (s *Scheduler) checkToolUpdate(now time.Time) { if s.tools == nil || s.settingsSvc == nil { return } settings, err := s.settingsSvc.GetAll() if err != nil { slog.Error("scheduler: failed to read settings for tool update", "err", err) return } if !settings.ToolAutoUpdate || now.Sub(settings.ToolsCheckedAt) < toolUpdateInterval { return } // Stamped before the run, so a failure waits a day instead of retrying every // tick. It also keeps the goroutine below off the database, which matters // because shutdown does not wait for it. if err := s.settingsSvc.SetToolsCheckedAt(now); err != nil { slog.Error("scheduler: failed to record tool check time", "err", err) return } // An update takes minutes. Blocking the tick would delay every subscription, // so it runs detached. Shutdown cancels ctx and kills the child. Both // updaters replace the binary by rename, so a killed run leaves the old one // intact. go func() { summary, err := s.tools.Update(s.ctx) if err != nil { slog.Error("scheduled tool update failed", "err", err) return } slog.Info("scheduled tool update finished", "result", summary) }() } func (s *Scheduler) run(sub *models.Subscription, now time.Time) { next, err := s.subscriptionSvc.ComputeNextRun(sub, now) if err != nil { // Don't keep retrying a broken schedule every tick; push it out a day. slog.Warn("scheduler: invalid schedule, deferring", "subscription_id", sub.ID, "err", err) next = now.Add(24 * time.Hour) } // A run longer than the interval would stack duplicate downloads. next_run_at // still advances, so this is not re-evaluated every tick. if active, aErr := s.downloadSvc.HasActiveForSubscription(sub.ID); aErr != nil { slog.Error("scheduler: active-run check failed", "subscription_id", sub.ID, "err", aErr) } else if active { lastRun := now if sub.LastRunAt.Valid { lastRun = sub.LastRunAt.Time } slog.Info("scheduler: active run in progress, skipping", "subscription_id", sub.ID, "next_run", next.Format(time.RFC3339)) if mErr := s.subscriptionSvc.MarkRun(sub.ID, lastRun, next, "skipped"); mErr != nil { slog.Error("scheduler: failed to mark run", "subscription_id", sub.ID, "err", mErr) } return } d, err := s.downloadSvc.CreateForSubscription(sub) if err != nil { slog.Error("scheduler: failed to queue subscription", "subscription_id", sub.ID, "err", err) if mErr := s.subscriptionSvc.MarkRun(sub.ID, now, next, "error"); mErr != nil { slog.Error("scheduler: failed to mark run", "subscription_id", sub.ID, "err", mErr) } return } s.pool.Submit(d) if err := s.subscriptionSvc.MarkRun(sub.ID, now, next, "queued"); err != nil { slog.Error("scheduler: failed to mark run", "subscription_id", sub.ID, "err", err) } slog.Info("scheduler: queued subscription", "subscription_id", sub.ID, "download_id", d.ID, "next_run", next.Format(time.RFC3339)) }