initial commit
This commit is contained in:
@@ -0,0 +1,148 @@
|
||||
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()
|
||||
}
|
||||
Reference in New Issue
Block a user