mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-01 15:07:17 +08:00
fix(region): migrate snapshot policy resource (#23773)
This commit is contained in:
+18
-16
@@ -215,7 +215,7 @@ func (manager *SDiskManager) ListItemFilter(
|
||||
}
|
||||
|
||||
if query.BindingSnapshotpolicy != nil {
|
||||
spjsq := SnapshotPolicyDiskManager.Query("disk_id").SubQuery()
|
||||
spjsq := SnapshotPolicyResourceManager.Query("resource_id").Equals("resource_type", api.SNAPSHOT_POLICY_TYPE_DISK).SubQuery()
|
||||
if *query.BindingSnapshotpolicy {
|
||||
q = q.In("id", spjsq)
|
||||
} else {
|
||||
@@ -248,7 +248,7 @@ func (manager *SDiskManager) ListItemFilter(
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
sq := SnapshotPolicyDiskManager.Query("disk_id").Equals("snapshotpolicy_id", query.SnapshotpolicyId)
|
||||
sq := SnapshotPolicyResourceManager.Query("resource_id").Equals("resource_type", api.SNAPSHOT_POLICY_TYPE_DISK).Equals("snapshotpolicy_id", query.SnapshotpolicyId).SubQuery()
|
||||
q = q.In("id", sq)
|
||||
}
|
||||
|
||||
@@ -2465,11 +2465,13 @@ func (self *SDisk) Delete(ctx context.Context, userCred mcclient.TokenCredential
|
||||
func (self *SDisk) RealDelete(ctx context.Context, userCred mcclient.TokenCredential) error {
|
||||
// diskbackups := DiskBackupManager.Query("id").Equals("disk_id", self.Id)
|
||||
guestdisks := GuestdiskManager.Query("row_id").Equals("disk_id", self.Id)
|
||||
diskpolicies := SnapshotPolicyDiskManager.Query("row_id").Equals("disk_id", self.Id)
|
||||
err := SnapshotPolicyResourceManager.RemoveByResource(self.Id, api.SNAPSHOT_POLICY_TYPE_DISK)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "RemoveByResource")
|
||||
}
|
||||
pairs := []purgePair{
|
||||
// {manager: DiskBackupManager, key: "id", q: diskbackups},
|
||||
{manager: GuestdiskManager, key: "row_id", q: guestdisks},
|
||||
{manager: SnapshotPolicyDiskManager, key: "row_id", q: diskpolicies},
|
||||
}
|
||||
for i := range pairs {
|
||||
err := pairs[i].purgeAll(ctx)
|
||||
@@ -2664,16 +2666,16 @@ func (manager *SDiskManager) FetchCustomizeColumns(
|
||||
}
|
||||
|
||||
policySQ := SnapshotPolicyManager.Query().SubQuery()
|
||||
dps := SnapshotPolicyDiskManager.Query().SubQuery()
|
||||
dps := SnapshotPolicyResourceManager.Query().Equals("resource_type", api.SNAPSHOT_POLICY_TYPE_DISK).SubQuery()
|
||||
|
||||
q = policySQ.Query(
|
||||
policySQ.Field("id"),
|
||||
policySQ.Field("name"),
|
||||
policySQ.Field("time_points"),
|
||||
policySQ.Field("repeat_weekdays"),
|
||||
dps.Field("disk_id"),
|
||||
dps.Field("resource_id"),
|
||||
).Join(dps, sqlchemy.Equals(dps.Field("snapshotpolicy_id"), policySQ.Field("id"))).
|
||||
Filter(sqlchemy.In(dps.Field("disk_id"), diskIds))
|
||||
Filter(sqlchemy.In(dps.Field("resource_id"), diskIds))
|
||||
|
||||
policyInfo := []struct {
|
||||
Id string
|
||||
@@ -2929,7 +2931,7 @@ func (manager *SDiskManager) CleanPendingDeleteDisks(ctx context.Context, userCr
|
||||
}
|
||||
}
|
||||
|
||||
func (manager *SDiskManager) GetNeedAutoSnapshotDisks() ([]SSnapshotPolicyDisk, error) {
|
||||
func (manager *SDiskManager) GetNeedAutoSnapshotDisks() ([]SSnapshotPolicyResource, error) {
|
||||
tz, _ := time.LoadLocation(options.Options.TimeZone)
|
||||
t := time.Now().In(tz)
|
||||
week := t.Weekday()
|
||||
@@ -2949,11 +2951,11 @@ func (manager *SDiskManager) GetNeedAutoSnapshotDisks() ([]SSnapshotPolicyDisk,
|
||||
),
|
||||
).SubQuery()
|
||||
disks := DiskManager.Query().SubQuery()
|
||||
q := SnapshotPolicyDiskManager.Query()
|
||||
q := SnapshotPolicyResourceManager.Query().Equals("resource_type", api.SNAPSHOT_POLICY_TYPE_DISK)
|
||||
q = q.Join(sq, sqlchemy.Equals(q.Field("snapshotpolicy_id"), sq.Field("id")))
|
||||
q = q.Join(disks, sqlchemy.Equals(q.Field("disk_id"), disks.Field("id")))
|
||||
ret := []SSnapshotPolicyDisk{}
|
||||
err := db.FetchModelObjects(SnapshotPolicyDiskManager, q, &ret)
|
||||
q = q.Join(disks, sqlchemy.Equals(q.Field("resource_id"), disks.Field("id")))
|
||||
ret := []SSnapshotPolicyResource{}
|
||||
err := db.FetchModelObjects(SnapshotPolicyResourceManager, q, &ret)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
@@ -2988,7 +2990,7 @@ func (manager *SDiskManager) AutoDiskSnapshot(ctx context.Context, userCred mccl
|
||||
}
|
||||
log.Infof("auto snapshot %d disks", len(disks))
|
||||
|
||||
guestDps := map[string][]SSnapshotPolicyDisk{}
|
||||
guestDps := map[string][]SSnapshotPolicyResource{}
|
||||
for i := 0; i < len(disks); i++ {
|
||||
disk, err := disks[i].GetDisk()
|
||||
if err != nil {
|
||||
@@ -2999,7 +3001,7 @@ func (manager *SDiskManager) AutoDiskSnapshot(ctx context.Context, userCred mccl
|
||||
if dps, ok := guestDps[guest.Id]; ok {
|
||||
guestDps[guest.Id] = append(dps, disks[i])
|
||||
} else {
|
||||
guestDps[guest.Id] = []SSnapshotPolicyDisk{disks[i]}
|
||||
guestDps[guest.Id] = []SSnapshotPolicyResource{disks[i]}
|
||||
}
|
||||
continue
|
||||
}
|
||||
@@ -3022,7 +3024,7 @@ func (manager *SDiskManager) AutoDiskSnapshot(ctx context.Context, userCred mccl
|
||||
}
|
||||
|
||||
func (manager *SDiskManager) OrderCreateDisksSnapshotsBySnapshotPolicy(
|
||||
ctx context.Context, userCred mcclient.TokenCredential, guest *SGuest, snapshotPolicyDisks []SSnapshotPolicyDisk,
|
||||
ctx context.Context, userCred mcclient.TokenCredential, guest *SGuest, snapshotPolicyDisks []SSnapshotPolicyResource,
|
||||
) error {
|
||||
params := jsonutils.NewDict()
|
||||
params.Set("snapshot_policy_disks", jsonutils.Marshal(snapshotPolicyDisks))
|
||||
@@ -3036,7 +3038,7 @@ func (manager *SDiskManager) OrderCreateDisksSnapshotsBySnapshotPolicy(
|
||||
|
||||
func (manager *SDiskManager) DoAutoSnapshot(
|
||||
ctx context.Context, userCred mcclient.TokenCredential,
|
||||
diskSnapshotPolicy *SSnapshotPolicyDisk, disk *SDisk, parentTaskId string,
|
||||
diskSnapshotPolicy *SSnapshotPolicyResource, disk *SDisk, parentTaskId string,
|
||||
) error {
|
||||
policy, err := diskSnapshotPolicy.GetSnapshotPolicy()
|
||||
if err != nil {
|
||||
|
||||
@@ -87,10 +87,10 @@ func InitDB() error {
|
||||
now := time.Now()
|
||||
err := manager.InitializeData()
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "%s InitializeData", manager.Keyword())
|
||||
return errors.Wrapf(err, "%s initializeData", manager.Keyword())
|
||||
}
|
||||
if cost := time.Now().Sub(now); cost > time.Duration(time.Second)*15 {
|
||||
log.Infof("%s InitializeData cost %s", manager.Keyword(), cost.Round(time.Second))
|
||||
log.Infof("%s initializeData cost %s", manager.Keyword(), cost.Round(time.Second))
|
||||
}
|
||||
}
|
||||
return nil
|
||||
|
||||
@@ -640,10 +640,10 @@ func (self *SZone) purgeStorages(ctx context.Context, managerId string) error {
|
||||
disks := DiskManager.Query("id").In("storage_id", storages.SubQuery())
|
||||
diskbackups := DiskBackupManager.Query("id").In("disk_id", disks.SubQuery())
|
||||
guestdisks := GuestdiskManager.Query("row_id").In("disk_id", disks.SubQuery())
|
||||
diskpolicies := SnapshotPolicyDiskManager.Query("row_id").In("disk_id", disks.SubQuery())
|
||||
diskpolicies := SnapshotPolicyResourceManager.Query("row_id").In("resource_id", disks.SubQuery())
|
||||
|
||||
pairs := []purgePair{
|
||||
{manager: SnapshotPolicyDiskManager, key: "row_id", q: diskpolicies},
|
||||
{manager: SnapshotPolicyResourceManager, key: "row_id", q: diskpolicies},
|
||||
{manager: GuestdiskManager, key: "row_id", q: guestdisks},
|
||||
{manager: SnapshotManager, key: "id", q: snapshots},
|
||||
{manager: DiskBackupManager, key: "id", q: diskbackups},
|
||||
@@ -668,11 +668,9 @@ func (self *SStorage) purge(ctx context.Context, userCred mcclient.TokenCredenti
|
||||
disks := DiskManager.Query("id").Equals("storage_id", self.Id)
|
||||
diskbackups := DiskBackupManager.Query("id").In("disk_id", disks.SubQuery())
|
||||
guestdisks := GuestdiskManager.Query("row_id").In("disk_id", disks.SubQuery())
|
||||
diskpolicies := SnapshotPolicyDiskManager.Query("row_id").In("disk_id", disks.SubQuery())
|
||||
|
||||
pairs := []purgePair{
|
||||
{manager: GuestdiskManager, key: "row_id", q: guestdisks},
|
||||
{manager: SnapshotPolicyDiskManager, key: "row_id", q: diskpolicies},
|
||||
{manager: SnapshotManager, key: "id", q: snapshots},
|
||||
{manager: DiskBackupManager, key: "id", q: diskbackups},
|
||||
{manager: DiskManager, key: "id", q: disks},
|
||||
|
||||
@@ -66,6 +66,22 @@ func (self *SSnapshotPolicyResource) GetServer() (*SGuest, error) {
|
||||
return guest.(*SGuest), nil
|
||||
}
|
||||
|
||||
func (self *SSnapshotPolicyResource) GetDisk() (*SDisk, error) {
|
||||
disk, err := DiskManager.FetchById(self.ResourceId)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return disk.(*SDisk), nil
|
||||
}
|
||||
|
||||
func (self *SSnapshotPolicyResource) GetSnapshotPolicy() (*SSnapshotPolicy, error) {
|
||||
policy, err := SnapshotPolicyManager.FetchById(self.SnapshotpolicyId)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return policy.(*SSnapshotPolicy), nil
|
||||
}
|
||||
|
||||
func (man *SSnapshotPolicyResourceManager) RemoveBySnapshotpolicy(id string) error {
|
||||
_, err := sqlchemy.GetDB().Exec(
|
||||
fmt.Sprintf(
|
||||
|
||||
@@ -21,6 +21,7 @@ import (
|
||||
|
||||
"yunion.io/x/cloudmux/pkg/cloudprovider"
|
||||
"yunion.io/x/jsonutils"
|
||||
"yunion.io/x/log"
|
||||
"yunion.io/x/pkg/errors"
|
||||
"yunion.io/x/pkg/util/compare"
|
||||
"yunion.io/x/pkg/utils"
|
||||
@@ -213,24 +214,11 @@ func (manager *SSnapshotPolicyManager) FetchCustomizeColumns(
|
||||
policyIds[i] = policy.Id
|
||||
}
|
||||
|
||||
q := SnapshotPolicyDiskManager.Query().In("snapshotpolicy_id", policyIds)
|
||||
pds := []SSnapshotPolicyDisk{}
|
||||
err := q.All(&pds)
|
||||
if err != nil {
|
||||
return rows
|
||||
}
|
||||
pdMap := map[string][]SSnapshotPolicyDisk{}
|
||||
for _, pd := range pds {
|
||||
_, ok := pdMap[pd.SnapshotpolicyId]
|
||||
if !ok {
|
||||
pdMap[pd.SnapshotpolicyId] = []SSnapshotPolicyDisk{}
|
||||
}
|
||||
pdMap[pd.SnapshotpolicyId] = append(pdMap[pd.SnapshotpolicyId], pd)
|
||||
}
|
||||
q = SnapshotPolicyResourceManager.Query().In("snapshotpolicy_id", policyIds)
|
||||
q := SnapshotPolicyResourceManager.Query().In("snapshotpolicy_id", policyIds)
|
||||
sprs := []SSnapshotPolicyResource{}
|
||||
err = q.All(&sprs)
|
||||
err := q.All(&sprs)
|
||||
if err != nil {
|
||||
log.Errorf("query snapshot policy resources error: %v", err)
|
||||
return rows
|
||||
}
|
||||
sprmap := map[string][]SSnapshotPolicyResource{}
|
||||
@@ -242,10 +230,12 @@ func (manager *SSnapshotPolicyManager) FetchCustomizeColumns(
|
||||
sprmap[sp.SnapshotpolicyId] = append(sprmap[sp.SnapshotpolicyId], sp)
|
||||
}
|
||||
for i := range rows {
|
||||
disks := pdMap[policyIds[i]]
|
||||
rows[i].BindingDiskCount = len(disks)
|
||||
resources := sprmap[policyIds[i]]
|
||||
rows[i].BindingResourceCount = len(resources)
|
||||
sp := objs[i].(*SSnapshotPolicy)
|
||||
if sp.Type == api.SNAPSHOT_POLICY_TYPE_DISK {
|
||||
rows[i].BindingDiskCount = len(resources)
|
||||
}
|
||||
}
|
||||
|
||||
return rows
|
||||
@@ -435,15 +425,11 @@ func (sp *SSnapshotPolicy) Delete(ctx context.Context, userCred mcclient.TokenCr
|
||||
}
|
||||
|
||||
func (sp *SSnapshotPolicy) RealDelete(ctx context.Context, userCred mcclient.TokenCredential) error {
|
||||
err := SnapshotPolicyDiskManager.RemoveBySnapshotpolicy(sp.Id)
|
||||
err := SnapshotPolicyResourceManager.RemoveBySnapshotpolicy(sp.Id)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "delete snapshot policy disks for policy %s", sp.Name)
|
||||
return errors.Wrapf(err, "RemoveBySnapshotpolicy for policy %s", sp.Name)
|
||||
}
|
||||
err = SnapshotPolicyResourceManager.RemoveBySnapshotpolicy(sp.Id)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "delete snapshot policy resources for policy %s", sp.Name)
|
||||
}
|
||||
return db.DeleteModel(ctx, userCred, sp)
|
||||
return sp.SVirtualResourceBase.Delete(ctx, userCred)
|
||||
}
|
||||
|
||||
func (sp *SSnapshotPolicy) StartBindDisksTask(ctx context.Context, userCred mcclient.TokenCredential, diskIds []string) error {
|
||||
@@ -682,8 +668,8 @@ func (manager *SSnapshotPolicyManager) OrderByExtraFields(
|
||||
}
|
||||
|
||||
if db.NeedOrderQuery([]string{input.OrderByBindDiskCount}) {
|
||||
sdQ := SnapshotPolicyDiskManager.Query()
|
||||
sdSQ := sdQ.AppendField(sdQ.Field("snapshotpolicy_id"), sqlchemy.COUNT("disk_count")).GroupBy("snapshotpolicy_id").SubQuery()
|
||||
sdQ := SnapshotPolicyResourceManager.Query().Equals("resource_type", api.SNAPSHOT_POLICY_TYPE_DISK)
|
||||
sdSQ := sdQ.AppendField(sdQ.Field("snapshotpolicy_id"), sqlchemy.COUNT("disk_count")).GroupBy(sdQ.Field("snapshotpolicy_id")).SubQuery()
|
||||
q = q.LeftJoin(sdSQ, sqlchemy.Equals(sdSQ.Field("snapshotpolicy_id"), q.Field("id")))
|
||||
q = q.AppendField(q.QueryFields()...)
|
||||
q = q.AppendField(sdSQ.Field("disk_count"))
|
||||
@@ -789,7 +775,7 @@ func (self *SSnapshotPolicy) GetProvider(ctx context.Context) (cloudprovider.ICl
|
||||
}
|
||||
|
||||
func (self *SSnapshotPolicy) GetUnbindDisks(diskIds []string) ([]SDisk, error) {
|
||||
sq := SnapshotPolicyDiskManager.Query("disk_id").Equals("snapshotpolicy_id", self.Id).SubQuery()
|
||||
sq := SnapshotPolicyResourceManager.Query("resource_id").Equals("resource_type", api.SNAPSHOT_POLICY_TYPE_DISK).Equals("snapshotpolicy_id", self.Id).SubQuery()
|
||||
q := DiskManager.Query().In("id", diskIds)
|
||||
q = q.Filter(sqlchemy.NotIn(q.Field("id"), sq))
|
||||
ret := []SDisk{}
|
||||
@@ -801,7 +787,7 @@ func (self *SSnapshotPolicy) GetUnbindDisks(diskIds []string) ([]SDisk, error) {
|
||||
}
|
||||
|
||||
func (self *SSnapshotPolicy) GetBindDisks(diskIds []string) ([]SDisk, error) {
|
||||
sq := SnapshotPolicyDiskManager.Query("disk_id").Equals("snapshotpolicy_id", self.Id).SubQuery()
|
||||
sq := SnapshotPolicyResourceManager.Query("resource_id").Equals("resource_type", api.SNAPSHOT_POLICY_TYPE_DISK).Equals("snapshotpolicy_id", self.Id).SubQuery()
|
||||
q := DiskManager.Query().In("id", diskIds)
|
||||
q = q.Filter(sqlchemy.In(q.Field("id"), sq))
|
||||
ret := []SDisk{}
|
||||
@@ -813,9 +799,9 @@ func (self *SSnapshotPolicy) GetBindDisks(diskIds []string) ([]SDisk, error) {
|
||||
}
|
||||
|
||||
func (self *SSnapshotPolicy) GetDisks() ([]SDisk, error) {
|
||||
sq := SnapshotPolicyDiskManager.Query().Equals("snapshotpolicy_id", self.Id).SubQuery()
|
||||
sq := SnapshotPolicyResourceManager.Query().Equals("resource_type", api.SNAPSHOT_POLICY_TYPE_DISK).Equals("snapshotpolicy_id", self.Id).SubQuery()
|
||||
q := DiskManager.Query()
|
||||
q = q.Join(sq, sqlchemy.Equals(q.Field("id"), sq.Field("disk_id")))
|
||||
q = q.Join(sq, sqlchemy.Equals(q.Field("id"), sq.Field("resource_id")))
|
||||
ret := []SDisk{}
|
||||
err := db.FetchModelObjects(DiskManager, q, &ret)
|
||||
if err != nil {
|
||||
@@ -826,11 +812,12 @@ func (self *SSnapshotPolicy) GetDisks() ([]SDisk, error) {
|
||||
|
||||
func (sp *SSnapshotPolicy) BindDisks(ctx context.Context, disks []SDisk) error {
|
||||
for i := range disks {
|
||||
spd := &SSnapshotPolicyDisk{}
|
||||
spd.SetModelManager(SnapshotPolicyDiskManager, spd)
|
||||
spd.DiskId = disks[i].Id
|
||||
spd := &SSnapshotPolicyResource{}
|
||||
spd.SetModelManager(SnapshotPolicyResourceManager, spd)
|
||||
spd.ResourceId = disks[i].Id
|
||||
spd.ResourceType = api.SNAPSHOT_POLICY_TYPE_DISK
|
||||
spd.SnapshotpolicyId = sp.Id
|
||||
err := SnapshotPolicyDiskManager.TableSpec().Insert(ctx, spd)
|
||||
err := SnapshotPolicyResourceManager.TableSpec().Insert(ctx, spd)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
@@ -847,8 +834,8 @@ func (sp *SSnapshotPolicy) UnbindDisks(diskIds []string) error {
|
||||
}
|
||||
_, err := sqlchemy.GetDB().Exec(
|
||||
fmt.Sprintf(
|
||||
"delete from %s where snapshotpolicy_id = ? and disk_id in (%s)",
|
||||
SnapshotPolicyDiskManager.TableSpec().Name(), strings.Join(placeholders, ","),
|
||||
"delete from %s where snapshotpolicy_id = ? and resource_id in (%s)",
|
||||
SnapshotPolicyResourceManager.TableSpec().Name(), strings.Join(placeholders, ","),
|
||||
), vars...,
|
||||
)
|
||||
return err
|
||||
@@ -863,26 +850,22 @@ func (sp *SSnapshotPolicy) SyncDisks(ctx context.Context, userCred mcclient.Toke
|
||||
return errors.Wrapf(err, "GetApplyDiskIds")
|
||||
}
|
||||
{
|
||||
sq := SnapshotPolicyDiskManager.Query("disk_id").Equals("snapshotpolicy_id", sp.Id).SubQuery()
|
||||
sq := SnapshotPolicyResourceManager.Query("resource_id").Equals("resource_type", api.SNAPSHOT_POLICY_TYPE_DISK).Equals("snapshotpolicy_id", sp.Id).SubQuery()
|
||||
q := DiskManager.Query().In("id", sq).NotIn("external_id", extIds)
|
||||
needCancel := []SDisk{}
|
||||
err = db.FetchModelObjects(DiskManager, q, &needCancel)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "db.FetchModelObjects")
|
||||
}
|
||||
diskIds := []string{}
|
||||
for _, disk := range needCancel {
|
||||
diskIds = append(diskIds, disk.Id)
|
||||
}
|
||||
if len(diskIds) > 0 {
|
||||
err = sp.UnbindDisks(diskIds)
|
||||
err = SnapshotPolicyResourceManager.RemoveByResource(disk.Id, api.SNAPSHOT_POLICY_TYPE_DISK)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "UnbindDisks")
|
||||
return errors.Wrapf(err, "RemoveByResource")
|
||||
}
|
||||
}
|
||||
}
|
||||
{
|
||||
sq := SnapshotPolicyDiskManager.Query("disk_id").Equals("snapshotpolicy_id", sp.Id).SubQuery()
|
||||
sq := SnapshotPolicyResourceManager.Query("resource_id").Equals("resource_type", api.SNAPSHOT_POLICY_TYPE_DISK).Equals("snapshotpolicy_id", sp.Id).SubQuery()
|
||||
storages := StorageManager.Query().Equals("manager_id", sp.ManagerId).SubQuery()
|
||||
q := DiskManager.Query()
|
||||
q = q.Join(storages, sqlchemy.Equals(q.Field("storage_id"), storages.Field("id")))
|
||||
@@ -897,9 +880,16 @@ func (sp *SSnapshotPolicy) SyncDisks(ctx context.Context, userCred mcclient.Toke
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "db.FetchModelObjects")
|
||||
}
|
||||
err = sp.BindDisks(ctx, needApply)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "BindDisks")
|
||||
for _, disk := range needApply {
|
||||
spd := &SSnapshotPolicyResource{}
|
||||
spd.SetModelManager(SnapshotPolicyResourceManager, spd)
|
||||
spd.SnapshotpolicyId = sp.Id
|
||||
spd.ResourceId = disk.Id
|
||||
spd.ResourceType = api.SNAPSHOT_POLICY_TYPE_DISK
|
||||
err := SnapshotPolicyResourceManager.TableSpec().Insert(ctx, spd)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "Insert")
|
||||
}
|
||||
}
|
||||
}
|
||||
return nil
|
||||
|
||||
@@ -15,11 +15,12 @@
|
||||
package models
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
|
||||
"yunion.io/x/pkg/errors"
|
||||
"yunion.io/x/sqlchemy"
|
||||
|
||||
api "yunion.io/x/onecloud/pkg/apis/compute"
|
||||
"yunion.io/x/onecloud/pkg/cloudcommon/db"
|
||||
)
|
||||
|
||||
@@ -62,38 +63,41 @@ type SSnapshotPolicyDisk struct {
|
||||
SDiskResourceBase `width:"36" charset:"ascii" nullable:"false" list:"user" create:"required" index:"true"`
|
||||
}
|
||||
|
||||
func (self *SSnapshotPolicyDisk) GetDisk() (*SDisk, error) {
|
||||
disk, err := DiskManager.FetchById(self.DiskId)
|
||||
func (manager *SSnapshotPolicyDiskManager) InitializeData() error {
|
||||
disks := DiskManager.Query("id").SubQuery()
|
||||
policies := SnapshotPolicyManager.Query("id").SubQuery()
|
||||
q := manager.Query().In("disk_id", disks).In("snapshotpolicy_id", policies)
|
||||
|
||||
pds := []SSnapshotPolicyDisk{}
|
||||
err := db.FetchModelObjects(manager, q, &pds)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "FetchById(%s)", self.DiskId)
|
||||
return err
|
||||
}
|
||||
return disk.(*SDisk), nil
|
||||
}
|
||||
|
||||
func (self *SSnapshotPolicyDisk) GetSnapshotPolicy() (*SSnapshotPolicy, error) {
|
||||
policy, err := SnapshotPolicyManager.FetchById(self.SnapshotpolicyId)
|
||||
if err != nil {
|
||||
return nil, errors.Wrapf(err, "FetchById(%s)", self.SnapshotpolicyId)
|
||||
for i := range pds {
|
||||
pd := &pds[i]
|
||||
migrateData := &SSnapshotPolicyResource{
|
||||
SnapshotpolicyId: pd.SnapshotpolicyId,
|
||||
ResourceId: pd.DiskId,
|
||||
ResourceType: api.SNAPSHOT_POLICY_TYPE_DISK,
|
||||
}
|
||||
migrateData.SetModelManager(SnapshotPolicyResourceManager, migrateData)
|
||||
cnt, err := SnapshotPolicyResourceManager.Query().
|
||||
Equals("resource_type", api.SNAPSHOT_POLICY_TYPE_DISK).
|
||||
Equals("resource_id", pd.DiskId).
|
||||
Equals("snapshotpolicy_id", pd.SnapshotpolicyId).CountWithError()
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "Count")
|
||||
}
|
||||
if cnt == 0 {
|
||||
err = SnapshotPolicyResourceManager.TableSpec().Insert(context.Background(), migrateData)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "Insert %s", migrateData.Keyword())
|
||||
}
|
||||
}
|
||||
err = db.Purge(manager, "row_id", []string{fmt.Sprintf("%d", pd.RowId)}, true)
|
||||
if err != nil {
|
||||
return errors.Wrapf(err, "Purge %d", pd.RowId)
|
||||
}
|
||||
}
|
||||
return policy.(*SSnapshotPolicy), nil
|
||||
}
|
||||
|
||||
func (man *SSnapshotPolicyDiskManager) RemoveByDisk(id string) error {
|
||||
_, err := sqlchemy.GetDB().Exec(
|
||||
fmt.Sprintf(
|
||||
"delete from %s where disk_id = ?",
|
||||
man.TableSpec().Name(),
|
||||
), id,
|
||||
)
|
||||
return err
|
||||
}
|
||||
|
||||
func (man *SSnapshotPolicyDiskManager) RemoveBySnapshotpolicy(id string) error {
|
||||
_, err := sqlchemy.GetDB().Exec(
|
||||
fmt.Sprintf(
|
||||
"delete from %s where snapshotpolicy_id = ?",
|
||||
man.TableSpec().Name(),
|
||||
), id,
|
||||
)
|
||||
return err
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -1241,10 +1241,10 @@ func (manager *SSnapshotManager) CleanupSnapshots(ctx context.Context, userCred
|
||||
|
||||
{
|
||||
sq = SnapshotPolicyManager.Query().Equals("type", api.SNAPSHOT_POLICY_TYPE_DISK).GT("retention_count", 0).SubQuery()
|
||||
spd := SnapshotPolicyDiskManager.Query().SubQuery()
|
||||
spd := SnapshotPolicyResourceManager.Query().Equals("resource_type", api.SNAPSHOT_POLICY_TYPE_DISK).SubQuery()
|
||||
q = sq.Query(
|
||||
sq.Field("retention_count"),
|
||||
spd.Field("disk_id"),
|
||||
spd.Field("resource_id").Label("disk_id"),
|
||||
)
|
||||
q = q.Join(spd, sqlchemy.Equals(q.Field("id"), spd.Field("snapshotpolicy_id")))
|
||||
|
||||
|
||||
@@ -104,7 +104,6 @@ func InitHandlers(app *appsrv.Application) {
|
||||
models.WafRuleStatementManager,
|
||||
models.BillingResourceCheckManager,
|
||||
|
||||
models.SnapshotPolicyDiskManager,
|
||||
models.LoadbalancerSecurityGroupManager,
|
||||
|
||||
models.HostFileJointsManager,
|
||||
|
||||
@@ -197,7 +197,7 @@ func (self *DiskDeleteTask) startPendingDeleteDisk(ctx context.Context, disk *mo
|
||||
self.OnGuestDiskDeleteCompleteFailed(ctx, disk, jsonutils.NewString("pending delete disk failed"))
|
||||
return
|
||||
}
|
||||
err = models.SnapshotPolicyDiskManager.RemoveByDisk(disk.Id)
|
||||
err = models.SnapshotPolicyResourceManager.RemoveByResource(disk.Id, api.SNAPSHOT_POLICY_TYPE_DISK)
|
||||
if err != nil {
|
||||
self.OnGuestDiskDeleteCompleteFailed(ctx, disk,
|
||||
jsonutils.NewString("detach all snapshotpolicies of disk failed"))
|
||||
|
||||
@@ -85,7 +85,7 @@ func (self *GuestDisksSnapshotPolicyExecuteTask) OnInit(ctx context.Context, obj
|
||||
}
|
||||
|
||||
func (self *GuestDisksSnapshotPolicyExecuteTask) OnDiskSnapshot(ctx context.Context, obj db.IStandaloneModel, data jsonutils.JSONObject) {
|
||||
snapshotPolicyDisks := make([]models.SSnapshotPolicyDisk, 0)
|
||||
snapshotPolicyDisks := make([]models.SSnapshotPolicyResource, 0)
|
||||
self.Params.Unmarshal(&snapshotPolicyDisks, "snapshot_policy_disks")
|
||||
if len(snapshotPolicyDisks) == 0 {
|
||||
self.SetStageComplete(ctx, nil)
|
||||
|
||||
Reference in New Issue
Block a user