From f52557a3f80a8bda19ca448c24f040db80ef3f26 Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Sat, 11 Aug 2018 22:21:25 +0800 Subject: [PATCH 1/4] =?UTF-8?q?=E4=BF=AE=E6=AD=A3=EF=BC=9A=E9=98=BF?= =?UTF-8?q?=E9=87=8C=E4=BA=91=E4=B8=BB=E6=9C=BA=E6=B8=85=E7=90=86=E5=9B=9E?= =?UTF-8?q?=E6=94=B6=E7=AB=99=E5=8A=9F=E8=83=BD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pkg/cloudcommon/cronman/cronman.go | 94 ++++++++++++++++++++++++++++++ pkg/compute/models/disks.go | 31 +++++++++- pkg/compute/models/guests.go | 26 +++++++++ pkg/compute/models/vpcs.go | 2 +- pkg/compute/options/options.go | 4 +- pkg/compute/service/service.go | 9 +++ 6 files changed, 161 insertions(+), 5 deletions(-) create mode 100644 pkg/cloudcommon/cronman/cronman.go diff --git a/pkg/cloudcommon/cronman/cronman.go b/pkg/cloudcommon/cronman/cronman.go new file mode 100644 index 0000000000..5d2ad425e1 --- /dev/null +++ b/pkg/cloudcommon/cronman/cronman.go @@ -0,0 +1,94 @@ +package cronman + +import ( + "time" + "runtime/debug" + "context" + + "yunion.io/x/log" + + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/mcclient/auth" + "reflect" + "runtime" +) + +const ( + DEFAULT_CRON_INTERVAL = 60*time.Second // default resolution is 1 monutes +) + +type SCronJobManager struct { + checkInterval time.Duration + timer *time.Timer + jobs []SCronJob +} + +type SCronJob struct { + name string + runInterval time.Duration + job func(ctx context.Context, userCred mcclient.TokenCredential) + lastRun time.Time +} + +func NewCronJobManager(interval time.Duration) *SCronJobManager { + if interval == 0 { + interval = DEFAULT_CRON_INTERVAL + } + cron := SCronJobManager{checkInterval: interval, jobs: make([]SCronJob, 0)} + return &cron +} + +func (self *SCronJobManager) Start() { + if self.timer != nil { + return + } + self.timer = time.AfterFunc(self.checkInterval, self.runCronJobs) +} + +func (self *SCronJobManager) Stop() { + if self.timer != nil { + self.timer.Stop() + self.timer = nil + } +} + +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() { + 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) + } +} + +func runJob(name string, job func(ctx context.Context, userCred mcclient.TokenCredential)) { + defer func() { + if r := recover(); r != nil { + log.Errorf("CronJob task %s run error: %s", name, r) + debug.PrintStack() + } + }() + + ctx := context.Background() + userCred := auth.AdminCredential() + job(ctx, userCred) +} + diff --git a/pkg/compute/models/disks.go b/pkg/compute/models/disks.go index 578f388a8f..4d0b190436 100644 --- a/pkg/compute/models/disks.go +++ b/pkg/compute/models/disks.go @@ -6,11 +6,10 @@ import ( "fmt" "path" "strings" + "time" "yunion.io/x/jsonutils" "yunion.io/x/log" - "yunion.io/x/onecloud/pkg/httperrors" - "yunion.io/x/onecloud/pkg/mcclient" "yunion.io/x/pkg/tristate" "yunion.io/x/pkg/util/compare" "yunion.io/x/pkg/util/fileutils" @@ -24,6 +23,8 @@ import ( "yunion.io/x/onecloud/pkg/cloudcommon/db/taskman" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/compute/options" + "yunion.io/x/onecloud/pkg/httperrors" + "yunion.io/x/onecloud/pkg/mcclient" ) const ( @@ -631,4 +632,30 @@ func (self *SDisk) PerformCancelDelete(ctx context.Context, userCred mcclient.To return nil, err } return nil, nil +} + +func (manager *SDiskManager) getExpiredPendingDeleteDisks() ([]SDisk) { + deadline := time.Now().Add(time.Duration(options.Options.PendingDeleteExpireSeconds)*time.Second) + + q := manager.Query() + q = q.IsTrue("pending_deleted").LT("pending_deleted_at", deadline).Limit(options.Options.PendingDeleteMaxCleanBatchSize) + + disks := make([]SDisk, 0) + err := db.FetchModelObjects(DiskManager, q, &disks) + if err != nil { + log.Errorf("fetch disks error %s", err) + return nil + } + + return disks +} + +func (manager *SDiskManager) CleanPendingDeleteDisks(ctx context.Context, userCred mcclient.TokenCredential) { + disks := manager.getExpiredPendingDeleteDisks() + if disks == nil { + return + } + for i := 0; i < len(disks); i += 1 { + disks[i].StartDiskDeleteTask(ctx, userCred, "", false) + } } \ No newline at end of file diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index a0f6a14fe1..8c93ed56fe 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -2660,3 +2660,29 @@ func (manager *SGuestManager) getIpsByExit(ips []string, isExitOnly bool) []stri } return extRet } + +func (manager *SGuestManager) getExpiredPendingDeleteGuests() ([]SGuest) { + deadline := time.Now().Add(time.Duration(options.Options.PendingDeleteExpireSeconds)*time.Second) + + q := manager.Query() + q = q.IsTrue("pending_deleted").LT("pending_deleted_at", deadline).In("hypervisor", []string{"aliyun"}).Limit(options.Options.PendingDeleteMaxCleanBatchSize) + + guests := make([]SGuest, 0) + err := db.FetchModelObjects(GuestManager, q, &guests) + if err != nil { + log.Errorf("fetch guests error %s", err) + return nil + } + + return guests +} + +func (manager *SGuestManager) CleanPendingDeleteServers(ctx context.Context, userCred mcclient.TokenCredential) { + guests := manager.getExpiredPendingDeleteGuests() + if guests == nil { + return + } + for i := 0; i < len(guests); i += 1 { + guests[i].StartDeleteGuestTask(ctx, userCred, "", false, true) + } +} \ No newline at end of file diff --git a/pkg/compute/models/vpcs.go b/pkg/compute/models/vpcs.go index e4fd13ad5c..bba2aed865 100644 --- a/pkg/compute/models/vpcs.go +++ b/pkg/compute/models/vpcs.go @@ -281,7 +281,7 @@ func (self *SVpc) markAllNetworksUnknown(userCred mcclient.TokenCredential) erro if wires == nil || len(wires) == 0 { return nil } - for i := 0; i <= len(wires); i += 1 { + for i := 0; i < len(wires); i += 1 { wires[i].markNetworkUnknown(userCred) } return nil diff --git a/pkg/compute/options/options.go b/pkg/compute/options/options.go index d62e7c0d50..2d12536217 100644 --- a/pkg/compute/options/options.go +++ b/pkg/compute/options/options.go @@ -22,8 +22,8 @@ type ComputeOptions struct { DefaultDiskSize int `default:"30720" help:"Default disk size in MB if not specified, default to 30GiB"` EnablePendingDelete bool `default:"true" help:"Turn on/off pending delete VM and disk, default is on"` - PendingDeleteExpireSeconds int64 `default:"259200" help:"How long a pending delete VM cleaned automatically, default 3 days"` - PendingDeleteMaxCleanBatchSize int32 `default:"50" help:"How many pending delete servers can be clean in a batch"` + PendingDeleteExpireSeconds int `default:"259200" help:"How long a pending delete VM cleaned automatically, default 3 days"` + PendingDeleteMaxCleanBatchSize int `default:"50" help:"How many pending delete servers can be clean in a batch"` ImageCacheStoragePolicy string `default:"least_used" choices:"best_fit|least_used" help:"Policy to choose storage for image cache, best_fit or least_used"` MetricsRetentionDays int32 `default:"30" help:"Retention days for monitoring metrics in influxdb"` diff --git a/pkg/compute/service/service.go b/pkg/compute/service/service.go index 31d227e934..572b29c08c 100644 --- a/pkg/compute/service/service.go +++ b/pkg/compute/service/service.go @@ -18,6 +18,8 @@ import ( _ "yunion.io/x/onecloud/pkg/util/aliyun/provider" _ "yunion.io/x/onecloud/pkg/util/esxi/provider" + "yunion.io/x/onecloud/pkg/cloudcommon/cronman" + "time" ) func StartService() { @@ -46,6 +48,13 @@ func StartService() { if db.CheckSync(options.Options.AutoSyncTable) { err := models.InitDB() if err == nil { + + cron := cronman.NewCronJobManager(0) + cron.AddJob("CleanPendingDeleteServers", time.Duration(options.Options.PendingDeleteExpireSeconds)*time.Second, models.GuestManager.CleanPendingDeleteServers) + cron.AddJob("CleanPendingDeleteDisks", time.Duration(options.Options.PendingDeleteExpireSeconds)*time.Second, models.DiskManager.CleanPendingDeleteDisks) + cron.Start() + defer cron.Stop() + cloudcommon.ServeForever(app, &options.Options.Options) } else { log.Errorf("InitDB fail: %s", err) From e112524185a6d5baa8cd2cff429e850120c97eae Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Sat, 11 Aug 2018 22:24:06 +0800 Subject: [PATCH 2/4] minor updates --- pkg/compute/service/service.go | 7 +++---- 1 file changed, 3 insertions(+), 4 deletions(-) diff --git a/pkg/compute/service/service.go b/pkg/compute/service/service.go index 572b29c08c..4ae117e023 100644 --- a/pkg/compute/service/service.go +++ b/pkg/compute/service/service.go @@ -2,24 +2,23 @@ package service import ( "os" + "time" _ "github.com/go-sql-driver/mysql" "yunion.io/x/log" + "yunion.io/x/onecloud/pkg/cloudcommon" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/compute" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/compute/options" + "yunion.io/x/onecloud/pkg/cloudcommon/cronman" _ "yunion.io/x/onecloud/pkg/compute/tasks" - _ "yunion.io/x/onecloud/pkg/compute/guestdrivers" - _ "yunion.io/x/onecloud/pkg/util/aliyun/provider" _ "yunion.io/x/onecloud/pkg/util/esxi/provider" - "yunion.io/x/onecloud/pkg/cloudcommon/cronman" - "time" ) func StartService() { From 1f8bf7fa04a45e4398dbcfe5ef4ca4e61f77f010 Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Sat, 11 Aug 2018 22:50:53 +0800 Subject: [PATCH 3/4] =?UTF-8?q?=E4=BF=AE=E6=AD=A3=EF=BC=9A=E5=A2=9E?= =?UTF-8?q?=E5=8A=A0PendingDeleteCheckSeconds=E5=8F=82=E6=95=B0=EF=BC=8C?= =?UTF-8?q?=E8=AE=BE=E7=BD=AE=E5=A4=9A=E4=B9=85=E8=BF=9B=E8=A1=8C=E4=B8=80?= =?UTF-8?q?=E6=AC=A1=E6=89=AB=E6=8F=8F?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pkg/compute/options/options.go | 3 ++- pkg/compute/service/service.go | 4 ++-- 2 files changed, 4 insertions(+), 3 deletions(-) diff --git a/pkg/compute/options/options.go b/pkg/compute/options/options.go index 2d12536217..c873819ebd 100644 --- a/pkg/compute/options/options.go +++ b/pkg/compute/options/options.go @@ -22,7 +22,8 @@ type ComputeOptions struct { DefaultDiskSize int `default:"30720" help:"Default disk size in MB if not specified, default to 30GiB"` EnablePendingDelete bool `default:"true" help:"Turn on/off pending delete VM and disk, default is on"` - PendingDeleteExpireSeconds int `default:"259200" help:"How long a pending delete VM cleaned automatically, default 3 days"` + PendingDeleteCheckSeconds int `default:"3600" help:"How long to wait to scan pending delete VM or disks, default is 1 hour"` + PendingDeleteExpireSeconds int `default:"259200" help:"How long a pending delete VM/disks cleaned automatically, default 3 days"` PendingDeleteMaxCleanBatchSize int `default:"50" help:"How many pending delete servers can be clean in a batch"` ImageCacheStoragePolicy string `default:"least_used" choices:"best_fit|least_used" help:"Policy to choose storage for image cache, best_fit or least_used"` diff --git a/pkg/compute/service/service.go b/pkg/compute/service/service.go index 4ae117e023..c26b0bdab2 100644 --- a/pkg/compute/service/service.go +++ b/pkg/compute/service/service.go @@ -49,8 +49,8 @@ func StartService() { if err == nil { cron := cronman.NewCronJobManager(0) - cron.AddJob("CleanPendingDeleteServers", time.Duration(options.Options.PendingDeleteExpireSeconds)*time.Second, models.GuestManager.CleanPendingDeleteServers) - cron.AddJob("CleanPendingDeleteDisks", time.Duration(options.Options.PendingDeleteExpireSeconds)*time.Second, models.DiskManager.CleanPendingDeleteDisks) + cron.AddJob("CleanPendingDeleteServers", time.Duration(options.Options.PendingDeleteCheckSeconds)*time.Second, models.GuestManager.CleanPendingDeleteServers) + cron.AddJob("CleanPendingDeleteDisks", time.Duration(options.Options.PendingDeleteCheckSeconds)*time.Second, models.DiskManager.CleanPendingDeleteDisks) cron.Start() defer cron.Stop() From f7d815eb641d531eaa5e21d6ff07113d056bdcd4 Mon Sep 17 00:00:00 2001 From: Qiu Jian Date: Sun, 12 Aug 2018 09:11:22 +0800 Subject: [PATCH 4/4] forget to update cronjob lastRun --- pkg/cloudcommon/cronman/cronman.go | 1 + 1 file changed, 1 insertion(+) diff --git a/pkg/cloudcommon/cronman/cronman.go b/pkg/cloudcommon/cronman/cronman.go index 5d2ad425e1..b2e443d983 100644 --- a/pkg/cloudcommon/cronman/cronman.go +++ b/pkg/cloudcommon/cronman/cronman.go @@ -76,6 +76,7 @@ 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 } }