diff --git a/pkg/cloudcommon/db/external.go b/pkg/cloudcommon/db/external.go index 3c8f8b702b..71f5e2164b 100644 --- a/pkg/cloudcommon/db/external.go +++ b/pkg/cloudcommon/db/external.go @@ -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 diff --git a/pkg/compute/guestdrivers/managedvirtual.go b/pkg/compute/guestdrivers/managedvirtual.go index e3076ce9b1..10f83edbe2 100644 --- a/pkg/compute/guestdrivers/managedvirtual.go +++ b/pkg/compute/guestdrivers/managedvirtual.go @@ -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() { diff --git a/pkg/compute/models/cloudproviderregions.go b/pkg/compute/models/cloudproviderregions.go index c165372084..b27880dc92 100644 --- a/pkg/compute/models/cloudproviderregions.go +++ b/pkg/compute/models/cloudproviderregions.go @@ -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) diff --git a/pkg/compute/models/cloudregions.go b/pkg/compute/models/cloudregions.go index 28f8c59d00..4537d62909 100644 --- a/pkg/compute/models/cloudregions.go +++ b/pkg/compute/models/cloudregions.go @@ -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 } diff --git a/pkg/compute/models/dbinstance_backups.go b/pkg/compute/models/dbinstance_backups.go index bb17a5f699..ebb79a1797 100644 --- a/pkg/compute/models/dbinstance_backups.go +++ b/pkg/compute/models/dbinstance_backups.go @@ -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 { diff --git a/pkg/compute/models/dbinstancenetworks.go b/pkg/compute/models/dbinstancenetworks.go index 649c4b8436..62061868f1 100644 --- a/pkg/compute/models/dbinstancenetworks.go +++ b/pkg/compute/models/dbinstancenetworks.go @@ -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) diff --git a/pkg/compute/models/dbinstances.go b/pkg/compute/models/dbinstances.go index 7a28f6ac7d..7f48c86cfc 100644 --- a/pkg/compute/models/dbinstances.go +++ b/pkg/compute/models/dbinstances.go @@ -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") } diff --git a/pkg/compute/models/disks.go b/pkg/compute/models/disks.go index 90eaf72ed1..da89b76c63 100644 --- a/pkg/compute/models/disks.go +++ b/pkg/compute/models/disks.go @@ -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() { diff --git a/pkg/compute/models/elasticcache_instances.go b/pkg/compute/models/elasticcache_instances.go index 57410e0a84..19c3d40edf 100644 --- a/pkg/compute/models/elasticcache_instances.go +++ b/pkg/compute/models/elasticcache_instances.go @@ -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") } diff --git a/pkg/compute/models/elasticips.go b/pkg/compute/models/elasticips.go index d72abe9e65..e3c8cf8032 100644 --- a/pkg/compute/models/elasticips.go +++ b/pkg/compute/models/elasticips.go @@ -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 } diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index 979c7f7710..ad7770e8ae 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -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) diff --git a/pkg/compute/models/host_recycle.go b/pkg/compute/models/host_recycle.go index 92c8a1355e..7d86dd2d99 100644 --- a/pkg/compute/models/host_recycle.go +++ b/pkg/compute/models/host_recycle.go @@ -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 diff --git a/pkg/compute/models/hosts.go b/pkg/compute/models/hosts.go index b12ebe39ea..6b8b1016ad 100644 --- a/pkg/compute/models/hosts.go +++ b/pkg/compute/models/hosts.go @@ -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 diff --git a/pkg/compute/models/loadbalancerawscachedlbb.go b/pkg/compute/models/loadbalancerawscachedlbb.go index 702e29f96e..642ad4f190 100644 --- a/pkg/compute/models/loadbalancerawscachedlbb.go +++ b/pkg/compute/models/loadbalancerawscachedlbb.go @@ -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 } diff --git a/pkg/compute/models/loadbalancerawscachedlbbg.go b/pkg/compute/models/loadbalancerawscachedlbbg.go index 653379bd28..a40a2b8a66 100644 --- a/pkg/compute/models/loadbalancerawscachedlbbg.go +++ b/pkg/compute/models/loadbalancerawscachedlbbg.go @@ -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) } diff --git a/pkg/compute/models/loadbalancerbackends.go b/pkg/compute/models/loadbalancerbackends.go index 29c5387b32..ebb86ebe9e 100644 --- a/pkg/compute/models/loadbalancerbackends.go +++ b/pkg/compute/models/loadbalancerbackends.go @@ -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 } diff --git a/pkg/compute/models/loadbalancercachedacls.go b/pkg/compute/models/loadbalancercachedacls.go index cc757a31f7..b31d61f15f 100644 --- a/pkg/compute/models/loadbalancercachedacls.go +++ b/pkg/compute/models/loadbalancercachedacls.go @@ -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") } diff --git a/pkg/compute/models/loadbalancerhuaweicachedlbb.go b/pkg/compute/models/loadbalancerhuaweicachedlbb.go index 711d78a85e..ba2ae3d219 100644 --- a/pkg/compute/models/loadbalancerhuaweicachedlbb.go +++ b/pkg/compute/models/loadbalancerhuaweicachedlbb.go @@ -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 } diff --git a/pkg/compute/models/loadbalancerlisteners.go b/pkg/compute/models/loadbalancerlisteners.go index 2f3c42466e..a9bd36be34 100644 --- a/pkg/compute/models/loadbalancerlisteners.go +++ b/pkg/compute/models/loadbalancerlisteners.go @@ -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") } diff --git a/pkg/compute/models/loadbalancerqcloudcachedlbb.go b/pkg/compute/models/loadbalancerqcloudcachedlbb.go index 2efbea1f7c..d4e2d7674e 100644 --- a/pkg/compute/models/loadbalancerqcloudcachedlbb.go +++ b/pkg/compute/models/loadbalancerqcloudcachedlbb.go @@ -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 } diff --git a/pkg/compute/models/loadbalancers.go b/pkg/compute/models/loadbalancers.go index e9dfc5e0b7..62d9fe1bf8 100644 --- a/pkg/compute/models/loadbalancers.go +++ b/pkg/compute/models/loadbalancers.go @@ -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) diff --git a/pkg/compute/models/natstable.go b/pkg/compute/models/natstable.go index 0ad34d9825..b9327196c1 100644 --- a/pkg/compute/models/natstable.go +++ b/pkg/compute/models/natstable.go @@ -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 } diff --git a/pkg/compute/models/networkinterfacenetwork.go b/pkg/compute/models/networkinterfacenetwork.go index 65b96a0e4b..760d51a10b 100644 --- a/pkg/compute/models/networkinterfacenetwork.go +++ b/pkg/compute/models/networkinterfacenetwork.go @@ -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) } diff --git a/pkg/compute/models/networkinterfaces.go b/pkg/compute/models/networkinterfaces.go index 7926a9e1c3..c03453ed67 100644 --- a/pkg/compute/models/networkinterfaces.go +++ b/pkg/compute/models/networkinterfaces.go @@ -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) } diff --git a/pkg/compute/models/networks.go b/pkg/compute/models/networks.go index 8c0c89c482..c12f821292 100644 --- a/pkg/compute/models/networks.go +++ b/pkg/compute/models/networks.go @@ -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 } diff --git a/pkg/compute/models/snapshots.go b/pkg/compute/models/snapshots.go index d2ed35a76f..89a7cbdbce 100644 --- a/pkg/compute/models/snapshots.go +++ b/pkg/compute/models/snapshots.go @@ -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 { diff --git a/pkg/compute/models/storagecaches.go b/pkg/compute/models/storagecaches.go index 9e12ba8aa0..7fa02981fc 100644 --- a/pkg/compute/models/storagecaches.go +++ b/pkg/compute/models/storagecaches.go @@ -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) diff --git a/pkg/compute/models/vpcs.go b/pkg/compute/models/vpcs.go index 69de08b3b6..fbc3774a90 100644 --- a/pkg/compute/models/vpcs.go +++ b/pkg/compute/models/vpcs.go @@ -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 diff --git a/pkg/compute/models/wires.go b/pkg/compute/models/wires.go index eeb4860c38..0a8252feeb 100644 --- a/pkg/compute/models/wires.go +++ b/pkg/compute/models/wires.go @@ -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 } diff --git a/pkg/compute/models/zones.go b/pkg/compute/models/zones.go index 1ffcb3040d..c9a4b105a3 100644 --- a/pkg/compute/models/zones.go +++ b/pkg/compute/models/zones.go @@ -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 diff --git a/pkg/compute/regiondrivers/huawei.go b/pkg/compute/regiondrivers/huawei.go index 9595df2e47..9f15a74597 100644 --- a/pkg/compute/regiondrivers/huawei.go +++ b/pkg/compute/regiondrivers/huawei.go @@ -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") }