package scheduler import ( "context" "fmt" "log/slog" "os" "os/signal" "sync" "syscall" "time" "github.com/robfig/cron/v3" "git.misaka.ren/M1saka/docker_backup/internal/backup" "git.misaka.ren/M1saka/docker_backup/internal/config" "git.misaka.ren/M1saka/docker_backup/internal/storage" ) // Scheduler runs periodic backups using cron expressions. type Scheduler struct { cfg *config.Config dryRun bool logger *slog.Logger locks map[string]*sync.Mutex } // New creates a new scheduler from the config. func New(cfg *config.Config, dryRun bool) (*Scheduler, error) { return &Scheduler{ cfg: cfg, dryRun: dryRun, logger: slog.New(slog.NewTextHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelInfo})), locks: make(map[string]*sync.Mutex), }, nil } // Run starts the cron scheduler and blocks until a shutdown signal is received. func (s *Scheduler) Run(ctx context.Context) error { c := cron.New() engine := backup.NewEngine(s.cfg.Global.BackupDir, s.cfg.Global.TempDir, s.dryRun, s.logger) for _, proj := range s.cfg.Projects { if proj.Cron == "" { continue } proj := proj // capture for closure if _, ok := s.locks[proj.Name]; !ok { s.locks[proj.Name] = &sync.Mutex{} } projectLock := s.locks[proj.Name] _, err := c.AddFunc(proj.Cron, func() { if !projectLock.TryLock() { s.logger.Warn("scheduled backup skipped because previous run is still active", "project", proj.Name) return } defer projectLock.Unlock() // Each cron job uses its own context with a generous timeout. jobCtx, cancel := context.WithTimeout(context.Background(), 2*time.Hour) defer cancel() s.logger.Info("scheduled backup starting", "project", proj.Name) results := engine.Run(jobCtx, []config.ProjectConfig{proj}) for _, r := range results { if r.Error != nil { s.logger.Error("scheduled backup failed", "project", r.ProjectName, "error", r.Error) continue } s.logger.Info("scheduled backup done", "project", r.ProjectName, "file", r.BackupPath, "size", r.FileSize) if s.dryRun { s.logger.Info("dry-run: skipping retention cleanup and remote upload", "project", proj.Name) continue } // Apply retention. policy := backup.RetentionPolicy{ Count: proj.Retention.Count, Days: proj.Retention.Days, } if n, err := backup.ApplyRetention(s.cfg.Global.BackupDir, proj.Name, policy); err != nil { s.logger.Warn("retention cleanup failed", "project", proj.Name, "error", err) } else if n > 0 { s.logger.Info("retention cleaned", "project", proj.Name, "removed", n) } // Upload to remote if configured. if s.cfg.HasRemote() { backends, err := storage.FromConfig(s.cfg.Remote) if err != nil { s.logger.Warn("remote init failed", "error", err) continue } for _, b := range backends { if err := b.Upload(jobCtx, proj.Name, r.BackupPath); err != nil { s.logger.Warn("remote upload failed", "project", proj.Name, "error", err) } else { s.logger.Info("remote upload done", "project", proj.Name) } } } } }) if err != nil { return fmt.Errorf("add cron for %s (%s): %w", proj.Name, proj.Cron, err) } s.logger.Info("scheduled", "project", proj.Name, "cron", proj.Cron) } c.Start() // Wait for shutdown signal. sigCh := make(chan os.Signal, 1) signal.Notify(sigCh, syscall.SIGINT, syscall.SIGTERM) select { case sig := <-sigCh: s.logger.Info("shutting down", "signal", sig) case <-ctx.Done(): s.logger.Info("shutting down", "reason", ctx.Err()) } // Gracefully stop cron (wait for running jobs). stopCtx := c.Stop() <-stopCtx.Done() s.logger.Info("daemon stopped") return nil }