mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-21 14:19:49 +08:00
Merge branch 'release/2.0.0' of ssh://git.yunion.io/~qiujian/onecloud into hotfix/qj-conflict-resolve-20180813
Conflicts: pkg/compute/models/disks.go pkg/compute/service/service.go
This commit is contained in:
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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"`
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user