Merge pull request #2406 from wanyaoqi/feature/wyq/suporrt-custom-snapshot-policy

feature: snapshot support custom snapshot policy
This commit is contained in:
yunion-ci-robot
2019-08-21 15:29:06 +08:00
committed by GitHub
19 changed files with 310 additions and 157 deletions
+50 -12
View File
@@ -57,6 +57,15 @@ func (t *Timer2) Next(now time.Time) time.Time {
return time.Date(next.Year(), next.Month(), next.Day(), t.hour, t.min, t.sec, 0, next.Location())
}
type TimerHour struct {
hour, min, sec int
}
func (t *TimerHour) Next(now time.Time) time.Time {
next := now.Add(time.Hour * time.Duration(t.hour))
return time.Date(next.Year(), next.Month(), next.Day(), next.Hour(), t.min, t.sec, 0, next.Location())
}
type SCronJob struct {
Name string
job TCronJobFunction
@@ -67,6 +76,14 @@ type SCronJob struct {
type CronJobTimerHeap []*SCronJob
func (c CronJobTimerHeap) String() string {
var s string
for i := 0; i < len(c); i++ {
s += c[i].Name + " : " + c[i].Next.String() + "\n"
}
return s
}
func (cjth CronJobTimerHeap) Len() int {
return len(cjth)
}
@@ -117,11 +134,11 @@ func GetCronJobManager(idDbWorker bool) *SCronJobManager {
return manager
}
func (self *SCronJobManager) AddJob1(name string, interval time.Duration, jobFunc TCronJobFunction) {
self.AddJob1WithStartRun(name, interval, jobFunc, false)
func (self *SCronJobManager) AddJobAtIntervals(name string, interval time.Duration, jobFunc TCronJobFunction) {
self.AddJobAtIntervalsWithStartRun(name, interval, jobFunc, false)
}
func (self *SCronJobManager) AddJob1WithStartRun(name string, interval time.Duration, jobFunc TCronJobFunction, startRun bool) {
func (self *SCronJobManager) AddJobAtIntervalsWithStartRun(name string, interval time.Duration, jobFunc TCronJobFunction, startRun bool) {
t := Timer1{
dur: interval,
}
@@ -138,7 +155,7 @@ func (self *SCronJobManager) AddJob1WithStartRun(name string, interval time.Dura
}
}
func (self *SCronJobManager) AddJob2(name string, day, hour, min, sec int, jobFunc TCronJobFunction, startRun bool) {
func (self *SCronJobManager) AddJobEveryFewDays(name string, day, hour, min, sec int, jobFunc TCronJobFunction, startRun bool) {
t := Timer2{
day: day,
hour: hour,
@@ -158,6 +175,25 @@ func (self *SCronJobManager) AddJob2(name string, day, hour, min, sec int, jobFu
}
}
func (self *SCronJobManager) AddJobEveryFewHour(name string, hour, min, sec int, jobFunc TCronJobFunction, startRun bool) {
t := TimerHour{
hour: hour,
min: min,
sec: sec,
}
job := SCronJob{
Name: name,
job: jobFunc,
Timer: &t,
StartRun: startRun,
}
if !self.running {
self.jobs = append(self.jobs, &job)
} else {
self.add <- &job
}
}
func (self *SCronJobManager) Next(now time.Time) {
for _, job := range self.jobs {
job.Next = job.Timer.Next(now)
@@ -197,14 +233,7 @@ func (self *SCronJobManager) run() {
}
select {
case now = <-timer.C:
for i, job := range self.jobs {
if job.Next.After(now) || job.Next.IsZero() {
break
}
job.runJob(false)
job.Next = job.Timer.Next(now)
heap.Fix(&self.jobs, i)
}
self.runJob(now)
case newJob := <-self.add:
now = time.Now()
newJob.Next = newJob.Timer.Next(now)
@@ -216,6 +245,15 @@ 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()) {
self.jobs[0].runJob(false)
self.jobs[0].Next = self.jobs[0].Timer.Next(now)
heap.Fix(&self.jobs, 0)
self.runJob(now)
}
}
func (job *SCronJob) runJob(isStart bool) {
manager.workers.Run(func() {
job.runJobInWorker(isStart)
+77 -28
View File
@@ -1795,61 +1795,107 @@ func (manager *SDiskManager) CleanPendingDeleteDisks(ctx context.Context, userCr
}
}
func (manager *SDiskManager) getAutoSnapshotDisks() []SDisk {
q := manager.Query().SubQuery()
dest := make([]SDisk, 0)
err := q.Query().Filter(sqlchemy.Equals(q.Field("auto_snapshot"), true)).All(&dest)
if err != nil {
return nil
func (manager *SDiskManager) getAutoSnapshotDisksId() ([]SSnapshotPolicyDisk, error) {
t := time.Now()
week := t.Weekday()
if week == 0 { // sunday is zero
week += 7
}
return dest
timePoint := t.Hour()
sps, err := SnapshotPolicyManager.GetSnapshotPoliciesAt(uint32(week), uint32(timePoint))
if err != nil {
return nil, err
}
if len(sps) == 0 {
return nil, nil
}
spds := make([]SSnapshotPolicyDisk, 0)
spdq := SnapshotPolicyDiskManager.Query()
spdq.Filter(sqlchemy.In(spdq.Field("snapshotpolicy_id"), sps))
err = spdq.All(&spds)
if err != nil {
return nil, err
}
return spds, nil
}
func (manager *SDiskManager) AutoDiskSnapshot(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) {
disks := manager.getAutoSnapshotDisks()
if len(disks) == 0 {
spds, err := manager.getAutoSnapshotDisksId()
if err != nil {
log.Errorf("Get auto snapshot disks id failed: %s", err)
return
}
if len(spds) == 0 {
log.Infof("CronJob AutoDiskSnapshot: No disk need create snapshot")
return
}
for i := 0; i < len(disks); i++ {
disks[i].SetModelManager(DiskManager, &disks[i])
now := time.Now()
for i := 0; i < len(spds); i++ {
disk := manager.FetchDiskById(spds[i].DiskId)
snapshotPolicy := SnapshotPolicyManager.FetchSnapshotPolicyById(spds[i].SnapshotpolicyId)
var (
err error
snapCount int
guests = disks[i].GetGuests()
snapshotName = "Auto-" + disks[i].Name + time.Now().Format("2006-01-02#15:04:05")
err error
snapCount int
cleanOverdueSnapshots bool
guests = disk.GetGuests()
snapshotName = "Auto-" + disk.Name + time.Now().Format("2006-01-02#15:04:05")
)
if len(guests) == 1 && !utils.IsInStringArray(guests[0].Status, []string{api.VM_RUNNING, api.VM_READY}) {
if utils.IsInStringArray(disk.GetStorage().StorageType, []string{api.STORAGE_LOCAL, api.STORAGE_GPFS, api.STORAGE_NFS}) &&
len(guests) == 1 && !utils.IsInStringArray(guests[0].Status, []string{api.VM_RUNNING, api.VM_READY}) {
err = fmt.Errorf("Guest(%s) in status(%s) cannot do snapshot action", guests[0].Id, guests[0].Status)
goto onFail
}
if err := disks[i].CreateSnpashotAuto(ctx, userCred, snapshotName); err != nil {
if err := disk.CreateSnpashotAuto(ctx, userCred, snapshotName); err != nil {
err = fmt.Errorf("Create snapshot auto failed %s", err)
goto onFail
}
snapCount, err = SnapshotManager.Query().Equals("fake_deleted", false).
Equals("disk_id", disks[i].Id).Equals("created_by", api.SNAPSHOT_AUTO).CountWithError()
// if auto snapshot count gt max auto snapshot count, do clean overdued snapshots
snapCount, err = SnapshotManager.Query().Equals("fake_deleted", false).Equals("disk_id", disk.Id).
Equals("created_by", api.SNAPSHOT_AUTO).CountWithError()
if err != nil {
err = fmt.Errorf("GetSnapshotCount fail %s", err)
goto onFail
}
cleanOverdueSnapshots = snapCount > (options.Options.DefaultMaxSnapshotCount - options.Options.DefaultMaxManualSnapshotCount)
log.Infof("Auto snapshot count %v, max auto snapshot count %v",
snapCount, options.Options.DefaultMaxSnapshotCount-options.Options.DefaultMaxManualSnapshotCount)
if snapCount > (options.Options.DefaultMaxSnapshotCount - options.Options.DefaultMaxManualSnapshotCount) {
disks[i].CleanOverdueSnapshots(ctx, userCred)
// else if snapshot is overdued, do clean overdued snapshots
if snapshotPolicy.RetentionDays > 0 && !cleanOverdueSnapshots {
t := now.AddDate(0, 0, -1*snapshotPolicy.RetentionDays)
q := SnapshotManager.Query().Equals("fake_deleted", false).Equals("disk_id", disk.Id).
Equals("created_by", api.SNAPSHOT_AUTO).LT("created_at", t)
q.DebugQuery()
snapCount, err = SnapshotManager.Query().Equals("fake_deleted", false).Equals("disk_id", disk.Id).
Equals("created_by", api.SNAPSHOT_AUTO).LT("created_at", t).CountWithError()
if err != nil {
err = fmt.Errorf("GetSnapshotCount fail %s", err)
goto onFail
}
cleanOverdueSnapshots = snapCount > 0
}
db.OpsLog.LogEvent(&disks[i], db.ACT_DISK_AUTO_SNAPSHOT, "disk auto snapshot "+snapshotName, userCred)
if cleanOverdueSnapshots {
disk.CleanOverdueSnapshots(ctx, userCred, snapshotPolicy, now)
}
db.OpsLog.LogEvent(disk, db.ACT_DISK_AUTO_SNAPSHOT, "disk auto snapshot "+snapshotName, userCred)
continue
onFail:
db.OpsLog.LogEvent(&disks[i], db.ACT_DISK_AUTO_SNAPSHOT_FAIL, err.Error(), userCred)
db.OpsLog.LogEvent(disk, db.ACT_DISK_AUTO_SNAPSHOT_FAIL, err.Error(), userCred)
reason := fmt.Sprintf("Disk auto create snapshot failed: %s", err.Error())
notifyclient.NotifySystemError(disks[i].Id, disks[i].Name, db.ACT_DISK_AUTO_SNAPSHOT_FAIL, reason)
notifyclient.NotifySystemError(disk.Id, disk.Name, db.ACT_DISK_AUTO_SNAPSHOT_FAIL, reason)
}
}
func (self *SDisk) CreateSnpashotAuto(ctx context.Context, userCred mcclient.TokenCredential, snapshotName string) error {
// TODO: snapshot quota is not enough, default is 10, or is need check
quotaPlatform := self.GetQuotaPlatformID()
pendingUsage := &SQuota{Snapshot: 1}
_, err := QuotaManager.CheckQuota(ctx, userCred, rbacutils.ScopeProject, self.GetOwnerId(), quotaPlatform, pendingUsage)
if err != nil {
return httperrors.NewOutOfQuotaError("Check set pending quota error %s", err)
}
snap, err := SnapshotManager.CreateSnapshot(ctx, userCred, api.SNAPSHOT_AUTO, self.Id, "", "", snapshotName)
if err != nil {
return err
@@ -1858,8 +1904,11 @@ func (self *SDisk) CreateSnpashotAuto(ctx context.Context, userCred mcclient.Tok
return snap.StartSnapshotCreateTask(ctx, userCred, nil)
}
func (self *SDisk) CleanOverdueSnapshots(ctx context.Context, userCred mcclient.TokenCredential) error {
if task, err := taskman.TaskManager.NewTask(ctx, "DiskCleanOverduedSnapshots", self, userCred, nil, "", "", nil); err != nil {
func (self *SDisk) CleanOverdueSnapshots(ctx context.Context, userCred mcclient.TokenCredential, sp *SSnapshotPolicy, now time.Time) error {
kwargs := jsonutils.NewDict()
kwargs.Set("snapshotpolicy_id", jsonutils.NewString(sp.Id))
kwargs.Set("start_time", jsonutils.NewTimeString(now))
if task, err := taskman.TaskManager.NewTask(ctx, "DiskCleanOverduedSnapshots", self, userCred, kwargs, "", "", nil); err != nil {
log.Errorln(err)
return err
} else {
+1 -1
View File
@@ -79,7 +79,7 @@ type IRegionDriver interface {
ValidateCreateEipData(ctx context.Context, userCred mcclient.TokenCredential, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error)
// Region Driver Snapshot Policy Apis
ValidateCreateSnapshotPolicyData(ctx context.Context, userCred mcclient.TokenCredential, data *compute.SSnapshotPolicyCreateInput) error
ValidateCreateSnapshotPolicyData(context.Context, mcclient.TokenCredential, *compute.SSnapshotPolicyCreateInput, mcclient.IIdentityProvider, *jsonutils.JSONDict) error
RequestCreateSnapshotPolicy(ctx context.Context, userCred mcclient.TokenCredential, sp *SSnapshotPolicy, task taskman.ITask) error
RequestDeleteSnapshotPolicy(ctx context.Context, userCred mcclient.TokenCredential, sp *SSnapshotPolicy, task taskman.ITask) error
+40 -10
View File
@@ -17,10 +17,12 @@ package models
import (
"context"
"fmt"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/util/compare"
"yunion.io/x/pkg/utils"
"yunion.io/x/sqlchemy"
api "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
@@ -46,9 +48,11 @@ type SSnapshotPolicy struct {
RetentionDays int `nullable:"false" list:"user" get:"user" create:"required"`
RepeatWeekdays uint8 `charset:"utf8" create:"required"`
TimePoints uint32 `charset:"utf8" create:"required"`
IsActivated bool `list:"user" get:"user" create:"optional" default:"true"`
// 0~6, 0 is Monday
RepeatWeekdays uint8 `charset:"utf8" create:"required"`
// 0~23
TimePoints uint32 `charset:"utf8" create:"required"`
IsActivated bool `list:"user" get:"user" create:"optional" default:"true"`
}
var SnapshotPolicyManager *SSnapshotPolicyManager
@@ -79,12 +83,6 @@ func (manager *SSnapshotPolicyManager) ValidateCreateData(ctx context.Context, u
return nil, err
}
managerIdV := validators.NewModelIdOrNameValidator("manager", "cloudprovider", nil)
if err := managerIdV.Validate(data); err != nil {
return nil, err
}
input.ManagerId, _ = data.GetString("manager_id")
cloudregionV := validators.NewModelIdOrNameValidator("cloudregion", "cloudregion", ownerId)
err = cloudregionV.Validate(data)
if err != nil {
@@ -93,7 +91,7 @@ func (manager *SSnapshotPolicyManager) ValidateCreateData(ctx context.Context, u
cloudregion := cloudregionV.Model.(*SCloudregion)
input.CloudregionId = cloudregion.GetId()
err = cloudregion.GetDriver().ValidateCreateSnapshotPolicyData(ctx, userCred, input)
err = cloudregion.GetDriver().ValidateCreateSnapshotPolicyData(ctx, userCred, input, ownerId, data)
if err != nil {
return nil, err
}
@@ -417,3 +415,35 @@ func (self *SSnapshotPolicy) preCheck(
}
return diskIds, nil
}
func (manager *SSnapshotPolicyManager) GetSnapshotPoliciesAt(week, timePoint uint32) ([]string, error) {
q := manager.Query("id")
q = q.Filter(sqlchemy.Equals(sqlchemy.AND_Val("", q.Field("repeat_weekdays"), 1<<week), 1<<week))
q = q.Filter(sqlchemy.Equals(sqlchemy.AND_Val("", q.Field("time_points"), 1<<timePoint), 1<<timePoint))
q = q.Equals("is_activated", true)
q.DebugQuery()
sps := make([]SSnapshotPolicy, 0)
err := q.All(&sps)
if err != nil {
return nil, err
}
if len(sps) > 0 {
ret := make([]string, len(sps))
for i := 0; i < len(sps); i++ {
ret[i] = sps[i].Id
}
return ret, nil
}
return nil, nil
}
func (manager *SSnapshotPolicyManager) FetchSnapshotPolicyById(spId string) *SSnapshotPolicy {
sp, err := manager.FetchById(spId)
if err != nil {
log.Errorf("FetchBId fail %s", err)
return nil
}
return sp.(*SSnapshotPolicy)
}
+3 -9
View File
@@ -16,7 +16,6 @@ package models
import (
"context"
"yunion.io/x/onecloud/pkg/cloudcommon/validators"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
@@ -70,14 +69,9 @@ func (self *SSnapshotPolicyDisk) Detach(ctx context.Context, userCred mcclient.T
}
func (self *SSnapshotPolicyDiskManager) ValidateCreateData(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, data *jsonutils.JSONDict) (*jsonutils.JSONDict, error) {
cloudregionV := validators.NewModelIdOrNameValidator("cloudregion", "cloudregion", ownerId)
err := cloudregionV.Validate(data)
if err != nil {
return nil, err
}
cloudregion := cloudregionV.Model.(*SCloudregion)
diskId, _ := data.GetString(self.GetMasterManager().Keyword())
err = cloudregion.GetDriver().ValidateCreateSnapshopolicyDiskData(ctx, userCred, diskId)
diskId, _ := data.GetString(self.GetMasterFieldName())
disk := DiskManager.FetchDiskById(diskId)
err := disk.GetStorage().GetRegion().GetDriver().ValidateCreateSnapshopolicyDiskData(ctx, userCred, diskId)
if err != nil {
return nil, err
}
+1 -1
View File
@@ -534,7 +534,7 @@ func (self *SSnapshotManager) PerformDeleteDiskSnapshots(ctx context.Context, us
return nil, httperrors.NewNotFoundError("Disk %s dose not have snapshot", diskId)
}
for i := 0; i < len(snapshots); i++ {
if snapshots[i].CreatedBy == api.SNAPSHOT_MANUAL && snapshots[i].FakeDeleted == false {
if snapshots[i].FakeDeleted == false {
return nil, httperrors.NewBadRequestError("Can not delete disk snapshots, have manual snapshot")
}
}
-12
View File
@@ -26,7 +26,6 @@ import (
"yunion.io/x/onecloud/pkg/util/rand"
"yunion.io/x/pkg/utils"
"yunion.io/x/onecloud/pkg/apis/compute"
api "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
"yunion.io/x/onecloud/pkg/cloudcommon/validators"
@@ -834,17 +833,6 @@ func (self *SAliyunRegionDriver) ValidateCreateSnapshopolicyDiskData(ctx context
return nil
}
func (self *SAliyunRegionDriver) ValidateCreateSnapshotPolicyData(ctx context.Context, userCred mcclient.TokenCredential, data *compute.SSnapshotPolicyCreateInput) error {
err := self.SManagedVirtualizationRegionDriver.ValidateCreateSnapshotPolicyData(ctx, userCred, data)
if err != nil {
return err
}
if data.RetentionDays < -1 || data.RetentionDays == 0 || data.RetentionDays > 65535 {
return httperrors.NewInputParameterError("Retention days must in 1~65535 or -1")
}
return nil
}
func (self *SAliyunRegionDriver) ValidateSnapshotCreate(ctx context.Context, userCred mcclient.TokenCredential, disk *models.SDisk, data *jsonutils.JSONDict) error {
name, _ := data.GetString("name")
if strings.HasPrefix(name, "auto") || strings.HasPrefix(name, "http://") || strings.HasPrefix(name, "https://") {
+21 -3
View File
@@ -20,10 +20,11 @@ import (
"yunion.io/x/jsonutils"
"yunion.io/x/onecloud/pkg/apis/compute"
api "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/cloudcommon/db/taskman"
"yunion.io/x/onecloud/pkg/cloudprovider"
"yunion.io/x/onecloud/pkg/compute/models"
"yunion.io/x/onecloud/pkg/httperrors"
"yunion.io/x/onecloud/pkg/mcclient"
)
@@ -126,8 +127,25 @@ func (self *SBaseRegionDriver) RequestDeleteLoadbalancerListenerRule(ctx context
return fmt.Errorf("Not Implement RequestDeleteLoadbalancerListenerRule")
}
func (self *SBaseRegionDriver) ValidateCreateSnapshotPolicyData(ctx context.Context, userCred mcclient.TokenCredential, data *compute.SSnapshotPolicyCreateInput) error {
return fmt.Errorf("Not Implement ValidateCreateSnapshotPolicyData")
func (self *SBaseRegionDriver) ValidateCreateSnapshotPolicyData(ctx context.Context, userCred mcclient.TokenCredential, input *api.SSnapshotPolicyCreateInput, ownerId mcclient.IIdentityProvider, data *jsonutils.JSONDict) error {
var err error
if len(input.RepeatWeekdays) == 0 {
return httperrors.NewMissingParameterError("repeat_weekdays")
}
input.RepeatWeekdays, err = daysValidate(input.RepeatWeekdays, 1, 7)
if err != nil {
return httperrors.NewInputParameterError(err.Error())
}
if len(input.TimePoints) == 0 {
return httperrors.NewInputParameterError("time_points")
}
input.TimePoints, err = daysValidate(input.TimePoints, 0, 23)
if err != nil {
return httperrors.NewInputParameterError(err.Error())
}
return nil
}
func (self *SBaseRegionDriver) RequestCreateSnapshotPolicy(ctx context.Context, userCred mcclient.TokenCredential, sp *models.SSnapshotPolicy, task taskman.ITask) error {
+35
View File
@@ -773,3 +773,38 @@ func (self *SKVMRegionDriver) OnDiskReset(ctx context.Context, userCred mcclient
storage := disk.GetStorage()
return models.GetStorageDriver(storage.StorageType).OnDiskReset(ctx, userCred, disk, snapshot, data)
}
func (self *SKVMRegionDriver) ValidateCreateSnapshotPolicyData(ctx context.Context, userCred mcclient.TokenCredential, input *api.SSnapshotPolicyCreateInput, ownerId mcclient.IIdentityProvider, data *jsonutils.JSONDict) error {
err := self.SBaseRegionDriver.ValidateCreateSnapshotPolicyData(ctx, userCred, input, ownerId, data)
if err != nil {
return err
}
// TODO: kvm retention days
if input.RetentionDays < -1 || input.RetentionDays == 0 || input.RetentionDays > 10 {
return httperrors.NewInputParameterError("Retention days must in 1~10 or -1")
}
return nil
}
func (self *SKVMRegionDriver) RequestCreateSnapshotPolicy(ctx context.Context, userCred mcclient.TokenCredential, sp *models.SSnapshotPolicy, task taskman.ITask) error {
taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) {
return nil, nil
})
return nil
}
func (self *SKVMRegionDriver) ValidateCreateSnapshopolicyDiskData(ctx context.Context, userCred mcclient.TokenCredential, diskID string) error {
return nil
}
func (self *SKVMRegionDriver) RequestApplySnapshotPolicy(ctx context.Context, userCred mcclient.TokenCredential, sp *models.SSnapshotPolicy, task taskman.ITask, diskId string) error {
task.ScheduleRun(nil)
return nil
}
func (self *SKVMRegionDriver) RequestCancelSnapshotPolicy(ctx context.Context, userCred mcclient.TokenCredential, sp *models.SSnapshotPolicy, task taskman.ITask, diskId string) error {
taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) {
return nil, nil
})
return nil
}
+13 -22
View File
@@ -25,11 +25,11 @@ import (
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/utils"
"yunion.io/x/onecloud/pkg/apis/compute"
api "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
"yunion.io/x/onecloud/pkg/cloudcommon/db/lockman"
"yunion.io/x/onecloud/pkg/cloudcommon/db/taskman"
"yunion.io/x/onecloud/pkg/cloudcommon/validators"
"yunion.io/x/onecloud/pkg/cloudprovider"
"yunion.io/x/onecloud/pkg/compute/models"
"yunion.io/x/onecloud/pkg/httperrors"
@@ -1207,27 +1207,18 @@ func (self *SManagedVirtualizationRegionDriver) OnDiskReset(ctx context.Context,
return nil
}
func (self *SManagedVirtualizationRegionDriver) ValidateCreateSnapshotPolicyData(ctx context.Context, userCred mcclient.TokenCredential, data *compute.SSnapshotPolicyCreateInput) error {
var err error
if len(data.RepeatWeekdays) == 0 {
return httperrors.NewMissingParameterError("repeat_weekdays")
}
data.RepeatWeekdays, err = daysValidate(data.RepeatWeekdays, 1, 7)
if err != nil {
return httperrors.NewInputParameterError(err.Error())
}
if len(data.TimePoints) == 0 {
return httperrors.NewInputParameterError("time_points")
}
data.TimePoints, err = daysValidate(data.TimePoints, 0, 23)
if err != nil {
return httperrors.NewInputParameterError(err.Error())
}
return nil
}
func (self *SManagedVirtualizationRegionDriver) ValidateCreateSnapshopolicyDiskData(ctx context.Context, userCred mcclient.TokenCredential, diskID string) error {
return nil
}
func (self *SManagedVirtualizationRegionDriver) ValidateCreateSnapshotPolicyData(ctx context.Context, userCred mcclient.TokenCredential, input *api.SSnapshotPolicyCreateInput, ownerId mcclient.IIdentityProvider, data *jsonutils.JSONDict) error {
cloudregionV := validators.NewModelIdOrNameValidator("cloudregion", "cloudregion", ownerId)
err := cloudregionV.Validate(data)
if err != nil {
return err
}
cloudregion := cloudregionV.Model.(*models.SCloudregion)
input.CloudregionId = cloudregion.GetId()
return nil
}
-12
View File
@@ -21,7 +21,6 @@ import (
"yunion.io/x/jsonutils"
"yunion.io/x/onecloud/pkg/apis/compute"
api "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
"yunion.io/x/onecloud/pkg/cloudcommon/db/taskman"
@@ -744,14 +743,3 @@ func (self *SQcloudRegionDriver) ValidateCreateLoadbalancerBackendData(ctx conte
data.Set("cloudregion_id", jsonutils.NewString(lb.CloudregionId))
return data, nil
}
func (self *SQcloudRegionDriver) ValidateCreateSnapshotPolicyData(ctx context.Context, userCred mcclient.TokenCredential, data *compute.SSnapshotPolicyCreateInput) error {
err := self.SManagedVirtualizationRegionDriver.ValidateCreateSnapshotPolicyData(ctx, userCred, data)
if err != nil {
return err
}
if data.RetentionDays < -1 || data.RetentionDays == 0 || data.RetentionDays > 65535 {
return httperrors.NewInputParameterError("Retention days must in 1~65535 or -1")
}
return nil
}
+10 -10
View File
@@ -80,21 +80,21 @@ func StartService() {
if !opts.IsSlaveNode {
cron := cronman.GetCronJobManager(true)
cron.AddJob1("CleanPendingDeleteServers", time.Duration(opts.PendingDeleteCheckSeconds)*time.Second, models.GuestManager.CleanPendingDeleteServers)
cron.AddJob1("CleanPendingDeleteDisks", time.Duration(opts.PendingDeleteCheckSeconds)*time.Second, models.DiskManager.CleanPendingDeleteDisks)
cron.AddJob1("CleanPendingDeleteLoadbalancers", time.Duration(opts.LoadbalancerPendingDeleteCheckInterval)*time.Second, models.LoadbalancerAgentManager.CleanPendingDeleteLoadbalancers)
cron.AddJobAtIntervals("CleanPendingDeleteServers", time.Duration(opts.PendingDeleteCheckSeconds)*time.Second, models.GuestManager.CleanPendingDeleteServers)
cron.AddJobAtIntervals("CleanPendingDeleteDisks", time.Duration(opts.PendingDeleteCheckSeconds)*time.Second, models.DiskManager.CleanPendingDeleteDisks)
cron.AddJobAtIntervals("CleanPendingDeleteLoadbalancers", time.Duration(opts.LoadbalancerPendingDeleteCheckInterval)*time.Second, models.LoadbalancerAgentManager.CleanPendingDeleteLoadbalancers)
if opts.PrepaidExpireCheck {
cron.AddJob1("CleanExpiredPrepaidServers", time.Duration(opts.PrepaidExpireCheckSeconds)*time.Second, models.GuestManager.DeleteExpiredPrepaidServers)
cron.AddJobAtIntervals("CleanExpiredPrepaidServers", time.Duration(opts.PrepaidExpireCheckSeconds)*time.Second, models.GuestManager.DeleteExpiredPrepaidServers)
}
cron.AddJob1("StartHostPingDetectionTask", time.Duration(opts.HostOfflineDetectionInterval)*time.Second, models.HostManager.PingDetectionTask)
cron.AddJobAtIntervals("StartHostPingDetectionTask", time.Duration(opts.HostOfflineDetectionInterval)*time.Second, models.HostManager.PingDetectionTask)
cron.AddJob1WithStartRun("CalculateQuotaUsages", time.Duration(opts.CalculateQuotaUsageIntervalSeconds)*time.Second, models.QuotaManager.CalculateQuotaUsages, true)
cron.AddJobAtIntervalsWithStartRun("CalculateQuotaUsages", time.Duration(opts.CalculateQuotaUsageIntervalSeconds)*time.Second, models.QuotaManager.CalculateQuotaUsages, true)
cron.AddJob1WithStartRun("AutoSyncCloudaccountTask", time.Duration(opts.CloudAutoSyncIntervalSeconds)*time.Second, models.CloudaccountManager.AutoSyncCloudaccountTask, true)
cron.AddJobAtIntervalsWithStartRun("AutoSyncCloudaccountTask", time.Duration(opts.CloudAutoSyncIntervalSeconds)*time.Second, models.CloudaccountManager.AutoSyncCloudaccountTask, true)
cron.AddJob2("AutoDiskSnapshot", opts.AutoSnapshotDay, opts.AutoSnapshotHour, 0, 0, models.DiskManager.AutoDiskSnapshot, false)
cron.AddJob2("SyncSkus", opts.SyncSkusDay, opts.SyncSkusHour, 0, 0, models.SyncSkus, true)
cron.AddJob2("StorageSnapshotsRecycle", 1, 2, 0, 0, models.StorageManager.StorageSnapshotsRecycle, false)
cron.AddJobEveryFewDays("AutoDiskSnapshot", opts.AutoSnapshotDay, opts.AutoSnapshotHour, 0, 0, models.DiskManager.AutoDiskSnapshot, false)
cron.AddJobEveryFewDays("SyncSkus", opts.SyncSkusDay, opts.SyncSkusHour, 0, 0, models.SyncSkus, true)
cron.AddJobEveryFewDays("StorageSnapshotsRecycle", 1, 2, 0, 0, models.StorageManager.StorageSnapshotsRecycle, false)
cron.Start()
defer cron.Stop()
@@ -16,6 +16,7 @@ package tasks
import (
"context"
"fmt"
"yunion.io/x/jsonutils"
@@ -36,15 +37,44 @@ func init() {
func (self *DiskCleanOverduedSnapshots) OnInit(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
disk := obj.(*models.SDisk)
count, err := models.SnapshotManager.Query().Equals("disk_id", disk.Id).
Equals("created_by", compute.SNAPSHOT_AUTO).Equals("fake_deleted", false).CountWithError()
if err != nil {
self.SetStageFailed(ctx, err.Error())
spId, _ := self.Params.GetString("snapshotpolicy_id")
sp := models.SnapshotPolicyManager.FetchSnapshotPolicyById(spId)
if sp == nil {
self.SetStageFailed(ctx, "missing snapshot policy ???")
return
}
if count <= (options.Options.DefaultMaxSnapshotCount - options.Options.DefaultMaxManualSnapshotCount) {
now, err := self.Params.GetTime("start_time")
if err != nil {
self.SetStageFailed(ctx, "failed to get start time")
return
}
var (
snapCount int
cleanOverdueSnapshots bool
)
snapCount, err = models.SnapshotManager.Query().Equals("fake_deleted", false).Equals("disk_id", disk.Id).
Equals("created_by", compute.SNAPSHOT_AUTO).CountWithError()
if err != nil {
err = fmt.Errorf("GetSnapshotCount fail %s", err)
return
}
cleanOverdueSnapshots = snapCount > (options.Options.DefaultMaxSnapshotCount - options.Options.DefaultMaxManualSnapshotCount)
if sp.RetentionDays > 0 && !cleanOverdueSnapshots {
t := now.AddDate(0, 0, -1*sp.RetentionDays)
snapCount, err = models.SnapshotManager.Query().Equals("fake_deleted", false).Equals("disk_id", disk.Id).
Equals("created_by", compute.SNAPSHOT_AUTO).LT("created_at", t).CountWithError()
if err != nil {
self.SetStageFailed(ctx, err.Error())
return
}
cleanOverdueSnapshots = snapCount > 0
}
if !cleanOverdueSnapshots {
self.SetStageComplete(ctx, nil)
return
}
+15 -23
View File
@@ -83,9 +83,9 @@ type SnapshotPolicyApplyTask struct {
}
func (self *SnapshotPolicyApplyTask) taskFail(ctx context.Context, disk *models.SDisk, snapshotPolicyId, reason string) {
jointModel, err := db.FetchJointByIds(models.SnapshotPolicyDiskManager, self.Id, snapshotPolicyId, jsonutils.JSONNull)
jointModel, err := db.FetchJointByIds(models.SnapshotPolicyDiskManager, disk.Id, snapshotPolicyId, jsonutils.JSONNull)
if err != nil {
log.Errorf("Fetch SnapshotPolicy %s Disk %s joint model failed, need to delete", self.Id, snapshotPolicyId)
log.Errorf("Fetch SnapshotPolicy %s Disk %s joint model failed %s", disk.Id, snapshotPolicyId, reason)
return
}
snapshotPolicyDisk := jointModel.(*models.SSnapshotPolicyDisk)
@@ -104,29 +104,26 @@ func (self *SnapshotPolicyApplyTask) OnInit(ctx context.Context, obj db.IStandal
disk := obj.(*models.SDisk)
snapshotPolicyID, _ := self.Params.GetString("snapshot_policy_id")
iregion, err := disk.GetIRegion()
if err != nil {
self.taskFail(ctx, disk, snapshotPolicyID, fmt.Sprintf("failed to find iregion for snapshot policy %s: %s", disk.Id, err.Error()))
return
}
// fetch disk model by diksID
model, err := models.SnapshotPolicyManager.FetchById(snapshotPolicyID)
if err != nil {
self.taskFail(ctx, disk, snapshotPolicyID, fmt.Sprintf("failed to fetch disk by id %s: %s", snapshotPolicyID, err.Error()))
return
}
snapshotPolicy := model.(*models.SSnapshotPolicy)
self.SetStage("OnSnapshotPolicyApply", nil)
if err := iregion.ApplySnapshotPolicyToDisks(snapshotPolicy.ExternalId, disk.ExternalId); err != nil {
if err := disk.GetStorage().GetRegion().GetDriver().
RequestApplySnapshotPolicy(ctx, self.UserCred, snapshotPolicy, self, disk.ExternalId); err != nil {
self.taskFail(ctx, disk, snapshotPolicyID, fmt.Sprintf("faile to attach snapshot policy %s and disk %s: %s", snapshotPolicy.Id, disk.Id, err.Error()))
}
}
func (self *SnapshotPolicyApplyTask) OnSnapshotPolicyApply(ctx context.Context, disk *models.SDisk, data jsonutils.JSONObject) {
db.OpsLog.LogEvent(disk, db.ACT_APPLY_SNAPSHOT_POLICY, "", self.UserCred)
logclient.AddActionLogWithStartable(self, disk, logclient.ACT_APPLY_SNAPSHOT_POLICY, "", self.UserCred, true)
self.SetStageComplete(ctx, nil)
}
// -=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-=-
@@ -136,9 +133,9 @@ type SnapshotPolicyCancelTask struct {
}
func (self *SnapshotPolicyCancelTask) taskFail(ctx context.Context, disk *models.SDisk, snapshotPolicyId, reason string) {
jointModel, err := db.FetchJointByIds(models.SnapshotPolicyDiskManager, self.Id, snapshotPolicyId, jsonutils.JSONNull)
jointModel, err := db.FetchJointByIds(models.SnapshotPolicyDiskManager, disk.Id, snapshotPolicyId, jsonutils.JSONNull)
if err != nil {
log.Errorf("Fetch SnapshotPolicy %s Disk %s joint model failed, need to mark undelete", self.Id, snapshotPolicyId)
log.Errorf("Fetch SnapshotPolicy %s Disk %s joint model failed %s", disk.Id, snapshotPolicyId, reason)
return
}
snapshotPolicyDisk := jointModel.(*models.SSnapshotPolicyDisk)
@@ -156,25 +153,20 @@ func (self *SnapshotPolicyCancelTask) OnInit(ctx context.Context, obj db.IStanda
disk := obj.(*models.SDisk)
snapshotPolicyID, _ := self.GetParams().GetString("snapshot_policy_id")
// get region
iregion, err := disk.GetIRegion()
if err != nil {
self.taskFail(ctx, disk, snapshotPolicyID, fmt.Sprintf("failed to find iregion for disk %s: %s", disk.Name, err.Error()))
return
}
model, err := models.SnapshotPolicyManager.FetchById(snapshotPolicyID)
if err != nil {
self.taskFail(ctx, disk, snapshotPolicyID, fmt.Sprintf("failed to fetch disk by id %s: %s", snapshotPolicyID, err.Error()))
return
}
snapshotPolicy := model.(*models.SSnapshotPolicy)
self.SetStage("OnSnapshotPolicyApply", nil)
if err := iregion.CancelSnapshotPolicyToDisks(snapshotPolicy.ExternalId, disk.ExternalId); err != nil {
self.SetStage("OnSnapshotPolicyCancel", nil)
if err := disk.GetStorage().GetRegion().GetDriver().
RequestCancelSnapshotPolicy(ctx, self.UserCred, snapshotPolicy, self, disk.ExternalId); err != nil {
self.taskFail(ctx, disk, snapshotPolicyID, fmt.Sprintf("faile to detach snapshot policy %s and disk %s: %s", snapshotPolicy.Id, disk.Id, err.Error()))
}
}
func (self *SnapshotPolicyCancelTask) OnSnapshotPolicyCancel(ctx context.Context, disk *models.SDisk, data jsonutils.JSONObject) {
db.OpsLog.LogEvent(disk, db.ACT_CANCEL_SNAPSHOT_POLICY, "", self.UserCred)
logclient.AddActionLogWithStartable(self, disk, logclient.ACT_CANCEL_SNAPSHOT_POLICY, "", self.UserCred, true)
self.SetStageComplete(ctx, nil)
+1 -1
View File
@@ -100,7 +100,7 @@ func (host *SHostService) RunService() {
options.HostOptions.Address, options.HostOptions.Port+1000)
cronManager := cronman.GetCronJobManager(false)
cronManager.AddJob2(
cronManager.AddJobEveryFewDays(
"CleanRecycleDiskFiles", 1, 3, 0, 0, storageman.CleanRecycleDiskfiles, false)
cronManager.Start()
+1 -1
View File
@@ -315,7 +315,7 @@ func StartSnapshotRecycle(storage IStorage) {
if !fileutils2.Exists(storage.GetSnapshotDir()) {
procutils.NewCommand("mkdir", "-p", storage.GetSnapshotDir()).Run()
}
cronman.GetCronJobManager(false).AddJob1(
cronman.GetCronJobManager(false).AddJobAtIntervals(
"SnapshotRecycle", time.Hour*6,
func(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) {
snapshotRecycle(ctx, userCred, isStart, storage)
+2 -2
View File
@@ -101,8 +101,8 @@ func StartService() {
if !opts.IsSlaveNode {
cron := cronman.GetCronJobManager(true)
cron.AddJob1("CleanPendingDeleteImages", time.Duration(options.Options.PendingDeleteCheckSeconds)*time.Second, models.ImageManager.CleanPendingDeleteImages)
cron.AddJob1("CalculateQuotaUsages", time.Duration(opts.CalculateQuotaUsageIntervalSeconds)*time.Second, models.QuotaManager.CalculateQuotaUsages)
cron.AddJobAtIntervals("CleanPendingDeleteImages", time.Duration(options.Options.PendingDeleteCheckSeconds)*time.Second, models.ImageManager.CleanPendingDeleteImages)
cron.AddJobAtIntervals("CalculateQuotaUsages", time.Duration(opts.CalculateQuotaUsageIntervalSeconds)*time.Second, models.QuotaManager.CalculateQuotaUsages)
cron.Start()
}
+2 -2
View File
@@ -80,8 +80,8 @@ func StartService() {
if !opts.IsSlaveNode {
cron := cronman.GetCronJobManager(true)
cron.AddJob1WithStartRun("AutoSyncIdentityProviderTask", time.Duration(opts.AutoSyncIntervalSeconds)*time.Second, models.AutoSyncIdentityProviderTask, true)
cron.AddJob1("FetchProjectResourceCount", time.Duration(opts.FetchProjectResourceCountIntervalSeconds)*time.Second, cronjobs.FetchProjectResourceCount)
cron.AddJobAtIntervalsWithStartRun("AutoSyncIdentityProviderTask", time.Duration(opts.AutoSyncIntervalSeconds)*time.Second, models.AutoSyncIdentityProviderTask, true)
cron.AddJobAtIntervals("FetchProjectResourceCount", time.Duration(opts.FetchProjectResourceCountIntervalSeconds)*time.Second, cronjobs.FetchProjectResourceCount)
cron.Start()
defer cron.Stop()
+2 -2
View File
@@ -59,8 +59,8 @@ func StartService() {
/*if !opts.IsSlaveNode {
cron := cronman.GetCronJobManager(true)
cron.AddJob1("CleanPendingDeleteImages", time.Duration(options.Options.PendingDeleteCheckSeconds)*time.Second, models.ImageManager.CleanPendingDeleteImages)
cron.AddJob1("CalculateQuotaUsages", time.Duration(opts.CalculateQuotaUsageIntervalSeconds)*time.Second, models.QuotaManager.CalculateQuotaUsages)
cron.AddJobAtIntervals("CleanPendingDeleteImages", time.Duration(options.Options.PendingDeleteCheckSeconds)*time.Second, models.ImageManager.CleanPendingDeleteImages)
cron.AddJobAtIntervals("CalculateQuotaUsages", time.Duration(opts.CalculateQuotaUsageIntervalSeconds)*time.Second, models.QuotaManager.CalculateQuotaUsages)
cron.Start()
}*/