diff --git a/pkg/cloudcommon/cronman/cronman.go b/pkg/cloudcommon/cronman/cronman.go new file mode 100644 index 0000000000..b2e443d983 --- /dev/null +++ b/pkg/cloudcommon/cronman/cronman.go @@ -0,0 +1,95 @@ +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) + self.lastRun = now + } +} + +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 b9c027c43a..a5e552fd32 100644 --- a/pkg/compute/models/disks.go +++ b/pkg/compute/models/disks.go @@ -6,6 +6,7 @@ import ( "fmt" "path" "strings" + "time" "yunion.io/x/jsonutils" "yunion.io/x/log" @@ -844,3 +845,29 @@ func (self *SDisk) PerformCancelDelete(ctx context.Context, userCred mcclient.To } 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) + } +} diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index 81eb493186..5b563fed7f 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -3127,3 +3127,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) + } +} 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..c873819ebd 100644 --- a/pkg/compute/options/options.go +++ b/pkg/compute/options/options.go @@ -22,8 +22,9 @@ 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"` + 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"` 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 5a2a4ffdf1..f8cac5dbae 100644 --- a/pkg/compute/service/service.go +++ b/pkg/compute/service/service.go @@ -2,6 +2,7 @@ package service import ( "os" + "time" "yunion.io/x/log" "yunion.io/x/sqlchemy" @@ -14,6 +15,7 @@ import ( _ "yunion.io/x/onecloud/pkg/util/esxi/provider" "yunion.io/x/onecloud/pkg/cloudcommon" + "yunion.io/x/onecloud/pkg/cloudcommon/cronman" "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/compute" "yunion.io/x/onecloud/pkg/compute/models" @@ -50,6 +52,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.PendingDeleteCheckSeconds)*time.Second, models.GuestManager.CleanPendingDeleteServers) + cron.AddJob("CleanPendingDeleteDisks", time.Duration(options.Options.PendingDeleteCheckSeconds)*time.Second, models.DiskManager.CleanPendingDeleteDisks) + cron.Start() + defer cron.Stop() + cloudcommon.ServeForever(app, &options.Options.Options) } else { log.Errorf("InitDB fail: %s", err)