package services import ( "bookhoard/internal/database" "context" "fmt" "log" "strconv" "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 { GetSystemSetting(ctx context.Context, settingKey string) (string, error) ListLibraries(ctx context.Context) ([]database.ListLibrariesRow, error) GetLibraryFolders(ctx context.Context, libraryID pgtype.UUID) ([]database.LibraryFolders, error) } 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 func() { defer s.wg.Done() s.runSettingsChecker() }() } func (s *Scheduler) Stop() { s.cancel() s.mu.Lock() for _, timer := range s.timers { if timer != nil { 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 } autoScanEnabledStr, err := s.db.GetSystemSetting(ctx, "auto_scan_enabled") if err != nil { log.Printf("Error getting auto_scan_enabled setting: %v", err) return } autoScanEnabled, err := strconv.ParseBool(autoScanEnabledStr) if err != nil { log.Printf("Error parsing auto_scan_enabled: %v", err) return } if !autoScanEnabled { return } scanFrequencyStr, err := s.db.GetSystemSetting(ctx, "scan_frequency_minutes") if err != nil { log.Printf("Error getting scan_frequency_minutes setting: %v", err) return } scanFrequency, err := strconv.Atoi(scanFrequencyStr) if err != nil { log.Printf("Error parsing scan_frequency_minutes: %v", err) return } if scanFrequency < 15 { return } for _, library := range libraries { libraryIDStr := fmt.Sprintf("%x", library.ID.Bytes) s.mu.Lock() currentSetting, exists := s.scanSettings[libraryIDStr] scanSetting := ScanSetting{ UserID: libraryIDStr, Enabled: autoScanEnabled, Frequency: scanFrequency, } s.scanSettings[libraryIDStr] = scanSetting s.mu.Unlock() if !exists || currentSetting.Frequency != scanSetting.Frequency { s.scheduleLibraryScan(ctx, library.ID, libraryIDStr, 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, } }