Files
docker-compose-backup/internal/scheduler/scheduler.go
T
2026-07-04 17:20:28 +08:00

131 lines
3.7 KiB
Go

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
}