mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
- local disk snapshot create and delete
- auto snapshot everyday - rewrite cornman - make dep and some fix - fix conflict
This commit is contained in:
@@ -1,96 +1,192 @@
|
||||
package cronman
|
||||
|
||||
import (
|
||||
"container/heap"
|
||||
"context"
|
||||
"reflect"
|
||||
"runtime"
|
||||
"runtime/debug"
|
||||
"time"
|
||||
|
||||
"yunion.io/x/log"
|
||||
|
||||
"yunion.io/x/onecloud/pkg/appctx"
|
||||
"yunion.io/x/onecloud/pkg/mcclient"
|
||||
"yunion.io/x/onecloud/pkg/mcclient/auth"
|
||||
)
|
||||
|
||||
const (
|
||||
DEFAULT_CRON_INTERVAL = 60 * time.Second // default resolution is 1 monutes
|
||||
)
|
||||
var manager *SCronJobManager
|
||||
|
||||
type SCronJobManager struct {
|
||||
checkInterval time.Duration
|
||||
timer *time.Timer
|
||||
jobs []SCronJob
|
||||
func init() {
|
||||
manager = &SCronJobManager{
|
||||
jobs: make([]*SCronJob, 0),
|
||||
}
|
||||
}
|
||||
|
||||
type ICronTimer interface {
|
||||
Next(time.Time) time.Time
|
||||
}
|
||||
|
||||
type Timer1 struct {
|
||||
dur time.Duration
|
||||
}
|
||||
|
||||
func (t *Timer1) Next(now time.Time) time.Time {
|
||||
return now.Add(t.dur)
|
||||
}
|
||||
|
||||
type Timer2 struct {
|
||||
day, hour, min, sec int
|
||||
}
|
||||
|
||||
func (t *Timer2) Next(now time.Time) time.Time {
|
||||
next := now.Add(time.Hour * time.Duration(t.day) * 24)
|
||||
return time.Date(next.Year(), next.Month(), next.Day(), t.hour, t.min, t.sec, 0, next.Location())
|
||||
}
|
||||
|
||||
type SCronJob struct {
|
||||
name string
|
||||
runInterval time.Duration
|
||||
job func(ctx context.Context, userCred mcclient.TokenCredential)
|
||||
lastRun time.Time
|
||||
Name string
|
||||
job func(ctx context.Context, userCred mcclient.TokenCredential)
|
||||
Timer ICronTimer
|
||||
Next time.Time
|
||||
}
|
||||
|
||||
func NewCronJobManager(interval time.Duration) *SCronJobManager {
|
||||
if interval == 0 {
|
||||
interval = DEFAULT_CRON_INTERVAL
|
||||
type CronJobTimerHeap []*SCronJob
|
||||
|
||||
func (cjth CronJobTimerHeap) Len() int {
|
||||
return len(cjth)
|
||||
}
|
||||
|
||||
func (cjth CronJobTimerHeap) Swap(i, j int) {
|
||||
cjth[i], cjth[j] = cjth[j], cjth[i]
|
||||
}
|
||||
|
||||
func (cjth CronJobTimerHeap) Less(i, j int) bool {
|
||||
if cjth[i].Next.IsZero() {
|
||||
return false
|
||||
}
|
||||
if cjth[j].Next.IsZero() {
|
||||
return true
|
||||
}
|
||||
return cjth[i].Next.Before(cjth[j].Next)
|
||||
}
|
||||
|
||||
func (cjth *CronJobTimerHeap) Push(x interface{}) {
|
||||
*cjth = append(*cjth, x.(*SCronJob))
|
||||
}
|
||||
|
||||
func (cjth *CronJobTimerHeap) Pop() interface{} {
|
||||
old := *cjth
|
||||
n := old.Len()
|
||||
x := old[n-1]
|
||||
*cjth = old[0 : n-1]
|
||||
return x
|
||||
}
|
||||
|
||||
type SCronJobManager struct {
|
||||
jobs CronJobTimerHeap
|
||||
add chan *SCronJob
|
||||
stop chan struct{}
|
||||
running bool
|
||||
}
|
||||
|
||||
func GetCronJobManager() *SCronJobManager {
|
||||
return manager
|
||||
}
|
||||
|
||||
func (self *SCronJobManager) AddJob1(name string, interval time.Duration, jobFunc func(ctx context.Context, userCred mcclient.TokenCredential)) {
|
||||
t := Timer1{
|
||||
dur: interval,
|
||||
}
|
||||
job := SCronJob{
|
||||
Name: name,
|
||||
job: jobFunc,
|
||||
Timer: &t,
|
||||
}
|
||||
if !self.running {
|
||||
self.jobs = append(self.jobs, &job)
|
||||
} else {
|
||||
self.add <- &job
|
||||
}
|
||||
}
|
||||
|
||||
func (self *SCronJobManager) AddJob2(name string, day, hour, min, sec int, jobFunc func(ctx context.Context, userCred mcclient.TokenCredential)) {
|
||||
t := Timer2{
|
||||
day: day,
|
||||
hour: hour,
|
||||
min: min,
|
||||
sec: sec,
|
||||
}
|
||||
job := SCronJob{
|
||||
Name: name,
|
||||
job: jobFunc,
|
||||
Timer: &t,
|
||||
}
|
||||
if !self.running {
|
||||
self.jobs = append(self.jobs, &job)
|
||||
} else {
|
||||
self.add <- &job
|
||||
}
|
||||
}
|
||||
|
||||
func (self *SCronJobManager) Next(now time.Time) {
|
||||
for _, job := range self.jobs {
|
||||
job.Next = job.Timer.Next(now)
|
||||
}
|
||||
cron := SCronJobManager{checkInterval: interval, jobs: make([]SCronJob, 0)}
|
||||
return &cron
|
||||
}
|
||||
|
||||
func (self *SCronJobManager) Start() {
|
||||
if self.timer != nil {
|
||||
if self.running {
|
||||
return
|
||||
}
|
||||
self.timer = time.AfterFunc(self.checkInterval, self.runCronJobs)
|
||||
self.running = true
|
||||
go self.run()
|
||||
}
|
||||
|
||||
func (self *SCronJobManager) Stop() {
|
||||
if self.timer != nil {
|
||||
self.timer.Stop()
|
||||
self.timer = nil
|
||||
}
|
||||
close(self.stop)
|
||||
}
|
||||
|
||||
func getFunctionName(i interface{}) string {
|
||||
return runtime.FuncForPC(reflect.ValueOf(i).Pointer()).Name()
|
||||
}
|
||||
|
||||
func (self *SCronJobManager) AddJob(name string, interval time.Duration, jobFunc func(ctx context.Context, userCred mcclient.TokenCredential)) {
|
||||
// name := getFunctionName(jobFunc)
|
||||
log.Debugf("Add cronjob %s", name)
|
||||
job := SCronJob{name: name, job: jobFunc, runInterval: interval}
|
||||
self.jobs = append(self.jobs, job)
|
||||
}
|
||||
|
||||
func (self *SCronJobManager) runCronJobs() {
|
||||
func (self *SCronJobManager) run() {
|
||||
now := time.Now()
|
||||
for i := 0; i < len(self.jobs); i += 1 {
|
||||
self.jobs[i].run(now)
|
||||
}
|
||||
self.timer = nil
|
||||
self.Start() // schedule next run
|
||||
}
|
||||
|
||||
func (self *SCronJob) run(now time.Time) {
|
||||
if self.lastRun.IsZero() || now.Sub(self.lastRun) >= self.runInterval {
|
||||
log.Debugf("Run cronjob %s", self.name)
|
||||
go runJob(self.name, self.job)
|
||||
self.lastRun = now
|
||||
self.Next(now)
|
||||
heap.Init(&self.jobs)
|
||||
for {
|
||||
var timer *time.Timer
|
||||
if len(self.jobs) == 0 || self.jobs[0].Next.IsZero() {
|
||||
timer = time.NewTimer(100000 * time.Hour)
|
||||
} else {
|
||||
timer = time.NewTimer(self.jobs[0].Next.Sub(now))
|
||||
}
|
||||
select {
|
||||
case now = <-timer.C:
|
||||
for i, job := range self.jobs {
|
||||
if job.Next.After(now) || job.Next.IsZero() {
|
||||
break
|
||||
}
|
||||
go job.runJob()
|
||||
job.Next = job.Timer.Next(now)
|
||||
heap.Fix(&self.jobs, i)
|
||||
}
|
||||
case newJob := <-self.add:
|
||||
now = time.Now()
|
||||
newJob.Next = newJob.Timer.Next(now)
|
||||
heap.Push(&self.jobs, newJob)
|
||||
case <-self.stop:
|
||||
timer.Stop()
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func runJob(name string, job func(ctx context.Context, userCred mcclient.TokenCredential)) {
|
||||
func (job *SCronJob) runJob() {
|
||||
defer func() {
|
||||
if r := recover(); r != nil {
|
||||
log.Errorf("CronJob task %s run error: %s", name, r)
|
||||
log.Errorf("CronJob task %s run error: %s", job.Name, r)
|
||||
debug.PrintStack()
|
||||
}
|
||||
}()
|
||||
|
||||
log.Debugf("Cron job: %s started", job.Name)
|
||||
ctx := context.Background()
|
||||
ctx = context.WithValue(ctx, appctx.APP_CONTEXT_KEY_APPNAME, "Region-Corn-Service")
|
||||
userCred := auth.AdminCredential()
|
||||
job(ctx, userCred)
|
||||
job.job(ctx, userCred)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user