From e80be2b586eb991bdcc88461c1b288eac558bd49 Mon Sep 17 00:00:00 2001 From: Qu Xuan Date: Fri, 28 May 2021 11:29:57 +0800 Subject: [PATCH] fix(region): optimized disk eip snapshot project sync --- pkg/compute/models/cloudsync.go | 1 - pkg/compute/models/disks.go | 71 ++++++++++++++------------------ pkg/compute/models/elasticips.go | 25 +++++++---- pkg/compute/models/guests.go | 11 +++-- pkg/compute/models/networks.go | 2 - pkg/compute/models/snapshots.go | 13 +++--- pkg/multicloud/aws/network.go | 1 - 7 files changed, 58 insertions(+), 66 deletions(-) diff --git a/pkg/compute/models/cloudsync.go b/pkg/compute/models/cloudsync.go index 6de027f3be..ec6c0afd31 100644 --- a/pkg/compute/models/cloudsync.go +++ b/pkg/compute/models/cloudsync.go @@ -758,7 +758,6 @@ func syncHostVMs(ctx context.Context, userCred mcclient.TokenCredential, syncRes } func syncVMPeripherals(ctx context.Context, userCred mcclient.TokenCredential, local *SGuest, remote cloudprovider.ICloudVM, host *SHost, provider *SCloudprovider, driver cloudprovider.ICloudProvider) { - syncVirtualResourceMetadata(ctx, userCred, local, remote) err := syncVMNics(ctx, userCred, provider, host, local, remote) if err != nil { log.Errorf("syncVMNics error %s", err) diff --git a/pkg/compute/models/disks.go b/pkg/compute/models/disks.go index 8243baa570..d4478ec842 100644 --- a/pkg/compute/models/disks.go +++ b/pkg/compute/models/disks.go @@ -1248,42 +1248,34 @@ func (manager *SDiskManager) getDisksByStorage(storage *SStorage) ([]SDisk, erro } func (manager *SDiskManager) syncCloudDisk(ctx context.Context, userCred mcclient.TokenCredential, provider cloudprovider.ICloudProvider, vdisk cloudprovider.ICloudDisk, index int, syncOwnerId mcclient.IIdentityProvider, managerId string) (*SDisk, error) { - // ownerProjId := projectId - - lockman.LockClass(ctx, manager, db.GetLockClassKey(manager, syncOwnerId)) - defer lockman.ReleaseClass(ctx, manager, db.GetLockClassKey(manager, syncOwnerId)) - diskObj, err := db.FetchByExternalIdAndManagerId(manager, vdisk.GetGlobalId(), func(q *sqlchemy.SQuery) *sqlchemy.SQuery { sq := StorageManager.Query().SubQuery() return q.Join(sq, sqlchemy.Equals(sq.Field("id"), q.Field("storage_id"))).Filter(sqlchemy.Equals(sq.Field("manager_id"), managerId)) }) if err != nil { - if err == sql.ErrNoRows { - vstorage, err := vdisk.GetIStorage() - if err != nil { - return nil, errors.Wrapf(err, "unable to GetIStorage of vdisk %q", vdisk.GetName()) - } - - storageObj, err := db.FetchByExternalIdAndManagerId(StorageManager, vstorage.GetGlobalId(), func(q *sqlchemy.SQuery) *sqlchemy.SQuery { - return q.Equals("manager_id", managerId) - }) - if err != nil { - log.Errorf("cannot find storage of vdisk %s", err) - return nil, err - } - storage := storageObj.(*SStorage) - return manager.newFromCloudDisk(ctx, userCred, provider, vdisk, storage, -1, syncOwnerId) - } else { - return nil, err + if errors.Cause(err) != sql.ErrNoRows { + return nil, errors.Wrapf(err, "db.FetchByExternalIdAndManagerId") } - } else { - disk := diskObj.(*SDisk) - err = disk.syncWithCloudDisk(ctx, userCred, provider, vdisk, index, syncOwnerId, managerId) + vstorage, err := vdisk.GetIStorage() if err != nil { - return nil, err + return nil, errors.Wrapf(err, "unable to GetIStorage of vdisk %q", vdisk.GetName()) } - return disk, nil + + storageObj, err := db.FetchByExternalIdAndManagerId(StorageManager, vstorage.GetGlobalId(), func(q *sqlchemy.SQuery) *sqlchemy.SQuery { + return q.Equals("manager_id", managerId) + }) + if err != nil { + return nil, errors.Wrapf(err, "cannot find storage of vdisk %s", vdisk.GetName()) + } + storage := storageObj.(*SStorage) + return manager.newFromCloudDisk(ctx, userCred, provider, vdisk, storage, -1, syncOwnerId) } + disk := diskObj.(*SDisk) + err = disk.syncWithCloudDisk(ctx, userCred, provider, vdisk, index, syncOwnerId, managerId) + if err != nil { + return nil, errors.Wrapf(err, "syncWithCloudDisk") + } + return disk, nil } func (manager *SDiskManager) SyncDisks(ctx context.Context, userCred mcclient.TokenCredential, provider cloudprovider.ICloudProvider, storage *SStorage, disks []cloudprovider.ICloudDisk, syncOwnerId mcclient.IIdentityProvider) ([]SDisk, []cloudprovider.ICloudDisk, compare.SyncResult) { @@ -1327,7 +1319,6 @@ func (manager *SDiskManager) SyncDisks(ctx context.Context, userCred mcclient.To if err != nil { syncResult.UpdateError(err) } else { - syncVirtualResourceMetadata(ctx, userCred, &commondb[i], commonext[i]) localDisks = append(localDisks, commondb[i]) remoteDisks = append(remoteDisks, commonext[i]) syncResult.Update() @@ -1360,7 +1351,6 @@ func (manager *SDiskManager) SyncDisks(ctx context.Context, userCred mcclient.To if err != nil { syncResult.AddError(err) } else { - syncVirtualResourceMetadata(ctx, userCred, new, added[i]) localDisks = append(localDisks, *new) remoteDisks = append(remoteDisks, added[i]) syncResult.Add() @@ -1371,19 +1361,16 @@ func (manager *SDiskManager) SyncDisks(ctx context.Context, userCred mcclient.To } func (self *SDisk) syncDiskStorage(ctx context.Context, userCred mcclient.TokenCredential, idisk cloudprovider.ICloudDisk, managerId string) error { - extId := idisk.GetGlobalId() istorage, err := idisk.GetIStorage() if err != nil { - log.Errorf("failed to get istorage for disk %s error: %v", extId, err) - return err + return errors.Wrapf(err, "idisk.GetIStorage %s", idisk.GetGlobalId()) } storageExtId := istorage.GetGlobalId() storage, err := db.FetchByExternalIdAndManagerId(StorageManager, storageExtId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery { return q.Equals("manager_id", managerId) }) if err != nil { - log.Errorf("failed to found storage by istorage %s error: %v", storageExtId, err) - return err + return errors.Wrapf(err, "storage db.FetchByExternalIdAndManagerId(%s)", storageExtId) } diff, err := db.UpdateWithLock(ctx, self, func() error { self.StorageId = storage.GetId() @@ -1391,8 +1378,7 @@ func (self *SDisk) syncDiskStorage(ctx context.Context, userCred mcclient.TokenC return nil }) if err != nil { - log.Errorf("syncWithCloudDisk error %s", err) - return err + return errors.Wrapf(err, "db.UpdateWithLock") } db.OpsLog.LogSyncUpdate(self, diff, userCred) return nil @@ -1515,8 +1501,7 @@ func (self *SDisk) syncWithCloudDisk(ctx context.Context, userCred mcclient.Toke return nil }) if err != nil { - log.Errorf("syncWithCloudDisk error %s", err) - return err + return errors.Wrapf(err, "db.UpdateWithLock") } // sync disk's snapshotpolicy @@ -1534,7 +1519,13 @@ func (self *SDisk) syncWithCloudDisk(ctx context.Context, userCred mcclient.Toke } db.OpsLog.LogSyncUpdate(self, diff, userCred) - SyncCloudProject(userCred, self, syncOwnerId, extDisk, storage.ManagerId) + syncVirtualResourceMetadata(ctx, userCred, self, extDisk) + + if len(guests) == 0 { + SyncCloudProject(userCred, self, syncOwnerId, extDisk, storage.ManagerId) + } else { + self.SyncCloudProjectId(userCred, guests[0].GetOwnerId()) + } return nil } @@ -1596,6 +1587,8 @@ func (manager *SDiskManager) newFromCloudDisk(ctx context.Context, userCred mccl return nil, err } + syncVirtualResourceMetadata(ctx, userCred, &disk, extDisk) + SyncCloudProject(userCred, &disk, syncOwnerId, extDisk, storage.ManagerId) db.OpsLog.LogEvent(&disk, db.ACT_CREATE, disk.GetShortDesc(ctx), userCred) diff --git a/pkg/compute/models/elasticips.go b/pkg/compute/models/elasticips.go index eef1a0cbd8..92baeb8385 100644 --- a/pkg/compute/models/elasticips.go +++ b/pkg/compute/models/elasticips.go @@ -379,16 +379,14 @@ func (manager *SElasticipManager) SyncEips(ctx context.Context, userCred mcclien if err != nil { syncResult.UpdateError(err) } else { - syncVirtualResourceMetadata(ctx, userCred, &commondb[i], commonext[i]) syncResult.Update() } } for i := 0; i < len(added); i += 1 { - new, err := manager.newFromCloudEip(ctx, userCred, added[i], provider, region, syncOwnerId) + _, err := manager.newFromCloudEip(ctx, userCred, added[i], provider, region, syncOwnerId) if err != nil { syncResult.AddError(err) } else { - syncVirtualResourceMetadata(ctx, userCred, new, added[i]) syncResult.Add() } } @@ -493,8 +491,7 @@ func (self *SElasticip) SyncWithCloudEip(ctx context.Context, userCred mcclient. return nil }) if err != nil { - log.Errorf("SyncWithCloudEip fail %s", err) - return err + return errors.Wrapf(err, "db.UpdateWithLock") } db.OpsLog.LogSyncUpdate(self, diff, userCred) @@ -503,7 +500,12 @@ func (self *SElasticip) SyncWithCloudEip(ctx context.Context, userCred mcclient. return errors.Wrap(err, "fail to sync associated instance of EIP") } syncVirtualResourceMetadata(ctx, userCred, self, ext) - SyncCloudProject(userCred, self, syncOwnerId, ext, self.ManagerId) + + if res := self.GetAssociateResource(); res != nil { + self.SyncCloudProjectId(userCred, res.GetOwnerId()) + } else { + SyncCloudProject(userCred, self, syncOwnerId, ext, self.ManagerId) + } return nil } @@ -556,14 +558,19 @@ func (manager *SElasticipManager) newFromCloudEip(ctx context.Context, userCred return nil, errors.Wrapf(err, "newFromCloudEip") } - syncVirtualResourceMetadata(ctx, userCred, &eip, extEip) - SyncCloudProject(userCred, &eip, syncOwnerId, extEip, eip.ManagerId) - err = eip.SyncInstanceWithCloudEip(ctx, userCred, extEip) if err != nil { return nil, errors.Wrap(err, "fail to sync associated instance of EIP") } + syncVirtualResourceMetadata(ctx, userCred, &eip, extEip) + + if res := eip.GetAssociateResource(); res != nil { + eip.SyncCloudProjectId(userCred, res.GetOwnerId()) + } else { + SyncCloudProject(userCred, &eip, syncOwnerId, extEip, eip.ManagerId) + } + db.OpsLog.LogEvent(&eip, db.ACT_CREATE, eip.GetShortDesc(ctx), userCred) return &eip, nil diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index 480f5057fe..b264b45440 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -2547,6 +2547,7 @@ func (self *SGuest) syncWithCloudVM(ctx context.Context, userCred mcclient.Token db.OpsLog.LogSyncUpdate(self, diff, userCred) + syncVirtualResourceMetadata(ctx, userCred, self, extVM) SyncCloudProject(userCred, self, syncOwnerId, extVM, host.ManagerId) if provider.GetFactory().IsSupportPrepaidResources() && recycle { @@ -3236,9 +3237,9 @@ func (self *SGuest) SyncVMDisks(ctx context.Context, userCred mcclient.TokenCred vdisk := needAdds[i].vdisk err := self.attach2Disk(ctx, needAdds[i].disk, userCred, vdisk.GetDriver(), vdisk.GetCacheMode(), vdisk.GetMountpoint()) if err != nil { - log.Errorf("attach2Disk error: %v", err) - result.AddError(err) + result.AddError(errors.Wrapf(err, "attach2Disk")) } else { + needAdds[i].disk.SyncCloudProjectId(userCred, self.GetOwnerId()) result.Add() } } @@ -4927,13 +4928,11 @@ func (self *SGuest) SyncVMEip(ctx context.Context, userCred mcclient.TokenCreden // add neip, err := ElasticipManager.getEipByExtEip(ctx, userCred, extEip, provider, self.getRegion(), syncOwnerId) if err != nil { - log.Errorf("getEipByExtEip error %v", err) - result.AddError(err) + result.AddError(errors.Wrapf(err, "getEipByExtEip")) } else { err = neip.AssociateInstance(ctx, userCred, api.EIP_ASSOCIATE_TYPE_SERVER, self) if err != nil { - log.Errorf("AssociateVM error %v", err) - result.AddError(err) + result.AddError(errors.Wrapf(err, "neip.AssociateInstance")) } else { result.Add() } diff --git a/pkg/compute/models/networks.go b/pkg/compute/models/networks.go index 9bd5b42355..5883cf1eb9 100644 --- a/pkg/compute/models/networks.go +++ b/pkg/compute/models/networks.go @@ -703,7 +703,6 @@ func (manager *SNetworkManager) SyncNetworks(ctx context.Context, userCred mccli if err != nil { syncResult.UpdateError(err) } else { - syncVirtualResourceMetadata(ctx, userCred, &commondb[i], commonext[i]) localNets = append(localNets, commondb[i]) remoteNets = append(remoteNets, commonext[i]) syncResult.Update() @@ -714,7 +713,6 @@ func (manager *SNetworkManager) SyncNetworks(ctx context.Context, userCred mccli if err != nil { syncResult.AddError(err) } else { - syncVirtualResourceMetadata(ctx, userCred, new, added[i]) localNets = append(localNets, *new) remoteNets = append(remoteNets, added[i]) syncResult.Add() diff --git a/pkg/compute/models/snapshots.go b/pkg/compute/models/snapshots.go index a09d54873b..dc1cedaf83 100644 --- a/pkg/compute/models/snapshots.go +++ b/pkg/compute/models/snapshots.go @@ -893,12 +893,11 @@ func (self *SSnapshot) SyncWithCloudSnapshot(ctx context.Context, userCred mccli syncVirtualResourceMetadata(ctx, userCred, self, ext) // bugfix for now: - disk, err := self.GetDisk() - if err != nil && err != sql.ErrNoRows { - return errors.Wrapf(err, "get disk of snapshot %s error", self.Id) - } - if err == nil { + disk, _ := self.GetDisk() + if disk != nil { self.SyncCloudProjectId(userCred, disk.GetOwnerId()) + } else { + SyncCloudProject(userCred, self, syncOwnerId, ext, disk.GetCloudprovider().Id) } return nil @@ -1005,16 +1004,14 @@ func (manager *SSnapshotManager) SyncSnapshots(ctx context.Context, userCred mcc if err != nil { syncResult.UpdateError(err) } else { - syncVirtualResourceMetadata(ctx, userCred, &commondb[i], commonext[i]) syncResult.Update() } } for i := 0; i < len(added); i += 1 { - local, err := manager.newFromCloudSnapshot(ctx, userCred, added[i], region, syncOwnerId, provider) + _, err := manager.newFromCloudSnapshot(ctx, userCred, added[i], region, syncOwnerId, provider) if err != nil { syncResult.AddError(err) } else { - syncVirtualResourceMetadata(ctx, userCred, local, added[i]) syncResult.Add() } } diff --git a/pkg/multicloud/aws/network.go b/pkg/multicloud/aws/network.go index d102bf138d..3376c5e77f 100644 --- a/pkg/multicloud/aws/network.go +++ b/pkg/multicloud/aws/network.go @@ -79,7 +79,6 @@ func (self *SNetwork) GetStatus() string { } func (self *SNetwork) Refresh() error { - log.Debugf("network refresh %s", self.NetworkId) new, err := self.wire.zone.region.getNetwork(self.NetworkId) if err != nil { return err