Files
2026-07-27 14:41:04 +02:00

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()
}