package scheduler import ( "context" "database/sql" "fmt" "log/slog" "os" "path/filepath" "strconv" "strings" "sync" "time" "github.com/syncserver/internal/config" "github.com/syncserver/internal/models" "github.com/syncserver/internal/syncengine" ) type Scheduler struct { db *sql.DB engine *syncengine.Engine cfg *config.Config stopCh chan struct{} wg sync.WaitGroup } func New(database interface{ SQLDB() *sql.DB }, engine *syncengine.Engine, cfg *config.Config) *Scheduler { return &Scheduler{ db: database.SQLDB(), engine: engine, cfg: cfg, stopCh: make(chan struct{}), } } func (s *Scheduler) Start() { s.wg.Add(1) go s.run() s.wg.Add(1) go s.cleanupRun() slog.Info("scheduler started") } func (s *Scheduler) Stop() { close(s.stopCh) s.wg.Wait() slog.Info("scheduler stopped") } func (s *Scheduler) run() { defer s.wg.Done() ticker := time.NewTicker(1 * time.Minute) defer ticker.Stop() for { select { case <-s.stopCh: return case <-ticker.C: s.tick() } } } func (s *Scheduler) tick() { scheduleRepo := models.NewScheduleRepository(s.db) now := time.Now().UTC() schedules, err := scheduleRepo.GetEnabledDue(now) if err != nil { slog.Error("scheduler: failed to get due schedules", "error", err) return } for _, sch := range schedules { pairRepo := models.NewSyncPairRepository(s.db) pair, err := pairRepo.GetByID(sch.SyncPairID) if err != nil || !pair.Enabled { continue } jobID, err := s.engine.CreateJob(sch.SyncPairID, "scheduled") if err != nil { slog.Error("scheduler: failed to create job", "schedule_id", sch.ID, "error", err) continue } ctx := context.Background() go func(jobID int64, pairID int64, schID int64) { defer func() { if r := recover(); r != nil { slog.Error("scheduler: job run panicked", "job_id", jobID, "panic", r) } }() if err := s.engine.Run(ctx, jobID, pairID); err != nil { slog.Warn("scheduler: job failed", "job_id", jobID, "error", err) } expr, _ := ParseCron(sch.CronExpr) if expr != nil { next := NextRun(expr, time.Now().UTC()) scheduleRepo.UpdateNextRun(schID, next) } }(jobID, sch.SyncPairID, sch.ID) } } func (s *Scheduler) cleanupRun() { defer s.wg.Done() ticker := time.NewTicker(24 * time.Hour) defer ticker.Stop() s.cleanup() for { select { case <-s.stopCh: return case <-ticker.C: s.cleanup() } } } func (s *Scheduler) cleanup() { retentionDays := s.cfg.Scheduler.RetentionDays if retentionDays <= 0 { return } before := time.Now().AddDate(0, 0, -retentionDays) if s.cfg.Scheduler.BackupDir != "" { s.backupDB(before) } logRepo := models.NewJobLogRepository(s.db) jobRepo := models.NewJobRepository(s.db) deletedLogs, err := logRepo.DeleteBefore(before) if err != nil { slog.Error("cleanup: failed to purge old job logs", "error", err) return } deletedJobs, err := jobRepo.DeleteFinishedBefore(before) if err != nil { slog.Error("cleanup: failed to purge old jobs", "error", err) return } if deletedLogs > 0 || deletedJobs > 0 { slog.Info("cleanup: purged old records", "logs_deleted", deletedLogs, "jobs_deleted", deletedJobs, "before", before.Format("2006-01-02")) s.purgeJobLogFiles(before) } s.purgeOldBackups() } func (s *Scheduler) backupDB(before time.Time) { backupDir := s.cfg.Scheduler.BackupDir if backupDir == "" { return } if err := os.MkdirAll(backupDir, 0700); err != nil { slog.Error("cleanup: failed to create backup dir", "error", err) return } ts := time.Now().UTC().Format("20060102-150405") backupPath := filepath.Join(backupDir, fmt.Sprintf("syncserver-%s.db", ts)) if _, err := s.db.Exec(fmt.Sprintf("VACUUM INTO '%s'", backupPath)); err != nil { slog.Error("cleanup: failed to vacuum into backup", "path", backupPath, "error", err) return } slog.Info("cleanup: database backup created", "path", backupPath) } func (s *Scheduler) purgeJobLogFiles(before time.Time) { logsDir := s.cfg.LogsDir() entries, err := os.ReadDir(logsDir) if err != nil { return } for _, entry := range entries { if entry.IsDir() || !strings.HasSuffix(entry.Name(), ".log") { continue } jobID := strings.TrimSuffix(entry.Name(), ".log") id, err := strconv.ParseInt(jobID, 10, 64) if err != nil { continue } jobRepo := models.NewJobRepository(s.db) job, err := jobRepo.GetByID(id) if err != nil || job == nil { os.Remove(filepath.Join(logsDir, entry.Name())) continue } if job.FinishedAt != nil && job.FinishedAt.Before(before) { os.Remove(filepath.Join(logsDir, entry.Name())) } } } func (s *Scheduler) purgeOldBackups() { backupDir := s.cfg.Scheduler.BackupDir retention := s.cfg.Scheduler.BackupRetentionDays if backupDir == "" || retention <= 0 { return } cutoff := time.Now().AddDate(0, 0, -retention) entries, err := os.ReadDir(backupDir) if err != nil { return } for _, entry := range entries { if entry.IsDir() || !strings.HasSuffix(entry.Name(), ".db") { continue } info, err := entry.Info() if err != nil { continue } if info.ModTime().Before(cutoff) { os.Remove(filepath.Join(backupDir, entry.Name())) } } }