From 199b4723979fe361dcd4ee5a29d3aa1ecce2448e Mon Sep 17 00:00:00 2001 From: wanyaoqi Date: Tue, 13 Aug 2019 21:13:49 +0800 Subject: [PATCH] - snapshot support custom snapshot policy --- pkg/cloudcommon/cronman/cronman.go | 62 +++++++++-- pkg/compute/models/disks.go | 105 +++++++++++++----- pkg/compute/models/regiondrivers.go | 2 +- pkg/compute/models/snapshotpolicy.go | 50 +++++++-- pkg/compute/models/snapshotpolicydisks.go | 12 +- pkg/compute/models/snapshots.go | 2 +- pkg/compute/regiondrivers/aliyun.go | 12 -- pkg/compute/regiondrivers/base.go | 24 +++- pkg/compute/regiondrivers/kvm.go | 35 ++++++ pkg/compute/regiondrivers/managedvirtual.go | 35 +++--- pkg/compute/regiondrivers/qcloud.go | 12 -- pkg/compute/service/service.go | 20 ++-- .../tasks/disk_clean_overdued_snapshots.go | 42 ++++++- .../tasks/snapshotpolicy_create_task.go | 38 +++---- pkg/hostman/host_services.go | 2 +- pkg/hostman/storageman/storage_base.go | 2 +- pkg/image/service/service.go | 4 +- pkg/keystone/service/service.go | 4 +- pkg/s3gateway/service/service.go | 4 +- 19 files changed, 310 insertions(+), 157 deletions(-) diff --git a/pkg/cloudcommon/cronman/cronman.go b/pkg/cloudcommon/cronman/cronman.go index c8d7bab81a..9542d5342f 100644 --- a/pkg/cloudcommon/cronman/cronman.go +++ b/pkg/cloudcommon/cronman/cronman.go @@ -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) diff --git a/pkg/compute/models/disks.go b/pkg/compute/models/disks.go index a441a4c6b2..3bd4cec76b 100644 --- a/pkg/compute/models/disks.go +++ b/pkg/compute/models/disks.go @@ -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 { diff --git a/pkg/compute/models/regiondrivers.go b/pkg/compute/models/regiondrivers.go index ab63bdb507..d0839b2cd0 100644 --- a/pkg/compute/models/regiondrivers.go +++ b/pkg/compute/models/regiondrivers.go @@ -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 diff --git a/pkg/compute/models/snapshotpolicy.go b/pkg/compute/models/snapshotpolicy.go index 94250adbea..1ac6d48948 100644 --- a/pkg/compute/models/snapshotpolicy.go +++ b/pkg/compute/models/snapshotpolicy.go @@ -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< 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) +} diff --git a/pkg/compute/models/snapshotpolicydisks.go b/pkg/compute/models/snapshotpolicydisks.go index c8cfee5809..9f7ac2933f 100644 --- a/pkg/compute/models/snapshotpolicydisks.go +++ b/pkg/compute/models/snapshotpolicydisks.go @@ -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 } diff --git a/pkg/compute/models/snapshots.go b/pkg/compute/models/snapshots.go index f8c728afda..e13b9f4b8f 100644 --- a/pkg/compute/models/snapshots.go +++ b/pkg/compute/models/snapshots.go @@ -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") } } diff --git a/pkg/compute/regiondrivers/aliyun.go b/pkg/compute/regiondrivers/aliyun.go index b5143c22d3..fa9442eb35 100644 --- a/pkg/compute/regiondrivers/aliyun.go +++ b/pkg/compute/regiondrivers/aliyun.go @@ -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://") { diff --git a/pkg/compute/regiondrivers/base.go b/pkg/compute/regiondrivers/base.go index b4f12f3396..daeb1cd420 100644 --- a/pkg/compute/regiondrivers/base.go +++ b/pkg/compute/regiondrivers/base.go @@ -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 { diff --git a/pkg/compute/regiondrivers/kvm.go b/pkg/compute/regiondrivers/kvm.go index 30e803ed45..a50e3f419d 100644 --- a/pkg/compute/regiondrivers/kvm.go +++ b/pkg/compute/regiondrivers/kvm.go @@ -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 +} diff --git a/pkg/compute/regiondrivers/managedvirtual.go b/pkg/compute/regiondrivers/managedvirtual.go index 2409a04995..f3755763bb 100644 --- a/pkg/compute/regiondrivers/managedvirtual.go +++ b/pkg/compute/regiondrivers/managedvirtual.go @@ -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" @@ -1202,27 +1202,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 +} diff --git a/pkg/compute/regiondrivers/qcloud.go b/pkg/compute/regiondrivers/qcloud.go index 02e0d04e8a..1e1b98ec10 100644 --- a/pkg/compute/regiondrivers/qcloud.go +++ b/pkg/compute/regiondrivers/qcloud.go @@ -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 -} diff --git a/pkg/compute/service/service.go b/pkg/compute/service/service.go index 942c78b45e..39a115d51b 100644 --- a/pkg/compute/service/service.go +++ b/pkg/compute/service/service.go @@ -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() diff --git a/pkg/compute/tasks/disk_clean_overdued_snapshots.go b/pkg/compute/tasks/disk_clean_overdued_snapshots.go index 087451bfff..dc75d5c88d 100644 --- a/pkg/compute/tasks/disk_clean_overdued_snapshots.go +++ b/pkg/compute/tasks/disk_clean_overdued_snapshots.go @@ -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 } diff --git a/pkg/compute/tasks/snapshotpolicy_create_task.go b/pkg/compute/tasks/snapshotpolicy_create_task.go index 3e82c1a4ce..39ca18b76c 100644 --- a/pkg/compute/tasks/snapshotpolicy_create_task.go +++ b/pkg/compute/tasks/snapshotpolicy_create_task.go @@ -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) diff --git a/pkg/hostman/host_services.go b/pkg/hostman/host_services.go index 8292be6ad6..c006fd3f37 100644 --- a/pkg/hostman/host_services.go +++ b/pkg/hostman/host_services.go @@ -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() diff --git a/pkg/hostman/storageman/storage_base.go b/pkg/hostman/storageman/storage_base.go index 9e0bc7fc05..f536f0308b 100644 --- a/pkg/hostman/storageman/storage_base.go +++ b/pkg/hostman/storageman/storage_base.go @@ -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) diff --git a/pkg/image/service/service.go b/pkg/image/service/service.go index af73cb75bd..0d7905bee2 100644 --- a/pkg/image/service/service.go +++ b/pkg/image/service/service.go @@ -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() } diff --git a/pkg/keystone/service/service.go b/pkg/keystone/service/service.go index 7aa6dc996e..4e88b21051 100644 --- a/pkg/keystone/service/service.go +++ b/pkg/keystone/service/service.go @@ -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() diff --git a/pkg/s3gateway/service/service.go b/pkg/s3gateway/service/service.go index 67e6e36fc8..9f34e8f194 100644 --- a/pkg/s3gateway/service/service.go +++ b/pkg/s3gateway/service/service.go @@ -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() }*/