149 lines
3.5 KiB
Go
149 lines
3.5 KiB
Go
package models
|
|
|
|
import (
|
|
"log"
|
|
"strings"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/robfig/cron/v3"
|
|
)
|
|
|
|
var (
|
|
schedMu sync.Mutex
|
|
schedCron *cron.Cron
|
|
entryTarget = map[cron.EntryID]targetRef{}
|
|
)
|
|
|
|
type targetRef struct {
|
|
scope string // "node" or "svc:<serviceID>"
|
|
targetID string
|
|
}
|
|
|
|
// StartScheduler starts the in-process cron scheduler and performs an initial
|
|
// sync from the current configuration. Call once from main.go before
|
|
// beego.Run() — not from the scan/info CLI paths.
|
|
func StartScheduler() {
|
|
schedMu.Lock()
|
|
schedCron = cron.New(cron.WithChain(cron.Recover(cron.DefaultLogger)))
|
|
schedCron.Start()
|
|
schedMu.Unlock()
|
|
SyncSchedule()
|
|
}
|
|
|
|
// SyncSchedule rebuilds all cron entries from the current backup target
|
|
// definitions (node + every service). Called once at startup and after any
|
|
// target add/update/delete so the running scheduler reflects the latest
|
|
// definitions.
|
|
func SyncSchedule() {
|
|
schedMu.Lock()
|
|
defer schedMu.Unlock()
|
|
if schedCron == nil {
|
|
return
|
|
}
|
|
for _, e := range schedCron.Entries() {
|
|
schedCron.Remove(e.ID)
|
|
}
|
|
entryTarget = map[cron.EntryID]targetRef{}
|
|
|
|
mu.RLock()
|
|
nodeTargets := make([]BackupTarget, len(cfg.Node.Backup.Targets))
|
|
copy(nodeTargets, cfg.Node.Backup.Targets)
|
|
services := make([]Service, len(cfg.Services))
|
|
copy(services, cfg.Services)
|
|
mu.RUnlock()
|
|
|
|
for _, t := range nodeTargets {
|
|
registerJob("node", t)
|
|
}
|
|
for _, s := range services {
|
|
if s.Backup == nil {
|
|
continue
|
|
}
|
|
for _, t := range s.Backup.Targets {
|
|
registerJob("svc:"+s.ID, t)
|
|
}
|
|
}
|
|
refreshNextRunLocked()
|
|
}
|
|
|
|
// registerJob schedules a single target's cron job. Caller must hold schedMu.
|
|
func registerJob(scope string, t BackupTarget) {
|
|
if t.Schedule == "" {
|
|
return
|
|
}
|
|
schedule, err := cron.ParseStandard(t.Schedule)
|
|
if err != nil {
|
|
// Already validated at CRUD time; defensively skip a bad schedule
|
|
// rather than let it wedge the whole sync.
|
|
return
|
|
}
|
|
targetID := t.ID
|
|
job := cron.NewChain(cron.SkipIfStillRunning(cron.DefaultLogger)).Then(cron.FuncJob(func() {
|
|
runScheduledTarget(scope, targetID)
|
|
}))
|
|
id := schedCron.Schedule(schedule, job)
|
|
entryTarget[id] = targetRef{scope: scope, targetID: targetID}
|
|
}
|
|
|
|
func runScheduledTarget(scope, targetID string) {
|
|
var err error
|
|
switch {
|
|
case scope == "node":
|
|
err = RunNodeBackup([]string{targetID})
|
|
case strings.HasPrefix(scope, "svc:"):
|
|
err = RunServiceBackup(strings.TrimPrefix(scope, "svc:"), []string{targetID})
|
|
}
|
|
if err != nil {
|
|
log.Printf("nodemaster: scheduled backup %s/%s failed: %v", scope, targetID, err)
|
|
}
|
|
schedMu.Lock()
|
|
refreshNextRunLocked()
|
|
schedMu.Unlock()
|
|
}
|
|
|
|
// refreshNextRunLocked writes each scheduled target's NextRun from the cron
|
|
// entries' computed next-fire time. Caller must hold schedMu.
|
|
func refreshNextRunLocked() {
|
|
if schedCron == nil {
|
|
return
|
|
}
|
|
updates := map[targetRef]time.Time{}
|
|
for _, e := range schedCron.Entries() {
|
|
if ref, ok := entryTarget[e.ID]; ok {
|
|
updates[ref] = e.Next
|
|
}
|
|
}
|
|
if len(updates) == 0 {
|
|
return
|
|
}
|
|
mu.Lock()
|
|
for ref, next := range updates {
|
|
n := next
|
|
if ref.scope == "node" {
|
|
for i, t := range cfg.Node.Backup.Targets {
|
|
if t.ID == ref.targetID {
|
|
cfg.Node.Backup.Targets[i].NextRun = &n
|
|
break
|
|
}
|
|
}
|
|
continue
|
|
}
|
|
svcID := strings.TrimPrefix(ref.scope, "svc:")
|
|
for i, s := range cfg.Services {
|
|
if s.ID != svcID || s.Backup == nil {
|
|
continue
|
|
}
|
|
for j, t := range s.Backup.Targets {
|
|
if t.ID == ref.targetID {
|
|
cfg.Services[i].Backup.Targets[j].NextRun = &n
|
|
break
|
|
}
|
|
}
|
|
break
|
|
}
|
|
}
|
|
_ = persist()
|
|
mu.Unlock()
|
|
}
|