From 225a2b4dbb69d627672d233b747c84b9371aa4e1 Mon Sep 17 00:00:00 2001 From: wanyaoqi Date: Tue, 3 Sep 2019 11:16:53 +0800 Subject: [PATCH] fix snapshot policy filter --- pkg/apis/compute/snapshot_const.go | 1 + pkg/cloudcommon/cronman/cronman.go | 24 ++++++++++++++++-------- pkg/compute/models/disks.go | 5 +++++ 3 files changed, 22 insertions(+), 8 deletions(-) diff --git a/pkg/apis/compute/snapshot_const.go b/pkg/apis/compute/snapshot_const.go index 65e51ec875..b3d8611907 100644 --- a/pkg/apis/compute/snapshot_const.go +++ b/pkg/apis/compute/snapshot_const.go @@ -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" diff --git a/pkg/cloudcommon/cronman/cronman.go b/pkg/cloudcommon/cronman/cronman.go index 9542d5342f..0a1d93daed 100644 --- a/pkg/cloudcommon/cronman/cronman.go +++ b/pkg/cloudcommon/cronman/cronman.go @@ -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) diff --git a/pkg/compute/models/disks.go b/pkg/compute/models/disks.go index 20b7aff041..f90324b30e 100644 --- a/pkg/compute/models/disks.go +++ b/pkg/compute/models/disks.go @@ -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