mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
Merge branch 'release/2.1.0' of ssh://git.yunion.io/~quxuan/onecloud into feature/qx-server-secgroup
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)
|
||||
}
|
||||
|
||||
@@ -172,6 +172,7 @@ func (self *SAliyunGuestDriver) RequestDeployGuestOnHost(ctx context.Context, gu
|
||||
|
||||
if len(guest.SecgrpId) > 0 {
|
||||
if err := iVM.SyncSecurityGroup(guest.SecgrpId, guest.GetSecgroupName(), guest.GetSecRules()); err != nil {
|
||||
log.Errorf("SyncSecurityGroup error: %v", err)
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -3192,3 +3192,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)
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -163,15 +163,20 @@ func (c *MinCounters) Add(counter Counter) {
|
||||
}
|
||||
|
||||
func (c *MinCounters) GetCount() int64 {
|
||||
n := EmptyCapacity
|
||||
for _, c0 := range c.counters {
|
||||
if len(c.counters) == 0 {
|
||||
return EmptyCapacity
|
||||
}
|
||||
minCount := c.counters[0].GetCount()
|
||||
if len(c.counters) == 1 {
|
||||
return minCount
|
||||
}
|
||||
for _, c0 := range c.counters[1:] {
|
||||
count := c0.GetCount()
|
||||
if count != EmptyCapacity && count < n {
|
||||
n = count
|
||||
if count < minCount {
|
||||
minCount = count
|
||||
}
|
||||
}
|
||||
|
||||
return n
|
||||
return minCount
|
||||
}
|
||||
|
||||
type Capacity struct {
|
||||
@@ -448,7 +453,9 @@ func (u *Unit) SetCapacity(id string, name string, capacity Counter) error {
|
||||
|
||||
// Capacity must >= 0
|
||||
if !validateCapacityInput(capacity) {
|
||||
return fmt.Errorf("Capacity invalid: %d", capacity)
|
||||
err := fmt.Errorf("Capacity counter invalid: %#v, count: %d", capacity, capacity.GetCount())
|
||||
log.Errorf("SetCapacity error: %v", err)
|
||||
return err
|
||||
}
|
||||
|
||||
log.V(10).Debugf("%q setCapacity id: %s, capacity: %d", name, id, capacity.GetCount())
|
||||
|
||||
@@ -451,6 +451,6 @@ func (self *SRegion) assignSecurityGroup(secgroupId, instanceId string) error {
|
||||
|
||||
func (self *SRegion) leaveSecurityGroup(secgroupId, instanceId string) error {
|
||||
params := map[string]string{"InstanceId": instanceId, "SecurityGroupId": secgroupId}
|
||||
_, err := self.ecsRequest("JoinSecurityGroup", params)
|
||||
_, err := self.ecsRequest("LeaveSecurityGroup", params)
|
||||
return err
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user