package services import ( "bookhoard/internal/database" "context" "fmt" "log" "sync" "time" "github.com/google/uuid" "github.com/jackc/pgx/v5/pgtype" ) type Scheduler struct { worker *Worker db Database timers map[string]*time.Timer mu sync.RWMutex scanSettings map[string]ScanSetting ctx context.Context cancel context.CancelFunc wg sync.WaitGroup } type Database interface { GetScanSettings(ctx context.Context, id pgtype.UUID) (database.GetScanSettingsRow, error) ListLibraries(ctx context.Context) ([]database.ListLibrariesRow, error) GetLibraryFolders(ctx context.Context, libraryID pgtype.UUID) ([]database.LibraryFolders, error) } type ScanSettingRow struct { ScanFrequencyMinutes pgtype.Int4 AutoScanEnabled pgtype.Bool } type Library struct { ID pgtype.UUID Name string } type LibraryFolder struct { ID pgtype.UUID LibraryID pgtype.UUID FolderPath string } type ScanSetting struct { UserID string Enabled bool Frequency int } func NewScheduler(worker *Worker, db Database) *Scheduler { ctx, cancel := context.WithCancel(context.Background()) return &Scheduler{ worker: worker, db: db, timers: make(map[string]*time.Timer), scanSettings: make(map[string]ScanSetting), ctx: ctx, cancel: cancel, } } func (s *Scheduler) Start() { s.wg.Add(1) go s.runSettingsChecker() } func (s *Scheduler) Stop() { s.cancel() s.mu.Lock() for _, timer := range s.timers { timer.Stop() } s.timers = make(map[string]*time.Timer) s.mu.Unlock() s.wg.Wait() } func (s *Scheduler) runSettingsChecker() { ticker := time.NewTicker(5 * time.Minute) defer ticker.Stop() for { select { case <-ticker.C: s.checkAndScheduleScans() case <-s.ctx.Done(): return } } } func (s *Scheduler) checkAndScheduleScans() { ctx, cancel := context.WithTimeout(s.ctx, 30*time.Second) defer cancel() libraries, err := s.db.ListLibraries(ctx) if err != nil { log.Printf("Error fetching libraries for scan scheduling: %v", err) return } for _, library := range libraries { settings, err := s.db.GetScanSettings(ctx, library.ID) if err != nil { libraryIDStr := fmt.Sprintf("%x", library.ID.Bytes) log.Printf("Error getting scan settings for library %s: %v", libraryIDStr, err) continue } if !settings.AutoScanEnabled.Bool || settings.ScanFrequencyMinutes.Int32 < 15 { continue } userID := fmt.Sprintf("%x", library.ID.Bytes) s.mu.Lock() currentSetting, exists := s.scanSettings[userID] scanSetting := ScanSetting{ UserID: userID, Enabled: settings.AutoScanEnabled.Bool, Frequency: int(settings.ScanFrequencyMinutes.Int32), } s.scanSettings[userID] = scanSetting s.mu.Unlock() if !exists || currentSetting.Frequency != scanSetting.Frequency { s.scheduleLibraryScan(ctx, library.ID, userID, scanSetting.Frequency) } } } func (s *Scheduler) scheduleLibraryScan(ctx context.Context, libraryID pgtype.UUID, userID string, frequencyMinutes int) { userIDStr := fmt.Sprintf("%x", libraryID.Bytes) s.mu.Lock() defer s.mu.Unlock() timerID := userIDStr + "-scan" if existingTimer, exists := s.timers[timerID]; exists { existingTimer.Stop() } duration := time.Duration(frequencyMinutes) * time.Minute timer := time.AfterFunc(duration, func() { s.triggerScheduledScan(ctx, libraryID, userIDStr) s.scheduleLibraryScan(ctx, libraryID, userID, frequencyMinutes) }) s.timers[timerID] = timer libraryIDStr := fmt.Sprintf("%x", libraryID.Bytes) log.Printf("Scheduled scan for library %s every %d minutes", libraryIDStr, frequencyMinutes) } func (s *Scheduler) triggerScheduledScan(ctx context.Context, libraryID pgtype.UUID, userID string) { libraryIDStr := fmt.Sprintf("%x", libraryID.Bytes) log.Printf("Triggering scheduled scan for library %s", libraryIDStr) folders, err := s.db.GetLibraryFolders(ctx, libraryID) if err != nil { log.Printf("Error getting folders for library %s: %v", libraryIDStr, err) return } if len(folders) == 0 { log.Printf("No folders configured for library %s, skipping scan", libraryIDStr) return } folderPaths := make([]string, len(folders)) for i, folder := range folders { folderPaths[i] = folder.FolderPath } job := &Job{ ID: uuid.New().String(), Type: JobTypeScan, Params: map[string]interface{}{ "library_id": fmt.Sprintf("%x", libraryID.Bytes), "folders": folderPaths, "admin_id": userID, "db": s.db, }, Status: JobStatusPending, Context: ctx, } if err := s.worker.EnqueueJob(job); err != nil { log.Printf("Error enqueuing scan job: %v", err) } else { libraryIDStr := fmt.Sprintf("%x", libraryID.Bytes) log.Printf("Enqueued scan job %s for library %s", job.ID, libraryIDStr) } } func (s *Scheduler) UpdateScanSettings(userID string, enabled bool, frequencyMinutes int) { s.mu.Lock() defer s.mu.Unlock() s.scanSettings[userID] = ScanSetting{ UserID: userID, Enabled: enabled, Frequency: frequencyMinutes, } }