mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-21 06:09:39 +08:00
Merge pull request #2988 from wanyaoqi/feature/wyq/cronjob-support-remove
feature: cronjob support remove job
This commit is contained in:
@@ -18,8 +18,10 @@ import (
|
||||
"container/heap"
|
||||
"context"
|
||||
"runtime/debug"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/pkg/errors"
|
||||
"yunion.io/x/log"
|
||||
|
||||
"yunion.io/x/onecloud/pkg/appctx"
|
||||
@@ -30,6 +32,7 @@ import (
|
||||
|
||||
var (
|
||||
DefaultAdminSessionGenerator = auth.AdminCredential
|
||||
ErrCronJobNameConflict = errors.New("Cron job Name Conflict")
|
||||
)
|
||||
|
||||
type TCronJobFunction func(ctx context.Context, userCred mcclient.TokenCredential, isStart bool)
|
||||
@@ -115,30 +118,53 @@ func (cjth *CronJobTimerHeap) Pop() interface{} {
|
||||
}
|
||||
|
||||
type SCronJobManager struct {
|
||||
jobs CronJobTimerHeap
|
||||
add chan *SCronJob
|
||||
stop chan struct{}
|
||||
running bool
|
||||
workers *appsrv.SWorkerManager
|
||||
jobs CronJobTimerHeap
|
||||
stop chan struct{}
|
||||
running bool
|
||||
workers *appsrv.SWorkerManager
|
||||
dataLock *sync.Mutex
|
||||
}
|
||||
|
||||
func GetCronJobManager(idDbWorker bool) *SCronJobManager {
|
||||
if manager == nil {
|
||||
manager = &SCronJobManager{
|
||||
jobs: make([]*SCronJob, 0),
|
||||
add: make(chan *SCronJob),
|
||||
workers: appsrv.NewWorkerManager("CronJobWorkers", 1, 1024, idDbWorker),
|
||||
jobs: make([]*SCronJob, 0),
|
||||
workers: appsrv.NewWorkerManager("CronJobWorkers", 1, 1024, idDbWorker),
|
||||
dataLock: new(sync.Mutex),
|
||||
}
|
||||
}
|
||||
|
||||
return manager
|
||||
}
|
||||
|
||||
func (self *SCronJobManager) AddJobAtIntervals(name string, interval time.Duration, jobFunc TCronJobFunction) {
|
||||
self.AddJobAtIntervalsWithStartRun(name, interval, jobFunc, false)
|
||||
func (self *SCronJobManager) IsNameUnique(name string) bool {
|
||||
for i := 0; i < len(self.jobs); i++ {
|
||||
if self.jobs[i].Name == name {
|
||||
return false
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
func (self *SCronJobManager) AddJobAtIntervalsWithStartRun(name string, interval time.Duration, jobFunc TCronJobFunction, startRun bool) {
|
||||
func (self *SCronJobManager) String() string {
|
||||
return self.jobs.String()
|
||||
}
|
||||
|
||||
func (self *SCronJobManager) AddJobAtIntervals(name string, interval time.Duration, jobFunc TCronJobFunction) error {
|
||||
return self.AddJobAtIntervalsWithStartRun(name, interval, jobFunc, false)
|
||||
}
|
||||
|
||||
func (self *SCronJobManager) AddJobAtIntervalsWithStartRun(name string, interval time.Duration, jobFunc TCronJobFunction, startRun bool) error {
|
||||
if interval <= 0 {
|
||||
return errors.New("AddJobAtIntervals: interval must > 0")
|
||||
}
|
||||
self.dataLock.Lock()
|
||||
defer self.dataLock.Unlock()
|
||||
|
||||
if !self.IsNameUnique(name) {
|
||||
return ErrCronJobNameConflict
|
||||
}
|
||||
|
||||
t := Timer1{
|
||||
dur: interval,
|
||||
}
|
||||
@@ -151,11 +177,30 @@ func (self *SCronJobManager) AddJobAtIntervalsWithStartRun(name string, interval
|
||||
if !self.running {
|
||||
self.jobs = append(self.jobs, &job)
|
||||
} else {
|
||||
self.add <- &job
|
||||
self.addJob(&job)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (self *SCronJobManager) AddJobEveryFewDays(name string, day, hour, min, sec int, jobFunc TCronJobFunction, startRun bool) {
|
||||
func (self *SCronJobManager) AddJobEveryFewDays(name string, day, hour, min, sec int, jobFunc TCronJobFunction, startRun bool) error {
|
||||
switch {
|
||||
case day <= 0:
|
||||
return errors.New("AddJobEveryFewDays: day must > 0")
|
||||
case hour <= 0:
|
||||
return errors.New("AddJobEveryFewDays: hour must > 0")
|
||||
case min <= 0:
|
||||
return errors.New("AddJobEveryFewDays: min must > 0")
|
||||
case sec <= 0:
|
||||
return errors.New("AddJobEveryFewDays: sec must > 0")
|
||||
}
|
||||
|
||||
self.dataLock.Lock()
|
||||
defer self.dataLock.Unlock()
|
||||
|
||||
if !self.IsNameUnique(name) {
|
||||
return ErrCronJobNameConflict
|
||||
}
|
||||
|
||||
t := Timer2{
|
||||
day: day,
|
||||
hour: hour,
|
||||
@@ -171,11 +216,28 @@ func (self *SCronJobManager) AddJobEveryFewDays(name string, day, hour, min, sec
|
||||
if !self.running {
|
||||
self.jobs = append(self.jobs, &job)
|
||||
} else {
|
||||
self.add <- &job
|
||||
self.addJob(&job)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (self *SCronJobManager) AddJobEveryFewHour(name string, hour, min, sec int, jobFunc TCronJobFunction, startRun bool) {
|
||||
func (self *SCronJobManager) AddJobEveryFewHour(name string, hour, min, sec int, jobFunc TCronJobFunction, startRun bool) error {
|
||||
switch {
|
||||
case hour <= 0:
|
||||
return errors.New("AddJobEveryFewHour: hour must > 0")
|
||||
case min <= 0:
|
||||
return errors.New("AddJobEveryFewHour: min must > 0")
|
||||
case sec <= 0:
|
||||
return errors.New("AddJobEveryFewHour: sec must > 0")
|
||||
}
|
||||
|
||||
self.dataLock.Lock()
|
||||
defer self.dataLock.Unlock()
|
||||
|
||||
if !self.IsNameUnique(name) {
|
||||
return ErrCronJobNameConflict
|
||||
}
|
||||
|
||||
t := TimerHour{
|
||||
hour: hour,
|
||||
min: min,
|
||||
@@ -190,11 +252,39 @@ func (self *SCronJobManager) AddJobEveryFewHour(name string, hour, min, sec int,
|
||||
if !self.running {
|
||||
self.jobs = append(self.jobs, &job)
|
||||
} else {
|
||||
self.add <- &job
|
||||
self.addJob(&job)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (self *SCronJobManager) Next(now time.Time) {
|
||||
func (self *SCronJobManager) addJob(newJob *SCronJob) {
|
||||
now := time.Now()
|
||||
newJob.Next = newJob.Timer.Next(now)
|
||||
if newJob.StartRun {
|
||||
newJob.runJob(true)
|
||||
}
|
||||
heap.Push(&self.jobs, newJob)
|
||||
}
|
||||
|
||||
func (self *SCronJobManager) Remove(name string) error {
|
||||
self.dataLock.Lock()
|
||||
defer self.dataLock.Unlock()
|
||||
|
||||
var jobIndex = -1
|
||||
for i := 0; i < len(self.jobs); i++ {
|
||||
if self.jobs[i].Name == name {
|
||||
jobIndex = i
|
||||
break
|
||||
}
|
||||
}
|
||||
if jobIndex == -1 {
|
||||
return errors.Errorf("job %s not found", name)
|
||||
}
|
||||
heap.Remove(&self.jobs, jobIndex)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (self *SCronJobManager) next(now time.Time) {
|
||||
for _, job := range self.jobs {
|
||||
job.Next = job.Timer.Next(now)
|
||||
}
|
||||
@@ -204,6 +294,8 @@ func (self *SCronJobManager) Start() {
|
||||
if self.running {
|
||||
return
|
||||
}
|
||||
self.dataLock.Lock()
|
||||
defer self.dataLock.Unlock()
|
||||
self.running = true
|
||||
self.init()
|
||||
go self.run()
|
||||
@@ -217,7 +309,7 @@ func (self *SCronJobManager) Stop() {
|
||||
|
||||
func (self *SCronJobManager) init() {
|
||||
now := time.Now()
|
||||
self.Next(now)
|
||||
self.next(now)
|
||||
heap.Init(&self.jobs)
|
||||
for i := 0; i < len(self.jobs); i += 1 {
|
||||
if self.jobs[i].StartRun {
|
||||
@@ -228,24 +320,19 @@ func (self *SCronJobManager) init() {
|
||||
}
|
||||
|
||||
func (self *SCronJobManager) run() {
|
||||
var timer *time.Timer
|
||||
var now = time.Now()
|
||||
for {
|
||||
now := time.Now()
|
||||
var timer *time.Timer
|
||||
self.dataLock.Lock()
|
||||
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))
|
||||
}
|
||||
self.dataLock.Unlock()
|
||||
select {
|
||||
case now = <-timer.C:
|
||||
self.runJob(now)
|
||||
case newJob := <-self.add:
|
||||
now = time.Now()
|
||||
newJob.Next = newJob.Timer.Next(now)
|
||||
if newJob.StartRun {
|
||||
newJob.runJob(true)
|
||||
}
|
||||
heap.Push(&self.jobs, newJob)
|
||||
self.runJobs(now)
|
||||
case <-self.stop:
|
||||
timer.Stop()
|
||||
return
|
||||
@@ -253,12 +340,15 @@ func (self *SCronJobManager) run() {
|
||||
}
|
||||
}
|
||||
|
||||
func (self *SCronJobManager) runJob(now time.Time) {
|
||||
if len(self.jobs) > 0 && !(self.jobs[0].Next.After(now) || self.jobs[0].Next.IsZero()) {
|
||||
self.jobs[0].runJob(false)
|
||||
self.jobs[0].Next = self.jobs[0].Timer.Next(now)
|
||||
heap.Fix(&self.jobs, 0)
|
||||
self.runJob(now)
|
||||
func (self *SCronJobManager) runJobs(now time.Time) {
|
||||
self.dataLock.Lock()
|
||||
defer self.dataLock.Unlock()
|
||||
for i := 0; i < len(self.jobs); i++ {
|
||||
if !(self.jobs[i].Next.After(now) || self.jobs[i].Next.IsZero()) {
|
||||
self.jobs[i].runJob(false)
|
||||
self.jobs[i].Next = self.jobs[i].Timer.Next(now)
|
||||
heap.Fix(&self.jobs, i)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -276,7 +366,7 @@ func (job *SCronJob) runJobInWorker(isStart bool) {
|
||||
}
|
||||
}()
|
||||
|
||||
// log.Debugf("Cron job: %s started", job.Name)
|
||||
log.Debugf("Cron job: %s started", job.Name)
|
||||
ctx := context.Background()
|
||||
ctx = context.WithValue(ctx, appctx.APP_CONTEXT_KEY_APPNAME, "Cron-Service")
|
||||
userCred := DefaultAdminSessionGenerator()
|
||||
|
||||
@@ -0,0 +1,42 @@
|
||||
// Copyright 2019 Yunion
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
package cronman
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"yunion.io/x/onecloud/pkg/mcclient"
|
||||
)
|
||||
|
||||
func TestSCronJobManager_AddRemoveJobs(t *testing.T) {
|
||||
manager := GetCronJobManager(false)
|
||||
testFunc := func(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) {}
|
||||
manager.AddJobAtIntervals("Test1", time.Second*100, testFunc)
|
||||
|
||||
manager.AddJobAtIntervals("Test2", time.Second*100, testFunc)
|
||||
manager.AddJobEveryFewDays("Test3", 1, 1, 1, 1, testFunc, false)
|
||||
manager.AddJobEveryFewDays("Test4", 1, 1, 1, 1, testFunc, false)
|
||||
manager.AddJobEveryFewDays("Test5", 1, 1, 1, 1, testFunc, false)
|
||||
manager.Start()
|
||||
t.Logf("Jobs \n%s", manager.String())
|
||||
manager.Remove("Test1")
|
||||
manager.Remove("Test2")
|
||||
manager.Remove("Test3")
|
||||
manager.AddJobAtIntervals("Test6", time.Second*100, testFunc)
|
||||
manager.AddJobEveryFewDays("Test7", 1, 1, 1, 1, testFunc, false)
|
||||
t.Logf("Jobs \n%s", manager.String())
|
||||
}
|
||||
Reference in New Issue
Block a user