mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
feat: add a retry mechanism to AutoSyncExtDiskSnapshot
This commit is contained in:
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user