Merge pull request #7005 from ioito/automated-cherry-pick-of-#6998-upstream-release-3.2

Automated cherry pick of #6998: fix: manager_id资源隔离获取
This commit is contained in:
Zexi Li
2020-07-02 20:20:23 +08:00
committed by GitHub
31 changed files with 310 additions and 121 deletions
+7
View File
@@ -66,7 +66,14 @@ func SetExternalId(model IExternalizedModel, userCred mcclient.TokenCredential,
}
func FetchByExternalId(manager IModelManager, idStr string) (IExternalizedModel, error) {
return FetchByExternalIdAndManagerId(manager, idStr, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q
})
}
func FetchByExternalIdAndManagerId(manager IModelManager, idStr string, filter func(q *sqlchemy.SQuery) *sqlchemy.SQuery) (IExternalizedModel, error) {
q := manager.Query().Equals("external_id", idStr)
q = filter(q)
count, err := q.CountWithError()
if err != nil {
return nil, err
+11 -2
View File
@@ -27,6 +27,7 @@ import (
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/util/osprofile"
"yunion.io/x/pkg/utils"
"yunion.io/x/sqlchemy"
billing_api "yunion.io/x/onecloud/pkg/apis/billing"
api "yunion.io/x/onecloud/pkg/apis/compute"
@@ -449,7 +450,9 @@ func (self *SManagedVirtualizedGuestDriver) RemoteDeployGuestForCreate(ctx conte
}
if hostId := iVM.GetIHostId(); len(hostId) > 0 {
host, err := db.FetchByExternalId(models.HostManager, hostId)
host, err := db.FetchByExternalIdAndManagerId(models.HostManager, hostId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", host.ManagerId)
})
if err != nil {
log.Warningf("failed to found new hostId(%s) for ivm %s(%s) error: %v", hostId, guest.Name, guest.Id, err)
} else if host.GetId() != guest.HostId {
@@ -886,7 +889,13 @@ func (self *SManagedVirtualizedGuestDriver) OnGuestDeployTaskDataReceived(ctx co
}
if len(diskInfo[i].StorageExternalId) > 0 {
storage, err := db.FetchByExternalId(models.StorageManager, diskInfo[i].StorageExternalId)
storage, err := db.FetchByExternalIdAndManagerId(models.StorageManager, diskInfo[i].StorageExternalId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
host := guest.GetHost()
if host != nil {
return q.Equals("manager_id", host.ManagerId)
}
return q
})
if err != nil {
log.Warningf("failed to found storage by externalId %s error: %v", diskInfo[i].StorageExternalId, err)
} else if disk.StorageId != storage.GetId() {
+1 -1
View File
@@ -122,7 +122,7 @@ func (self *SCloudproviderregion) GetAccount() *SCloudaccount {
func (self *SCloudproviderregion) GetRegion() *SCloudregion {
regionObj, err := CloudregionManager.FetchById(self.CloudregionId)
if err != nil {
log.Errorf("CloudproviderManager.FetchById fail %s", err)
log.Errorf("CloudregionManager.FetchById(%s) fail %s", self.CloudregionId, err)
return nil
}
return regionObj.(*SCloudregion)
+18 -5
View File
@@ -226,15 +226,28 @@ func (self *SCloudregion) getGuestCountInternal(increment bool) (int, error) {
return query.CountWithError()
}
func (self *SCloudregion) GetVpcCount() (int, error) {
func (self *SCloudregion) GetVpcQuery() *sqlchemy.SQuery {
vpcs := VpcManager.Query()
if self.Id == api.DEFAULT_REGION_ID {
return vpcs.Filter(sqlchemy.OR(sqlchemy.IsNull(vpcs.Field("cloudregion_id")),
sqlchemy.IsEmpty(vpcs.Field("cloudregion_id")),
sqlchemy.Equals(vpcs.Field("cloudregion_id"), self.Id))).CountWithError()
} else {
return vpcs.Equals("cloudregion_id", self.Id).CountWithError()
sqlchemy.Equals(vpcs.Field("cloudregion_id"), self.Id)))
}
return vpcs.Equals("cloudregion_id", self.Id)
}
func (self *SCloudregion) GetVpcCount() (int, error) {
return self.GetVpcQuery().CountWithError()
}
func (self *SCloudregion) GetVpcs() ([]SVpc, error) {
vpcs := []SVpc{}
q := self.GetVpcQuery()
err := db.FetchModelObjects(VpcManager, q, &vpcs)
if err != nil {
return nil, errors.Wrap(err, "db.FetchModelObjects")
}
return vpcs, nil
}
func (self *SCloudregion) GetDriver() IRegionDriver {
@@ -497,7 +510,7 @@ func (self *SCloudregion) AllowPerformDefaultVpc(ctx context.Context, userCred m
}
func (self *SCloudregion) PerformDefaultVpc(ctx context.Context, userCred mcclient.TokenCredential, query jsonutils.JSONObject, data jsonutils.JSONObject) (jsonutils.JSONObject, error) {
vpcs, err := VpcManager.getVpcsByRegion(self, nil)
vpcs, err := self.GetVpcs()
if err != nil {
return nil, err
}
+6 -2
View File
@@ -437,7 +437,9 @@ func (self *SDBInstanceBackup) SyncWithCloudDBInstanceBackup(
if dbinstanceId := extBackup.GetDBInstanceId(); len(dbinstanceId) > 0 {
//有可能云上删除了实例,未删除备份
_instance, err := db.FetchByExternalId(DBInstanceManager, dbinstanceId)
_instance, err := db.FetchByExternalIdAndManagerId(DBInstanceManager, dbinstanceId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", provider.Id)
})
if err == sql.ErrNoRows {
self.DBInstanceId = ""
}
@@ -493,7 +495,9 @@ func (manager *SDBInstanceBackupManager) newFromCloudDBInstanceBackup(
backup.ExternalId = extBackup.GetGlobalId()
if dbinstanceId := extBackup.GetDBInstanceId(); len(dbinstanceId) > 0 {
_dbinstance, err := db.FetchByExternalId(DBInstanceManager, dbinstanceId)
_dbinstance, err := db.FetchByExternalIdAndManagerId(DBInstanceManager, dbinstanceId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", provider.Id)
})
if err != nil {
log.Warningf("failed to found dbinstance for backup %s by externalId: %s error: %v", backup.Name, dbinstanceId, err)
} else {
+17 -4
View File
@@ -23,6 +23,7 @@ import (
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/util/compare"
"yunion.io/x/pkg/util/netutils"
"yunion.io/x/sqlchemy"
api "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
@@ -162,9 +163,15 @@ func (manager *SDBInstanceNetworkManager) SyncDBInstanceNetwork(ctx context.Cont
func (self *SDBInstanceNetwork) syncWithCloudDBNetwork(ctx context.Context, userCred mcclient.TokenCredential, dbinstance *SDBInstance, network *cloudprovider.SDBInstanceNetwork) error {
_, err := db.UpdateWithLock(ctx, self, func() error {
_localnetwork, err := db.FetchByExternalId(NetworkManager, network.NetworkId)
_localnetwork, err := db.FetchByExternalIdAndManagerId(NetworkManager, network.NetworkId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
wire := WireManager.Query().SubQuery()
vpc := VpcManager.Query().SubQuery()
return q.Join(wire, sqlchemy.Equals(wire.Field("id"), q.Field("wire_id"))).
Join(vpc, sqlchemy.Equals(vpc.Field("id"), wire.Field("vpc_id"))).
Filter(sqlchemy.Equals(vpc.Field("manager_id"), dbinstance.ManagerId))
})
if err != nil {
return errors.Wrapf(err, "FetchByExternalId")
return errors.Wrapf(err, "FetchByExternalIdAndManagerId")
}
localnetwork := _localnetwork.(*SNetwork)
self.NetworkId = localnetwork.Id
@@ -194,9 +201,15 @@ func (manager *SDBInstanceNetworkManager) newFromCloudDBNetwork(ctx context.Cont
dbNetwork.SetModelManager(manager, &dbNetwork)
dbNetwork.DBInstanceId = dbinstance.Id
_localnetwork, err := db.FetchByExternalId(NetworkManager, network.NetworkId)
_localnetwork, err := db.FetchByExternalIdAndManagerId(NetworkManager, network.NetworkId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
wire := WireManager.Query().SubQuery()
vpc := VpcManager.Query().SubQuery()
return q.Join(wire, sqlchemy.Equals(wire.Field("id"), q.Field("wire_id"))).
Join(vpc, sqlchemy.Equals(vpc.Field("id"), wire.Field("vpc_id"))).
Filter(sqlchemy.Equals(vpc.Field("manager_id"), dbinstance.ManagerId))
})
if err != nil {
return errors.Wrapf(err, "newFromCloudDBNetwork.FetchByExternalId")
return errors.Wrapf(err, "newFromCloudDBNetwork.FetchByExternalIdAndManagerId")
}
localnetwork := _localnetwork.(*SNetwork)
+12 -4
View File
@@ -1236,12 +1236,16 @@ func (manager *SDBInstanceManager) SyncDBInstanceMasterId(ctx context.Context, u
for _, instance := range cloudDBInstances {
masterId := instance.GetMasterInstanceId()
if len(masterId) > 0 {
master, err := db.FetchByExternalId(manager, masterId)
master, err := db.FetchByExternalIdAndManagerId(manager, masterId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", provider.Id)
})
if err != nil {
log.Errorf("failed to found master dbinstance by externalId: %s error: %v", masterId, err)
continue
}
slave, err := db.FetchByExternalId(manager, instance.GetGlobalId())
slave, err := db.FetchByExternalIdAndManagerId(manager, instance.GetGlobalId(), func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", provider.Id)
})
if err != nil {
log.Errorf("failed to found local dbinstance by externalId %s error: %v", instance.GetGlobalId(), err)
continue
@@ -1495,7 +1499,9 @@ func (self *SDBInstance) SyncWithCloudDBInstance(ctx context.Context, userCred m
if len(self.VpcId) == 0 {
if vpcId := extInstance.GetIVpcId(); len(vpcId) > 0 {
vpc, err := db.FetchByExternalId(VpcManager, vpcId)
vpc, err := db.FetchByExternalIdAndManagerId(VpcManager, vpcId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", provider.Id)
})
if err != nil {
return errors.Wrapf(err, "SyncWithCloudDBInstance.FetchVpcId")
}
@@ -1580,7 +1586,9 @@ func (manager *SDBInstanceManager) newFromCloudDBInstance(ctx context.Context, u
}
if vpcId := extInstance.GetIVpcId(); len(vpcId) > 0 {
vpc, err := db.FetchByExternalId(VpcManager, vpcId)
vpc, err := db.FetchByExternalIdAndManagerId(VpcManager, vpcId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", provider.Id)
})
if err != nil {
return nil, errors.Wrapf(err, "newFromCloudDBInstance.FetchVpcId")
}
+26 -11
View File
@@ -1217,13 +1217,16 @@ func (manager *SDiskManager) getDisksByStorage(storage *SStorage) ([]SDisk, erro
return disks, nil
}
func (manager *SDiskManager) syncCloudDisk(ctx context.Context, userCred mcclient.TokenCredential, provider cloudprovider.ICloudProvider, vdisk cloudprovider.ICloudDisk, index int, syncOwnerId mcclient.IIdentityProvider) (*SDisk, error) {
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.FetchByExternalId(manager, vdisk.GetGlobalId())
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()
@@ -1231,7 +1234,9 @@ func (manager *SDiskManager) syncCloudDisk(ctx context.Context, userCred mcclien
return nil, errors.Wrapf(err, "unable to GetIStorage of vdisk %q", vdisk.GetName())
}
storageObj, err := db.FetchByExternalId(StorageManager, vstorage.GetGlobalId())
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
@@ -1243,7 +1248,7 @@ func (manager *SDiskManager) syncCloudDisk(ctx context.Context, userCred mcclien
}
} else {
disk := diskObj.(*SDisk)
err = disk.syncWithCloudDisk(ctx, userCred, provider, vdisk, index, syncOwnerId)
err = disk.syncWithCloudDisk(ctx, userCred, provider, vdisk, index, syncOwnerId, managerId)
if err != nil {
return nil, err
}
@@ -1288,7 +1293,7 @@ func (manager *SDiskManager) SyncDisks(ctx context.Context, userCred mcclient.To
}
for i := 0; i < len(commondb); i += 1 {
err = commondb[i].syncWithCloudDisk(ctx, userCred, provider, commonext[i], -1, syncOwnerId)
err = commondb[i].syncWithCloudDisk(ctx, userCred, provider, commonext[i], -1, syncOwnerId, storage.ManagerId)
if err != nil {
syncResult.UpdateError(err)
} else {
@@ -1301,7 +1306,10 @@ func (manager *SDiskManager) SyncDisks(ctx context.Context, userCred mcclient.To
for i := 0; i < len(added); i += 1 {
extId := added[i].GetGlobalId()
_disk, err := db.FetchByExternalId(manager, extId)
_disk, err := db.FetchByExternalIdAndManagerId(manager, extId, 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"), storage.ManagerId))
})
if err != nil && err != sql.ErrNoRows {
//主要是显示duplicate err及 general err,方便排错
msg := fmt.Errorf("failed to found disk by external Id %s error: %v", extId, err)
@@ -1310,7 +1318,7 @@ func (manager *SDiskManager) SyncDisks(ctx context.Context, userCred mcclient.To
}
if _disk != nil {
disk := _disk.(*SDisk)
err = disk.syncDiskStorage(ctx, userCred, added[i])
err = disk.syncDiskStorage(ctx, userCred, added[i], storage.ManagerId)
if err != nil {
syncResult.UpdateError(err)
} else {
@@ -1332,7 +1340,7 @@ func (manager *SDiskManager) SyncDisks(ctx context.Context, userCred mcclient.To
return localDisks, remoteDisks, syncResult
}
func (self *SDisk) syncDiskStorage(ctx context.Context, userCred mcclient.TokenCredential, idisk cloudprovider.ICloudDisk) error {
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 {
@@ -1340,7 +1348,9 @@ func (self *SDisk) syncDiskStorage(ctx context.Context, userCred mcclient.TokenC
return err
}
storageExtId := istorage.GetGlobalId()
storage, err := db.FetchByExternalId(StorageManager, storageExtId)
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
@@ -1392,7 +1402,12 @@ func (self *SDisk) syncRemoveCloudDisk(ctx context.Context, userCred mcclient.To
iDisk, err := iregion.GetIDiskById(self.ExternalId)
if err == nil {
if storageId := iDisk.GetIStorageId(); len(storageId) > 0 {
storage, err := db.FetchByExternalId(StorageManager, storageId)
storage, err := db.FetchByExternalIdAndManagerId(StorageManager, storageId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
if s := self.GetStorage(); s != nil {
return q.Equals("manager_id", s.ManagerId)
}
return q
})
if err == nil {
_, err = db.Update(self, func() error {
self.StorageId = storage.GetId()
@@ -1418,7 +1433,7 @@ func (self *SDisk) syncRemoveCloudDisk(ctx context.Context, userCred mcclient.To
return self.RealDelete(ctx, userCred)
}
func (self *SDisk) syncWithCloudDisk(ctx context.Context, userCred mcclient.TokenCredential, provider cloudprovider.ICloudProvider, extDisk cloudprovider.ICloudDisk, index int, syncOwnerId mcclient.IIdentityProvider) error {
func (self *SDisk) syncWithCloudDisk(ctx context.Context, userCred mcclient.TokenCredential, provider cloudprovider.ICloudProvider, extDisk cloudprovider.ICloudDisk, index int, syncOwnerId mcclient.IIdentityProvider, managerId string) error {
recycle := false
guests := self.GetGuests()
if provider.GetFactory().IsSupportPrepaidResources() && len(guests) == 1 && guests[0].IsPrepaidRecycle() {
+10 -2
View File
@@ -582,7 +582,9 @@ func (manager *SElasticcacheManager) newFromCloudElasticcache(ctx context.Contex
}
if vpcId := extInstance.GetVpcId(); len(vpcId) > 0 {
vpc, err := db.FetchByExternalId(VpcManager, vpcId)
vpc, err := db.FetchByExternalIdAndManagerId(VpcManager, vpcId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", provider.Id)
})
if err != nil {
return nil, errors.Wrapf(err, "newFromCloudElasticcache.FetchVpcId")
}
@@ -590,7 +592,13 @@ func (manager *SElasticcacheManager) newFromCloudElasticcache(ctx context.Contex
}
if networkId := extInstance.GetNetworkId(); len(networkId) > 0 {
network, err := db.FetchByExternalId(NetworkManager, networkId)
network, err := db.FetchByExternalIdAndManagerId(NetworkManager, networkId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
wire := WireManager.Query().SubQuery()
vpc := VpcManager.Query().SubQuery()
return q.Join(wire, sqlchemy.Equals(wire.Field("id"), q.Field("wire_id"))).
Join(vpc, sqlchemy.Equals(vpc.Field("id"), wire.Field("vpc_id"))).
Filter(sqlchemy.Equals(vpc.Field("manager_id"), provider.Id))
})
if err != nil {
return nil, errors.Wrapf(err, "newFromCloudElasticcache.FetchNetworkId")
}
+20 -3
View File
@@ -414,7 +414,16 @@ func (self *SElasticip) SyncInstanceWithCloudEip(ctx context.Context, userCred m
return errors.Error("unsupported association type")
}
extRes, err := db.FetchByExternalId(manager, vmExtId)
extRes, err := db.FetchByExternalIdAndManagerId(manager, vmExtId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
switch ext.GetAssociationType() {
case api.EIP_ASSOCIATE_TYPE_SERVER:
sq := HostManager.Query().SubQuery()
return q.Join(sq, sqlchemy.Equals(sq.Field("id"), q.Field("host_id"))).Filter(sqlchemy.Equals(sq.Field("manager_id"), self.ManagerId))
case api.EIP_ASSOCIATE_TYPE_NAT_GATEWAY, api.EIP_ASSOCIATE_TYPE_LOADBALANCER:
return q.Equals("manager_id", self.ManagerId)
}
return q
})
if err != nil {
log.Errorf("fail to find vm by external ID %s", vmExtId)
return err
@@ -501,7 +510,13 @@ func (manager *SElasticipManager) newFromCloudEip(ctx context.Context, userCred
eip.CloudregionId = region.Id
eip.ChargeType = extEip.GetInternetChargeType()
if networkId := extEip.GetINetworkId(); len(networkId) > 0 {
network, err := db.FetchByExternalId(NetworkManager, networkId)
network, err := db.FetchByExternalIdAndManagerId(NetworkManager, networkId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
wire := WireManager.Query().SubQuery()
vpc := VpcManager.Query().SubQuery()
return q.Join(wire, sqlchemy.Equals(wire.Field("id"), q.Field("wire_id"))).
Join(vpc, sqlchemy.Equals(vpc.Field("id"), wire.Field("vpc_id"))).
Filter(sqlchemy.Equals(vpc.Field("manager_id"), provider.Id))
})
if err != nil {
msg := fmt.Sprintf("failed to found network by externalId %s error: %v", networkId, err)
log.Errorf(msg)
@@ -759,7 +774,9 @@ func (self *SElasticip) AssociateNatGateway(ctx context.Context, userCred mcclie
}
func (manager *SElasticipManager) getEipByExtEip(ctx context.Context, userCred mcclient.TokenCredential, extEip cloudprovider.ICloudEIP, provider *SCloudprovider, region *SCloudregion, syncOwnerId mcclient.IIdentityProvider) (*SElasticip, error) {
eipObj, err := db.FetchByExternalId(manager, extEip.GetGlobalId())
eipObj, err := db.FetchByExternalIdAndManagerId(manager, extEip.GetGlobalId(), func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", provider.Id)
})
if err == nil {
return eipObj.(*SElasticip), nil
}
+15 -3
View File
@@ -2292,7 +2292,13 @@ func (self *SGuest) syncRemoveCloudVM(ctx context.Context, userCred mcclient.Tok
iVM, err := iregion.GetIVMById(self.ExternalId)
if err == nil { //漂移归位
if hostId := iVM.GetIHostId(); len(hostId) > 0 {
host, err := db.FetchByExternalId(HostManager, hostId)
host, err := db.FetchByExternalIdAndManagerId(HostManager, hostId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
host := self.GetHost()
if host != nil {
return q.Equals("manager_id", host.ManagerId)
}
return q
})
if err == nil {
_, err = db.Update(self, func() error {
self.HostId = host.GetId()
@@ -2723,7 +2729,13 @@ func getCloudNicNetwork(vnic cloudprovider.ICloudNic, host *SHost, ipList []stri
// find network by IP
return host.getNetworkOfIPOnHost(ip)
}
localNetObj, err := db.FetchByExternalId(NetworkManager, vnet.GetGlobalId())
localNetObj, err := db.FetchByExternalIdAndManagerId(NetworkManager, vnet.GetGlobalId(), func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
vpc := VpcManager.Query().SubQuery()
wire := WireManager.Query().SubQuery()
return q.Join(wire, sqlchemy.Equals(q.Field("wire_id"), wire.Field("id"))).
Join(vpc, sqlchemy.Equals(wire.Field("vpc_id"), vpc.Field("id"))).
Filter(sqlchemy.Equals(vpc.Field("manager_id"), host.ManagerId))
})
if err != nil {
return nil, fmt.Errorf("Cannot find network of external_id %s: %v", vnet.GetGlobalId(), err)
}
@@ -2923,7 +2935,7 @@ func (self *SGuest) SyncVMDisks(ctx context.Context, userCred mcclient.TokenCred
if len(vdisks[i].GetGlobalId()) == 0 {
continue
}
disk, err := DiskManager.syncCloudDisk(ctx, userCred, provider, vdisks[i], i, syncOwnerId)
disk, err := DiskManager.syncCloudDisk(ctx, userCred, provider, vdisks[i], i, syncOwnerId, host.ManagerId)
if err != nil {
log.Errorf("syncCloudDisk error: %v", err)
result.Error(err)
+1 -1
View File
@@ -686,7 +686,7 @@ func (host *SHost) RebuildRecycledGuest(ctx context.Context, userCred mcclient.T
log.Errorf("disk.SetExternalId fail %s", err)
return err
}
err = disk.syncWithCloudDisk(ctx, userCred, iprovider, idisks[i], i, guest.GetOwnerId())
err = disk.syncWithCloudDisk(ctx, userCred, iprovider, idisks[i], i, guest.GetOwnerId(), host.ManagerId)
if err != nil {
log.Errorf("disk.syncWithCloudDisk fail %s", err)
return err
+14 -4
View File
@@ -1893,7 +1893,9 @@ func (self *SHost) Attach2Storage(ctx context.Context, userCred mcclient.TokenCr
}
func (self *SHost) newCloudHostStorage(ctx context.Context, userCred mcclient.TokenCredential, extStorage cloudprovider.ICloudStorage, provider *SCloudprovider) (*SStorage, error) {
storageObj, err := db.FetchByExternalId(StorageManager, extStorage.GetGlobalId())
storageObj, err := db.FetchByExternalIdAndManagerId(StorageManager, extStorage.GetGlobalId(), func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", provider.Id)
})
if err != nil {
if err == sql.ErrNoRows {
// no cloud storage found, this may happen for on-premise host
@@ -2008,7 +2010,10 @@ func (self *SHost) Attach2Wire(ctx context.Context, userCred mcclient.TokenCrede
}
func (self *SHost) newCloudHostWire(ctx context.Context, userCred mcclient.TokenCredential, extWire cloudprovider.ICloudWire) error {
wireObj, err := db.FetchByExternalId(WireManager, extWire.GetGlobalId())
wireObj, err := db.FetchByExternalIdAndManagerId(WireManager, extWire.GetGlobalId(), func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
sq := VpcManager.Query().SubQuery()
return q.Join(sq, sqlchemy.Equals(sq.Field("id"), q.Field("vpc_id"))).Filter(sqlchemy.Equals(sq.Field("manager_id"), self.ManagerId))
})
if err != nil {
log.Errorf("%s", err)
return nil
@@ -2076,7 +2081,10 @@ func (self *SHost) SyncHostVMs(ctx context.Context, userCred mcclient.TokenCrede
}
for i := 0; i < len(added); i += 1 {
vm, err := db.FetchByExternalId(GuestManager, added[i].GetGlobalId())
vm, err := db.FetchByExternalIdAndManagerId(GuestManager, added[i].GetGlobalId(), func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
sq := HostManager.Query().SubQuery()
return q.Join(sq, sqlchemy.Equals(sq.Field("id"), q.Field("host_id"))).Filter(sqlchemy.Equals(sq.Field("manager_id"), self.ManagerId))
})
if err != nil && err != sql.ErrNoRows {
log.Errorf("failed to found guest by externalId %s error: %v", added[i].GetGlobalId(), err)
continue
@@ -2088,7 +2096,9 @@ func (self *SHost) SyncHostVMs(ctx context.Context, userCred mcclient.TokenCrede
log.Errorf("failed to found ihost from vm %s", added[i].GetGlobalId())
continue
}
_host, err := db.FetchByExternalId(HostManager, ihost.GetGlobalId())
_host, err := db.FetchByExternalIdAndManagerId(HostManager, ihost.GetGlobalId(), func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", self.ManagerId)
})
if err != nil {
log.Errorf("failed to found host by externalId %s", ihost.GetGlobalId())
continue
@@ -20,6 +20,7 @@ import (
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/util/compare"
"yunion.io/x/sqlchemy"
api "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
@@ -194,7 +195,10 @@ func (lbb *SAwsCachedLb) syncRemoveCloudLoadbalancerBackend(ctx context.Context,
func (lbb *SAwsCachedLb) constructFieldsFromCloudLoadbalancerBackend(extLoadbalancerBackend cloudprovider.ICloudLoadbalancerBackend) error {
lbb.Status = extLoadbalancerBackend.GetStatus()
instance, err := db.FetchByExternalId(GuestManager, extLoadbalancerBackend.GetBackendId())
instance, err := db.FetchByExternalIdAndManagerId(GuestManager, extLoadbalancerBackend.GetBackendId(), func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
sq := HostManager.Query().SubQuery()
return q.Join(sq, sqlchemy.Equals(q.Field("host_id"), sq.Field("id"))).Filter(sqlchemy.Equals(sq.Field("manager_id"), lbb.ManagerId))
})
if err != nil {
return err
}
@@ -232,7 +232,9 @@ func (man *SAwsCachedLbbgManager) SyncLoadbalancerBackendgroups(ctx context.Cont
var elb *SLoadbalancer
elbId := commonext[i].GetLoadbalancerId()
if len(elbId) > 0 {
ielb, err := db.FetchByExternalId(LoadbalancerManager, elbId)
ielb, err := db.FetchByExternalIdAndManagerId(LoadbalancerManager, elbId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", provider.Id)
})
if err == nil {
elb = ielb.(*SLoadbalancer)
}
+7 -4
View File
@@ -522,7 +522,7 @@ func (man *SLoadbalancerBackendManager) SyncLoadbalancerBackends(ctx context.Con
return syncResult
}
func (lbb *SLoadbalancerBackend) constructFieldsFromCloudLoadbalancerBackend(extLoadbalancerBackend cloudprovider.ICloudLoadbalancerBackend) error {
func (lbb *SLoadbalancerBackend) constructFieldsFromCloudLoadbalancerBackend(extLoadbalancerBackend cloudprovider.ICloudLoadbalancerBackend, managerId string) error {
// lbb.Name = extLoadbalancerBackend.GetName()
lbb.Status = extLoadbalancerBackend.GetStatus()
@@ -532,7 +532,10 @@ func (lbb *SLoadbalancerBackend) constructFieldsFromCloudLoadbalancerBackend(ext
lbb.BackendType = extLoadbalancerBackend.GetBackendType()
lbb.BackendRole = extLoadbalancerBackend.GetBackendRole()
instance, err := db.FetchByExternalId(GuestManager, extLoadbalancerBackend.GetBackendId())
instance, err := db.FetchByExternalIdAndManagerId(GuestManager, extLoadbalancerBackend.GetBackendId(), func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
sq := HostManager.Query().SubQuery()
return q.Join(sq, sqlchemy.Equals(sq.Field("id"), q.Field("host_id"))).Filter(sqlchemy.Equals(sq.Field("manager_id"), managerId))
})
if err != nil {
return err
}
@@ -563,7 +566,7 @@ func (lbb *SLoadbalancerBackend) syncRemoveCloudLoadbalancerBackend(ctx context.
func (lbb *SLoadbalancerBackend) SyncWithCloudLoadbalancerBackend(ctx context.Context, userCred mcclient.TokenCredential, extLoadbalancerBackend cloudprovider.ICloudLoadbalancerBackend, syncOwnerId mcclient.IIdentityProvider, provider *SCloudprovider) error {
diff, err := db.UpdateWithLock(ctx, lbb, func() error {
return lbb.constructFieldsFromCloudLoadbalancerBackend(extLoadbalancerBackend)
return lbb.constructFieldsFromCloudLoadbalancerBackend(extLoadbalancerBackend, provider.Id)
})
if err != nil {
return err
@@ -613,7 +616,7 @@ func (man *SLoadbalancerBackendManager) newFromCloudLoadbalancerBackend(ctx cont
}
lbb.Name = newName
if err := lbb.constructFieldsFromCloudLoadbalancerBackend(extLoadbalancerBackend); err != nil {
if err := lbb.constructFieldsFromCloudLoadbalancerBackend(extLoadbalancerBackend, provider.Id); err != nil {
return nil, err
}
+4 -1
View File
@@ -324,7 +324,10 @@ func (acl *SCachedLoadbalancerAcl) SyncWithCloudLoadbalancerAcl(ctx context.Cont
} else {
ext_listener_id := extAcl.GetAclListenerID()
if len(ext_listener_id) > 0 {
ilistener, err := db.FetchByExternalId(LoadbalancerListenerManager, ext_listener_id)
ilistener, err := db.FetchByExternalIdAndManagerId(LoadbalancerListenerManager, ext_listener_id, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
sq := LoadbalancerManager.Query().SubQuery()
return q.Join(sq, sqlchemy.Equals(sq.Field("id"), q.Field("loadbalancer_id"))).Filter(sqlchemy.Equals(sq.Field("manager_id"), acl.ManagerId))
})
if err != nil {
return errors.Wrap(err, "cacheLoadbalancerAcl.sync.FetchByExternalId")
}
@@ -22,6 +22,7 @@ import (
"yunion.io/x/log"
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/util/compare"
"yunion.io/x/sqlchemy"
api "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
@@ -197,7 +198,10 @@ func (lbb *SHuaweiCachedLb) syncRemoveCloudLoadbalancerBackend(ctx context.Conte
func (lbb *SHuaweiCachedLb) constructFieldsFromCloudLoadbalancerBackend(extLoadbalancerBackend cloudprovider.ICloudLoadbalancerBackend) error {
lbb.Status = extLoadbalancerBackend.GetStatus()
instance, err := db.FetchByExternalId(GuestManager, extLoadbalancerBackend.GetBackendId())
instance, err := db.FetchByExternalIdAndManagerId(GuestManager, extLoadbalancerBackend.GetBackendId(), func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
sq := HostManager.Query().SubQuery()
return q.Join(sq, sqlchemy.Equals(sq.Field("id"), q.Field("host_id"))).Filter(sqlchemy.Equals(sq.Field("manager_id"), lbb.ManagerId))
})
if err != nil {
return err
}
@@ -286,17 +290,6 @@ func (man *SHuaweiCachedLbManager) newFromCloudLoadbalancerBackend(ctx context.C
}
func newLocalBackendFromCloudLoadbalancerBackend(ctx context.Context, userCred mcclient.TokenCredential, loadbalancerBackendgroup *SLoadbalancerBackendGroup, extLoadbalancerBackend cloudprovider.ICloudLoadbalancerBackend, syncOwnerId mcclient.IIdentityProvider) (*SLoadbalancerBackend, error) {
instance, err := db.FetchByExternalId(GuestManager, extLoadbalancerBackend.GetBackendId())
if err != nil {
return nil, err
}
guest := instance.(*SGuest)
//address, err := LoadbalancerBackendManager.GetGuestAddress(guest)
//if err != nil {
// return nil, err
//}
lbbgRegion := loadbalancerBackendgroup.GetRegion()
if lbbgRegion == nil {
return nil, errors.Wrap(httperrors.ErrInvalidStatus, "loadbalancerBackendgroup is not attached to any region")
@@ -306,6 +299,17 @@ func newLocalBackendFromCloudLoadbalancerBackend(ctx context.Context, userCred m
return nil, errors.Wrap(httperrors.ErrInvalidStatus, "loadbalancerBackendgroup is not attached to any cloudprovider")
}
instance, err := db.FetchByExternalIdAndManagerId(GuestManager, extLoadbalancerBackend.GetBackendId(), func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
sq := HostManager.Query().SubQuery()
return q.Join(sq, sqlchemy.Equals(sq.Field("id"), q.Field("host_id"))).Filter(sqlchemy.Equals(sq.Field("manager_id"), lbbgProvider.Id))
})
guest := instance.(*SGuest)
//address, err := LoadbalancerBackendManager.GetGuestAddress(guest)
//if err != nil {
// return nil, err
//}
man := LoadbalancerBackendManager
q := man.Query().IsFalse("pending_deleted")
q = q.Equals("weight", extLoadbalancerBackend.GetWeight()).Equals("port", extLoadbalancerBackend.GetPort())
@@ -350,7 +354,7 @@ func newLocalBackendFromCloudLoadbalancerBackend(ctx context.Context, userCred m
}
lbb.Name = newName
if err := lbb.constructFieldsFromCloudLoadbalancerBackend(extLoadbalancerBackend); err != nil {
if err := lbb.constructFieldsFromCloudLoadbalancerBackend(extLoadbalancerBackend, lbbgProvider.Id); err != nil {
return nil, err
}
+32 -12
View File
@@ -955,7 +955,9 @@ func (lblis *SLoadbalancerListener) constructFieldsFromCloudListener(userCred mc
lblis.AclType = extListener.GetAclType()
if aclID := extListener.GetAclId(); len(aclID) > 0 {
if _acl, err := db.FetchByExternalId(CachedLoadbalancerAclManager, aclID); err == nil {
if _acl, err := db.FetchByExternalIdAndManagerId(CachedLoadbalancerAclManager, aclID, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", lb.ManagerId)
}); err == nil {
acl := _acl.(*SCachedLoadbalancerAcl)
lblis.CachedAclId = acl.GetId()
lblis.AclId = acl.AclId
@@ -991,7 +993,9 @@ func (lblis *SLoadbalancerListener) constructFieldsFromCloudListener(userCred mc
lblis.TLSCipherPolicy = extListener.GetTLSCipherPolicy()
lblis.EnableHttp2 = extListener.HTTP2Enabled()
if certificateId := extListener.GetCertificateId(); len(certificateId) > 0 {
if _cert, err := db.FetchByExternalId(CachedLoadbalancerCertificateManager, certificateId); err == nil {
if _cert, err := db.FetchByExternalIdAndManagerId(CachedLoadbalancerCertificateManager, certificateId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", lb.ManagerId)
}); err == nil {
cert := _cert.(*SCachedLoadbalancerCertificate)
lblis.CachedCertificateId = cert.GetId()
lblis.CertificateId = cert.CertificateId
@@ -1025,7 +1029,9 @@ func (lblis *SLoadbalancerListener) constructFieldsFromCloudListener(userCred mc
switch lblis.GetProviderName() {
case api.CLOUD_PROVIDER_HUAWEI:
if len(groupId) > 0 {
group, err := db.FetchByExternalId(HuaweiCachedLbbgManager, groupId)
group, err := db.FetchByExternalIdAndManagerId(HuaweiCachedLbbgManager, groupId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", lb.ManagerId)
})
if err != nil {
if err == sql.ErrNoRows {
lblis.BackendGroupId = ""
@@ -1037,7 +1043,9 @@ func (lblis *SLoadbalancerListener) constructFieldsFromCloudListener(userCred mc
}
case api.CLOUD_PROVIDER_AWS:
if len(groupId) > 0 {
group, err := db.FetchByExternalId(AwsCachedLbbgManager, groupId)
group, err := db.FetchByExternalIdAndManagerId(AwsCachedLbbgManager, groupId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", lb.ManagerId)
})
if err != nil {
log.Errorf("Fetch aws loadbalancer backendgroup by external id %s failed: %s", groupId, err)
} else {
@@ -1060,7 +1068,9 @@ func (lblis *SLoadbalancerListener) constructFieldsFromCloudListener(userCred mc
lb := lblis.GetLoadbalancer()
if forward, _ := lb.LBInfo.Int("Forward"); forward == 1 {
// 应用型负载均衡
group, err := db.FetchByExternalId(QcloudCachedLbbgManager, groupId)
group, err := db.FetchByExternalIdAndManagerId(QcloudCachedLbbgManager, groupId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", lb.ManagerId)
})
if err != nil {
log.Errorf("Fetch qcloud loadbalancer backendgroup by external id %s failed: %s", groupId, err)
} else {
@@ -1068,7 +1078,10 @@ func (lblis *SLoadbalancerListener) constructFieldsFromCloudListener(userCred mc
}
} else {
// 传统型负载均衡
if group, err := db.FetchByExternalId(LoadbalancerBackendGroupManager, groupId); err == nil {
if group, err := db.FetchByExternalIdAndManagerId(LoadbalancerBackendGroupManager, groupId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
sq := LoadbalancerManager.Query().SubQuery()
return q.Join(sq, sqlchemy.Equals(sq.Field("id"), q.Field("loadbalancer_id"))).Filter(sqlchemy.Equals(sq.Field("manager_id"), lb.ManagerId))
}); err == nil {
lblis.BackendGroupId = group.GetId()
}
}
@@ -1076,13 +1089,16 @@ func (lblis *SLoadbalancerListener) constructFieldsFromCloudListener(userCred mc
default:
if len(lblis.BackendGroupId) == 0 && len(groupId) == 0 {
lblis.BackendGroupId = lb.BackendGroupId
} else if group, err := db.FetchByExternalId(LoadbalancerBackendGroupManager, groupId); err == nil {
} else if group, err := db.FetchByExternalIdAndManagerId(LoadbalancerBackendGroupManager, groupId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
sq := LoadbalancerManager.Query().SubQuery()
return q.Join(sq, sqlchemy.Equals(sq.Field("id"), q.Field("loadbalancer_id"))).Filter(sqlchemy.Equals(sq.Field("manager_id"), lb.ManagerId))
}); err == nil {
lblis.BackendGroupId = group.GetId()
}
}
}
func (lblis *SLoadbalancerListener) updateCachedLoadbalancerBackendGroupAssociate(ctx context.Context, extListener cloudprovider.ICloudLoadbalancerListener) error {
func (lblis *SLoadbalancerListener) updateCachedLoadbalancerBackendGroupAssociate(ctx context.Context, extListener cloudprovider.ICloudLoadbalancerListener, managerId string) error {
exteralLbbgId := extListener.GetBackendGroupId()
if len(exteralLbbgId) == 0 {
return nil
@@ -1090,7 +1106,9 @@ func (lblis *SLoadbalancerListener) updateCachedLoadbalancerBackendGroupAssociat
switch lblis.GetProviderName() {
case api.CLOUD_PROVIDER_HUAWEI:
_group, err := db.FetchByExternalId(HuaweiCachedLbbgManager, exteralLbbgId)
_group, err := db.FetchByExternalIdAndManagerId(HuaweiCachedLbbgManager, exteralLbbgId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", managerId)
})
if err != nil {
if err == sql.ErrNoRows {
lblis.BackendGroupId = ""
@@ -1115,7 +1133,9 @@ func (lblis *SLoadbalancerListener) updateCachedLoadbalancerBackendGroupAssociat
case api.CLOUD_PROVIDER_QCLOUD:
lb := lblis.GetLoadbalancer()
if forward, _ := lb.LBInfo.Int("Forward"); forward == 1 {
_group, err := db.FetchByExternalId(QcloudCachedLbbgManager, exteralLbbgId)
_group, err := db.FetchByExternalIdAndManagerId(QcloudCachedLbbgManager, exteralLbbgId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", managerId)
})
if err != nil {
if err == sql.ErrNoRows {
lblis.BackendGroupId = ""
@@ -1167,7 +1187,7 @@ func (lblis *SLoadbalancerListener) SyncWithCloudLoadbalancerListener(ctx contex
return err
}
err = lblis.updateCachedLoadbalancerBackendGroupAssociate(ctx, extListener)
err = lblis.updateCachedLoadbalancerBackendGroupAssociate(ctx, extListener, lb.ManagerId)
if err != nil {
return errors.Wrap(err, "LoadbalancerListener.SyncWithCloudLoadbalancerListener")
}
@@ -1199,7 +1219,7 @@ func (man *SLoadbalancerListenerManager) newFromCloudLoadbalancerListener(ctx co
return nil, err
}
err = lblis.updateCachedLoadbalancerBackendGroupAssociate(ctx, extListener)
err = lblis.updateCachedLoadbalancerBackendGroupAssociate(ctx, extListener, lb.ManagerId)
if err != nil {
return nil, errors.Wrap(err, "LoadbalancerListener.newFromCloudLoadbalancerListener")
}
@@ -21,6 +21,7 @@ import (
"yunion.io/x/jsonutils"
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/util/compare"
"yunion.io/x/sqlchemy"
api "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
@@ -94,7 +95,10 @@ func (lbb *SQcloudCachedLb) syncRemoveCloudLoadbalancerBackend(ctx context.Conte
func (lbb *SQcloudCachedLb) constructFieldsFromCloudLoadbalancerBackend(extLoadbalancerBackend cloudprovider.ICloudLoadbalancerBackend) error {
lbb.Status = extLoadbalancerBackend.GetStatus()
instance, err := db.FetchByExternalId(GuestManager, extLoadbalancerBackend.GetBackendId())
instance, err := db.FetchByExternalIdAndManagerId(GuestManager, extLoadbalancerBackend.GetBackendId(), func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
sq := HostManager.Query().SubQuery()
return q.Join(sq, sqlchemy.Equals(sq.Field("id"), q.Field("host_id"))).Filter(sqlchemy.Equals(sq.Field("manager_id"), lbb.ManagerId))
})
if err != nil {
return err
}
+18 -7
View File
@@ -750,11 +750,18 @@ func (man *SLoadbalancerManager) SyncLoadbalancers(ctx context.Context, userCred
return localLbs, remoteLbs, syncResult
}
func getExtLbNetworkIds(extLb cloudprovider.ICloudLoadbalancer) []string {
func getExtLbNetworkIds(extLb cloudprovider.ICloudLoadbalancer, managerId string) []string {
extNetworkIds := extLb.GetNetworkIds()
lbNetworkIds := []string{}
for _, networkId := range extNetworkIds {
if network, err := db.FetchByExternalId(NetworkManager, networkId); err == nil && network != nil {
network, err := db.FetchByExternalIdAndManagerId(NetworkManager, networkId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
wire := WireManager.Query().SubQuery()
vpc := VpcManager.Query().SubQuery()
return q.Join(wire, sqlchemy.Equals(wire.Field("id"), q.Field("wire_id"))).
Join(vpc, sqlchemy.Equals(vpc.Field("id"), wire.Field("vpc_id"))).
Filter(sqlchemy.Equals(vpc.Field("manager_id"), managerId))
})
if err == nil && network != nil {
lbNetworkIds = append(lbNetworkIds, network.GetId())
}
}
@@ -782,11 +789,13 @@ func (man *SLoadbalancerManager) newFromCloudLoadbalancer(ctx context.Context, u
lb.ChargeType = extLb.GetChargeType()
lb.EgressMbps = extLb.GetEgressMbps()
lb.ExternalId = extLb.GetGlobalId()
lbNetworkIds := getExtLbNetworkIds(extLb)
lbNetworkIds := getExtLbNetworkIds(extLb, lb.ManagerId)
lb.NetworkId = strings.Join(lbNetworkIds, ",")
if vpcId := extLb.GetVpcId(); len(vpcId) > 0 {
if vpc, err := db.FetchByExternalId(VpcManager, vpcId); err == nil && vpc != nil {
if vpc, err := db.FetchByExternalIdAndManagerId(VpcManager, vpcId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", provider.Id)
}); err == nil && vpc != nil {
lb.VpcId = vpc.GetId()
}
}
@@ -958,14 +967,16 @@ func (lb *SLoadbalancer) SyncWithCloudLoadbalancer(ctx context.Context, userCred
lb.EgressMbps = extLb.GetEgressMbps()
lb.ChargeType = extLb.GetChargeType()
lb.ManagerId = provider.Id
lbNetworkIds := getExtLbNetworkIds(extLb)
lbNetworkIds := getExtLbNetworkIds(extLb, lb.ManagerId)
lb.NetworkId = strings.Join(lbNetworkIds, ",")
if extLb.GetMetadata() != nil {
lb.LBInfo = extLb.GetMetadata()
}
if vpcId := extLb.GetVpcId(); len(vpcId) > 0 {
if vpc, err := db.FetchByExternalId(VpcManager, vpcId); err == nil && vpc != nil {
if vpc, err := db.FetchByExternalIdAndManagerId(VpcManager, vpcId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", provider.Id)
}); err == nil && vpc != nil {
lb.VpcId = vpc.GetId()
}
}
@@ -975,7 +986,7 @@ func (lb *SLoadbalancer) SyncWithCloudLoadbalancer(ctx context.Context, userCred
db.OpsLog.LogSyncUpdate(lb, diff, userCred)
networkIds := getExtLbNetworkIds(extLb)
networkIds := getExtLbNetworkIds(extLb, lb.ManagerId)
SyncCloudProject(userCred, lb, syncOwnerId, extLb, provider.Id)
lb.syncLoadbalancerNetwork(ctx, userCred, networkIds)
+18 -6
View File
@@ -261,7 +261,7 @@ func (manager *SNatSEntryManager) SyncNatSTable(ctx context.Context, userCred mc
}
for i := 0; i < len(commondb); i += 1 {
err := commondb[i].SyncWithCloudNatSTable(ctx, userCred, commonext[i], syncOwnerId)
err := commondb[i].SyncWithCloudNatSTable(ctx, userCred, commonext[i], syncOwnerId, provider.Id)
if err != nil {
result.UpdateError(err)
continue
@@ -271,7 +271,7 @@ func (manager *SNatSEntryManager) SyncNatSTable(ctx context.Context, userCred mc
}
for i := 0; i < len(added); i += 1 {
routeTableNew, err := manager.newFromCloudNatSTable(ctx, userCred, syncOwnerId, nat, added[i])
routeTableNew, err := manager.newFromCloudNatSTable(ctx, userCred, syncOwnerId, nat, added[i], provider.Id)
if err != nil {
result.AddError(err)
continue
@@ -293,13 +293,19 @@ func (self *SNatSEntry) syncRemoveCloudNatSTable(ctx context.Context, userCred m
return self.RealDelete(ctx, userCred)
}
func (self *SNatSEntry) SyncWithCloudNatSTable(ctx context.Context, userCred mcclient.TokenCredential, extEntry cloudprovider.ICloudNatSEntry, syncOwnerId mcclient.IIdentityProvider) error {
func (self *SNatSEntry) SyncWithCloudNatSTable(ctx context.Context, userCred mcclient.TokenCredential, extEntry cloudprovider.ICloudNatSEntry, syncOwnerId mcclient.IIdentityProvider, managerId string) error {
diff, err := db.UpdateWithLock(ctx, self, func() error {
self.Status = extEntry.GetStatus()
self.IP = extEntry.GetIP()
self.SourceCIDR = extEntry.GetSourceCIDR()
if extNetworkId := extEntry.GetNetworkId(); len(extNetworkId) > 0 {
network, err := db.FetchByExternalId(NetworkManager, extNetworkId)
network, err := db.FetchByExternalIdAndManagerId(NetworkManager, extNetworkId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
wire := WireManager.Query().SubQuery()
vpc := VpcManager.Query().SubQuery()
return q.Join(wire, sqlchemy.Equals(wire.Field("id"), q.Field("wire_id"))).
Join(vpc, sqlchemy.Equals(vpc.Field("id"), wire.Field("vpc_id"))).
Filter(sqlchemy.Equals(vpc.Field("manager_id"), managerId))
})
if err != nil {
return err
}
@@ -317,7 +323,7 @@ func (self *SNatSEntry) SyncWithCloudNatSTable(ctx context.Context, userCred mcc
return nil
}
func (manager *SNatSEntryManager) newFromCloudNatSTable(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, nat *SNatGateway, extEntry cloudprovider.ICloudNatSEntry) (*SNatSEntry, error) {
func (manager *SNatSEntryManager) newFromCloudNatSTable(ctx context.Context, userCred mcclient.TokenCredential, ownerId mcclient.IIdentityProvider, nat *SNatGateway, extEntry cloudprovider.ICloudNatSEntry, managerId string) (*SNatSEntry, error) {
table := SNatSEntry{}
table.SetModelManager(manager, &table)
@@ -330,7 +336,13 @@ func (manager *SNatSEntryManager) newFromCloudNatSTable(ctx context.Context, use
table.IP = extEntry.GetIP()
table.SourceCIDR = extEntry.GetSourceCIDR()
if extNetworkId := extEntry.GetNetworkId(); len(extNetworkId) > 0 {
network, err := db.FetchByExternalId(NetworkManager, extNetworkId)
network, err := db.FetchByExternalIdAndManagerId(NetworkManager, extNetworkId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
wire := WireManager.Query().SubQuery()
vpc := VpcManager.Query().SubQuery()
return q.Join(wire, sqlchemy.Equals(wire.Field("id"), q.Field("wire_id"))).
Join(vpc, sqlchemy.Equals(vpc.Field("id"), wire.Field("vpc_id"))).
Filter(sqlchemy.Equals(vpc.Field("manager_id"), managerId))
})
if err != nil {
return nil, err
}
@@ -22,6 +22,7 @@ import (
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/util/compare"
"yunion.io/x/pkg/util/netutils"
"yunion.io/x/sqlchemy"
"yunion.io/x/onecloud/pkg/cloudcommon/db"
"yunion.io/x/onecloud/pkg/cloudcommon/db/lockman"
@@ -182,7 +183,13 @@ func (manager *SNetworkinterfacenetworkManager) newFromCloudInterfaceAddress(ctx
address.SetModelManager(manager, &address)
networkId := ext.GetINetworkId()
_network, err := db.FetchByExternalId(NetworkManager, networkId)
_network, err := db.FetchByExternalIdAndManagerId(NetworkManager, networkId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
wire := WireManager.Query().SubQuery()
vpc := VpcManager.Query().SubQuery()
return q.Join(wire, sqlchemy.Equals(wire.Field("id"), q.Field("wire_id"))).
Join(vpc, sqlchemy.Equals(vpc.Field("id"), wire.Field("vpc_id"))).
Filter(sqlchemy.Equals(vpc.Field("manager_id"), networkinterface.ManagerId))
})
if err != nil {
return errors.Wrapf(err, "newFromCloudInterfaceAddress.FetchByExternalId(%s)", networkId)
}
+4 -1
View File
@@ -313,7 +313,10 @@ func (self *SNetworkInterface) SyncWithCloudNetworkInterface(ctx context.Context
func (self *SNetworkInterface) Associate(associateId string) error {
switch self.AssociateType {
case api.NETWORK_INTERFACE_ASSOCIATE_TYPE_SERVER:
guest, err := db.FetchByExternalId(GuestManager, associateId)
guest, err := db.FetchByExternalIdAndManagerId(GuestManager, associateId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
sq := HostManager.Query().SubQuery()
return q.Join(sq, sqlchemy.Equals(sq.Field("id"), q.Field("host_id"))).Filter(sqlchemy.Equals(q.Field("manager_id"), self.ManagerId))
})
if err != nil {
return errors.Wrapf(err, "failed to get guest for networkinterface %s associateId %s", self.Name, associateId)
}
+11 -1
View File
@@ -309,7 +309,17 @@ func (self *SNetwork) GetNetworkInterfacesCount() (int, error) {
}
func (manager *SNetworkManager) GetOrCreateClassicNetwork(wire *SWire) (*SNetwork, error) {
_network, err := db.FetchByExternalId(manager, wire.Id)
_network, err := db.FetchByExternalIdAndManagerId(manager, wire.Id, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
v := wire.GetVpc()
if v != nil {
wire := WireManager.Query().SubQuery()
vpc := VpcManager.Query().SubQuery()
return q.Join(wire, sqlchemy.Equals(wire.Field("id"), q.Field("wire_id"))).
Join(vpc, sqlchemy.Equals(vpc.Field("id"), wire.Field("vpc_id"))).
Filter(sqlchemy.Equals(vpc.Field("manager_id"), v.ManagerId))
}
return q
})
if err == nil {
return _network.(*SNetwork), nil
}
+4 -1
View File
@@ -877,7 +877,10 @@ func (manager *SSnapshotManager) newFromCloudSnapshot(ctx context.Context, userC
snapshot.ExternalId = extSnapshot.GetGlobalId()
var localDisk *SDisk
if len(extSnapshot.GetDiskId()) > 0 {
disk, err := db.FetchByExternalId(DiskManager, extSnapshot.GetDiskId())
disk, err := db.FetchByExternalIdAndManagerId(DiskManager, extSnapshot.GetDiskId(), func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
sq := StorageManager.Query().SubQuery()
return q.Join(sq, sqlchemy.Equals(q.Field("storage_id"), sq.Field("id"))).Filter(sqlchemy.Equals(sq.Field("manager_id"), provider.Id))
})
if err != nil {
log.Errorf("snapshot %s missing disk?", snapshot.Name)
} else {
+3 -1
View File
@@ -173,7 +173,9 @@ func (manager *SStoragecacheManager) SyncWithCloudStoragecache(ctx context.Conte
lockman.LockClass(ctx, manager, db.GetLockClassKey(manager, userCred))
defer lockman.ReleaseClass(ctx, manager, db.GetLockClassKey(manager, userCred))
localCacheObj, err := db.FetchByExternalId(manager, cloudCache.GetGlobalId())
localCacheObj, err := db.FetchByExternalIdAndManagerId(manager, cloudCache.GetGlobalId(), func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", provider.Id)
})
if err != nil {
if err == sql.ErrNoRows {
localCache, err := manager.newFromCloudStoragecache(ctx, userCred, cloudCache, provider)
+4 -15
View File
@@ -183,7 +183,9 @@ func (manager *SVpcManager) GetOrCreateVpcForClassicNetwork(host *SHost) (*SVpc,
region := host.GetRegion()
cloudprovider := host.GetCloudprovider()
externalId := manager.getVpcExternalIdForClassicNetwork(region.Id, cloudprovider.Id)
_vpc, err := db.FetchByExternalId(manager, externalId)
_vpc, err := db.FetchByExternalIdAndManagerId(manager, externalId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", host.ManagerId)
})
if err == nil {
return _vpc.(*SVpc), nil
}
@@ -310,19 +312,6 @@ func (manager *SVpcManager) FetchCustomizeColumns(
return rows
}
func (manager *SVpcManager) getVpcsByRegion(region *SCloudregion, provider *SCloudprovider) ([]SVpc, error) {
vpcs := make([]SVpc, 0)
q := manager.Query().Equals("cloudregion_id", region.Id)
if provider != nil {
q = q.Equals("manager_id", provider.Id)
}
err := db.FetchModelObjects(manager, q, &vpcs)
if err != nil {
return nil, err
}
return vpcs, nil
}
func (self *SVpc) setDefault(def bool) error {
var err error
if self.IsDefault != def {
@@ -342,7 +331,7 @@ func (manager *SVpcManager) SyncVPCs(ctx context.Context, userCred mcclient.Toke
remoteVPCs := make([]cloudprovider.ICloudVpc, 0)
syncResult := compare.SyncResult{}
dbVPCs, err := manager.getVpcsByRegion(region, provider)
dbVPCs, err := region.GetVpcs()
if err != nil {
syncResult.Error(err)
return nil, nil, syncResult
+4 -1
View File
@@ -196,7 +196,10 @@ func (manager *SWireManager) GetOrCreateWireForClassicNetwork(vpc *SVpc, zone *S
} else {
name = fmt.Sprintf("emulate for zone %s vpc %s classic network", zone.Name, vpc.Id)
}
_wire, err := db.FetchByExternalId(manager, externalId)
_wire, err := db.FetchByExternalIdAndManagerId(manager, externalId, func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
sq := VpcManager.Query().SubQuery()
return q.Join(sq, sqlchemy.Equals(sq.Field("id"), q.Field("id"))).Filter(sqlchemy.Equals(sq.Field("manager_id"), vpc.ManagerId))
})
if err == nil {
return _wire.(*SWire), nil
}
+1 -11
View File
@@ -183,16 +183,6 @@ func (zone *SZone) GetCloudRegionId() string {
}
}
func (manager *SZoneManager) GetZonesByRegion(region *SCloudregion) ([]SZone, error) {
zones := make([]SZone, 0)
q := manager.Query().Equals("cloudregion_id", region.Id)
err := db.FetchModelObjects(manager, q, &zones)
if err != nil {
return nil, err
}
return zones, nil
}
func (manager *SZoneManager) SyncZones(ctx context.Context, userCred mcclient.TokenCredential, region *SCloudregion, zones []cloudprovider.ICloudZone) ([]SZone, []cloudprovider.ICloudZone, compare.SyncResult) {
lockman.LockClass(ctx, manager, db.GetLockClassKey(manager, userCred))
defer lockman.ReleaseClass(ctx, manager, db.GetLockClassKey(manager, userCred))
@@ -201,7 +191,7 @@ func (manager *SZoneManager) SyncZones(ctx context.Context, userCred mcclient.To
remoteZones := make([]cloudprovider.ICloudZone, 0)
syncResult := compare.SyncResult{}
dbZones, err := manager.GetZonesByRegion(region)
dbZones, err := region.GetZones()
if err != nil {
syncResult.Error(err)
return nil, nil, syncResult
+4 -1
View File
@@ -27,6 +27,7 @@ import (
"yunion.io/x/pkg/errors"
"yunion.io/x/pkg/util/secrules"
"yunion.io/x/pkg/utils"
"yunion.io/x/sqlchemy"
billing_api "yunion.io/x/onecloud/pkg/apis/billing"
api "yunion.io/x/onecloud/pkg/apis/compute"
@@ -1908,7 +1909,9 @@ func (self *SHuaWeiRegionDriver) RequestCreateLoadbalancer(ctx context.Context,
return nil, errors.Wrap(err, "Huawei.RequestCreateLoadbalancer.Associate")
}
eip, err := db.FetchByExternalId(models.ElasticipManager, ieip.GetGlobalId())
eip, err := db.FetchByExternalIdAndManagerId(models.ElasticipManager, ieip.GetGlobalId(), func(q *sqlchemy.SQuery) *sqlchemy.SQuery {
return q.Equals("manager_id", lb.MarkUnDelete)
})
if err != nil {
return nil, errors.Wrap(err, "Huawei.RequestCreateLoadbalancer.FetchByExternalId")
}