From 5e82a7e506632007a2adbfac7141278c7c66abe8 Mon Sep 17 00:00:00 2001 From: rainzm Date: Mon, 12 Oct 2020 16:10:26 +0800 Subject: [PATCH] feat: add a retry mechanism to AutoSyncExtDiskSnapshot --- pkg/compute/models/disks.go | 40 +++++++++++++---- pkg/compute/models/snapshotpolicy.go | 54 ++++++++++++++++++++++- pkg/compute/models/snapshotpolicydisks.go | 44 ++++++++++++++++-- pkg/compute/options/options.go | 5 ++- pkg/compute/service/service.go | 3 +- 5 files changed, 130 insertions(+), 16 deletions(-) diff --git a/pkg/compute/models/disks.go b/pkg/compute/models/disks.go index 9d4eb8644f..7eb87551f5 100644 --- a/pkg/compute/models/disks.go +++ b/pkg/compute/models/disks.go @@ -28,6 +28,7 @@ import ( "yunion.io/x/pkg/tristate" "yunion.io/x/pkg/util/compare" "yunion.io/x/pkg/util/fileutils" + "yunion.io/x/pkg/util/sets" "yunion.io/x/pkg/utils" "yunion.io/x/sqlchemy" @@ -2454,28 +2455,49 @@ func (self *SDisk) UpdataSnapshotsBackingDisk(backingDiskId string) error { return nil } -func (manager *SDiskManager) AutoSyncExtDiskSnapshot(ctx context.Context, userCred mcclient.TokenCredential, - isStart bool) { +func (manager *SDiskManager) AutoSyncExtDiskSnapshot(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) { - spds, err := manager.getAutoSnapshotDisksId(true) + now := time.Now() + q := SnapshotPolicyDiskManager.Query().LE("next_sync_time", now) + spds := make([]SSnapshotPolicyDisk, 0) + err := db.FetchModelObjects(SnapshotPolicyDiskManager, q, &spds) if err != nil { - log.Errorf("Get auto snapshot ext disks id failed: %s", err) - return + log.Errorf("unable to FetchModelObjects: %v", err) } - if len(spds) == 0 { - log.Infof("CronJob AutoSyncExtDiskSnapshot: No external disk need sync snapshot") - return + // fetch all snapshotpolicy + spIdSet := sets.NewString() + for i := range spds { + spIdSet.Insert(spds[i].SnapshotpolicyId) + } + sps, err := SnapshotPolicyManager.FetchAllByIds(spIdSet.UnsortedList()) + if err != nil { + log.Errorf("unable to FetchAllByIds: %v", err) + } + spMap := make(map[string]*SSnapshotPolicy, len(sps)) + for i := range sps { + spMap[sps[i].GetId()] = &sps[i] } for i := 0; i < len(spds); i++ { - disk := manager.FetchDiskById(spds[i].DiskId) + spd := &spds[i] + disk := manager.FetchDiskById(spd.DiskId) syncResult := disk.syncSnapshots(ctx, userCred) if syncResult.IsError() { db.OpsLog.LogEvent(disk, db.ACT_DISK_AUTO_SYNC_SNAPSHOT_FAIL, syncResult.Result(), userCred) continue } + if syncResult.AddCnt == 0 { + continue + } db.OpsLog.LogEvent(disk, db.ACT_DISK_AUTO_SYNC_SNAPSHOT, "disk auto sync snapshot successfully", userCred) + _, err := db.Update(spd, func() error { + spd.NextSyncTime = spMap[spd.GetId()].ComputeNextSyncTime(now, spd.NextSyncTime) + return nil + }) + if err != nil { + log.Errorf("unable to update NextSyncTime for snapshotpolicydisk %q %q", spd.SnapshotpolicyId, spd.DiskId) + } } } diff --git a/pkg/compute/models/snapshotpolicy.go b/pkg/compute/models/snapshotpolicy.go index e7486e51d1..a188bcdc45 100644 --- a/pkg/compute/models/snapshotpolicy.go +++ b/pkg/compute/models/snapshotpolicy.go @@ -17,6 +17,8 @@ package models import ( "context" "fmt" + "sort" + "time" "yunion.io/x/jsonutils" "yunion.io/x/log" @@ -215,7 +217,10 @@ func (manager *SSnapshotPolicyManager) OnCreateComplete(ctx context.Context, ite userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, query jsonutils.JSONObject, data jsonutils.JSONObject) { for i := range items { sp := items[i].(*SSnapshotPolicy) - sp.SetStatus(userCred, api.SNAPSHOT_POLICY_READY, "create complete") + db.Update(sp, func() error { + sp.Status = api.SNAPSHOT_POLICY_READY + return nil + }) } } @@ -722,6 +727,53 @@ func (self *SSnapshotPolicyManager) TimePointsToIntArray(n uint32) []int { return bitmap.Uint2IntArray(n) } +func (sp *SSnapshotPolicy) ComputeNextSyncTime(base, lastSyncTime time.Time) time.Time { + if base.IsZero() { + base = time.Now() + } + base = base.Truncate(time.Hour) + + baseWeekday := int(base.Weekday()) + if baseWeekday == 0 { + baseWeekday = 7 + } + weekDays := SnapshotPolicyManager.RepeatWeekdaysToIntArray(sp.RepeatWeekdays) + weekDays = append(weekDays, weekDays[0]+7) + index := sort.SearchInts(weekDays, baseWeekday) + addDay := weekDays[index] - baseWeekday + nextTime := base.AddDate(0, 0, addDay) + + // find timePoint closest to the base + timePoints := SnapshotPolicyManager.TimePointsToIntArray(sp.TimePoints) + var newHour int + if addDay > 0 { + newHour = timePoints[0] + } else { + baseHour := base.Hour() + index := sort.SearchInts(timePoints, baseHour) + index = index % len(timePoints) + if timePoints[index] == baseHour { + index = index + 1 + newHour = timePoints[index] + } else { + newHour = timePoints[index] + } + } + nextTime = time.Date(nextTime.Year(), nextTime.Month(), nextTime.Day(), newHour, 0, 0, 0, base.Location()) + + if sp.RetentionDays <= 0 { + return nextTime + } + if lastSyncTime.IsZero() { + lastSyncTime = base + } + snapshotRentionExpired := lastSyncTime.AddDate(0, 0, sp.RetentionDays) + if snapshotRentionExpired.Before(nextTime) { + return snapshotRentionExpired + } + return nextTime +} + func (sp *SSnapshotPolicy) GenerateCreateSpParams() *cloudprovider.SnapshotPolicyInput { intWeekdays := SnapshotPolicyManager.RepeatWeekdaysToIntArray(sp.RepeatWeekdays) intTimePoints := SnapshotPolicyManager.TimePointsToIntArray(sp.TimePoints) diff --git a/pkg/compute/models/snapshotpolicydisks.go b/pkg/compute/models/snapshotpolicydisks.go index f32cb58a14..e701a10e2f 100644 --- a/pkg/compute/models/snapshotpolicydisks.go +++ b/pkg/compute/models/snapshotpolicydisks.go @@ -20,10 +20,12 @@ import ( "database/sql" "fmt" "strings" + "time" "yunion.io/x/jsonutils" "yunion.io/x/log" "yunion.io/x/pkg/errors" + "yunion.io/x/pkg/util/sets" "yunion.io/x/sqlchemy" api "yunion.io/x/onecloud/pkg/apis/compute" @@ -73,9 +75,8 @@ type SSnapshotPolicyDisk struct { SSnapshotPolicyResourceBase `width:"36" charset:"ascii" nullable:"false" list:"user" create:"required" index:"true"` SDiskResourceBase `width:"36" charset:"ascii" nullable:"false" list:"user" create:"required" index:"true"` - // SnapshotpolicyId string `width:"36" charset:"ascii" nullable:"false" list:"user" create:"required" index:"true"` - // DiskId string `width:"36" charset:"ascii" nullable:"false" list:"user" create:"required" index:"true"` - Status string `width:"36" charset:"ascii" nullable:"false" default:"init" list:"user" create:"optional"` + Status string `width:"36" charset:"ascii" nullable:"false" default:"init" list:"user" create:"optional"` + NextSyncTime time.Time } func (sd *SSnapshotPolicyDisk) SetStatus(userCred mcclient.TokenCredential, status string, reason string) error { @@ -171,6 +172,41 @@ func (m *SSnapshotPolicyDiskManager) FetchBySnapshotPolicyDisk(spId, diskId stri return &ret[0], nil } +func (sdm *SSnapshotPolicyDiskManager) InitalizeData() error { + q := sdm.Query().IsNullOrEmpty("next_sync_time") + var sds []SSnapshotPolicyDisk + err := db.FetchModelObjects(sdm, q, &sds) + if err != nil { + return err + } + + // fetch all snapshotpolicy + spIdSet := sets.NewString() + for i := range sds { + spIdSet.Insert(sds[i].SnapshotpolicyId) + } + sps, err := SnapshotPolicyManager.FetchAllByIds(spIdSet.UnsortedList()) + if err != nil { + return errors.Wrap(err, "FetchAllByIds") + } + spMap := make(map[string]*SSnapshotPolicy, len(sps)) + for i := range sps { + spMap[sps[i].GetId()] = &sps[i] + } + now := time.Now() + for i := range sds { + sd := &sds[i] + _, err := db.Update(sd, func() error { + sd.NextSyncTime = spMap[sd.SnapshotpolicyId].ComputeNextSyncTime(now, now) + return nil + }) + if err != nil { + return errors.Wrap(err, "db.Update") + } + } + return nil +} + func (m *SSnapshotPolicyDiskManager) FetchAllByDiskID(ctx context.Context, userCred mcclient.TokenCredential, diskID string) ([]SSnapshotPolicyDisk, error) { @@ -420,6 +456,8 @@ func (self *SSnapshotPolicyDiskManager) newSnapshotpolicyDisk(ctx context.Contex spd := SSnapshotPolicyDisk{} spd.SnapshotpolicyId = sp.GetId() spd.DiskId = disk.GetId() + now := time.Now() + spd.NextSyncTime = sp.ComputeNextSyncTime(now, now) spd.SetModelManager(self, &spd) lockman.LockJointObject(ctx, disk, sp) diff --git a/pkg/compute/options/options.go b/pkg/compute/options/options.go index 14ba767354..d55ed9149e 100644 --- a/pkg/compute/options/options.go +++ b/pkg/compute/options/options.go @@ -144,9 +144,10 @@ type ComputeOptions struct { EnableAutoRenameProject bool `help:"when it set true, auto create project will rename when cloud project name changed" default:"false"` - SyncStorageCapacityUsedIntervalMinutes int `help:"interval sync storage capacity used" default:"10"` + SyncStorageCapacityUsedIntervalMinutes int `help:"interval sync storage capacity used" default:"10"` + LockStorageFromCachedimage bool `help:"must use storage in where selected cachedimage when creating vm"` - LockStorageFromCachedimage bool `help:"must use storage in where selected cachedimage when creating vm"` + SyncExtDiskSnapshotIntervalMinutes int `help:"sync snapshot for external disk" default:"20"` SCapabilityOptions SASControllerOptions diff --git a/pkg/compute/service/service.go b/pkg/compute/service/service.go index a80d8ba3d1..20b8701b54 100644 --- a/pkg/compute/service/service.go +++ b/pkg/compute/service/service.go @@ -143,9 +143,10 @@ func StartService() { cron.AddJobAtIntervalsWithStartRun("SyncCapacityUsedForStorage", time.Duration(opts.SyncStorageCapacityUsedIntervalMinutes)*time.Minute, models.StorageManager.SyncCapacityUsedForStorage, true) + cron.AddJobAtIntervalsWithStartRun("AutoSyncExtDiskSnapshot", time.Duration(opts.SyncExtDiskSnapshotIntervalMinutes)*time.Minute, models.DiskManager.AutoSyncExtDiskSnapshot, true) + cron.AddJobEveryFewHour("AutoDiskSnapshot", 1, 5, 0, models.DiskManager.AutoDiskSnapshot, false) cron.AddJobEveryFewHour("SnapshotsCleanup", 1, 35, 0, models.SnapshotManager.CleanupSnapshots, false) - cron.AddJobEveryFewHour("AutoSyncExtDiskSnapshot", 1, 10, 0, models.DiskManager.AutoSyncExtDiskSnapshot, false) cron.AddJobEveryFewDays("SyncSkus", opts.SyncSkusDay, opts.SyncSkusHour, 0, 0, models.SyncServerSkus, true) cron.AddJobEveryFewDays("SyncDBInstanceSkus", opts.SyncSkusDay, opts.SyncSkusHour, 0, 0, models.SyncDBInstanceSkus, true) cron.AddJobEveryFewDays("SyncElasticCacheSkus", opts.SyncSkusDay, opts.SyncSkusHour, 0, 0, models.SyncElasticCacheSkus, true)