fix snapshot policy filter

This commit is contained in:
wanyaoqi
2019-09-04 11:52:16 +08:00
parent 867ab0e747
commit 225a2b4dbb
3 changed files with 22 additions and 8 deletions
+1
View File
@@ -39,6 +39,7 @@ const (
SNAPSHOT_POLICY_CANCEL = "canceling"
SNAPSHOT_POLICY_CANCEL_FAILED = "cancel_failed"
SNAPSHOT_POLICY_DISK_INIT = "init"
SNAPSHOT_POLICY_DISK_READY = "ready"
SNAPSHOT_POLICY_DISK_DELETING = "deleting"
SNAPSHOT_POLICY_DISK_DELETE_FAILED = "delete_failed"
+16 -8
View File
@@ -205,6 +205,7 @@ func (self *SCronJobManager) Start() {
return
}
self.running = true
self.init()
go self.run()
}
@@ -214,29 +215,36 @@ func (self *SCronJobManager) Stop() {
}
}
func (self *SCronJobManager) run() {
func (self *SCronJobManager) init() {
now := time.Now()
self.Next(now)
heap.Init(&self.jobs)
for i := 0; i < len(self.jobs); i += 1 {
if self.jobs[i].StartRun {
self.jobs[i].StartRun = false
self.jobs[i].runJob(true)
}
}
}
func (self *SCronJobManager) run() {
for {
now := time.Now()
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))
}
for i := 0; i < len(self.jobs); i += 1 {
if self.jobs[i].StartRun {
self.jobs[i].StartRun = false
self.jobs[i].runJob(true)
}
}
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)
case <-self.stop:
timer.Stop()
@@ -246,7 +254,7 @@ 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()) {
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)
+5
View File
@@ -1908,7 +1908,12 @@ func (manager *SDiskManager) getAutoSnapshotDisksId() ([]SSnapshotPolicyDisk, er
spds := make([]SSnapshotPolicyDisk, 0)
spdq := SnapshotPolicyDiskManager.Query()
spdq.NotEquals("status", api.SNAPSHOT_POLICY_DISK_INIT)
spdq.Filter(sqlchemy.In(spdq.Field("snapshotpolicy_id"), sps))
diskQ := DiskManager.Query().SubQuery()
spdq.Join(diskQ, sqlchemy.Equals(spdq.Field("disk_id"), diskQ.Field("id")))
spdq.Filter(sqlchemy.IsNullOrEmpty(diskQ.Field("external_id")))
err = spdq.All(&spds)
if err != nil {
return nil, err