diff --git a/pkg/cloudprovider/loadbalancerlistener.go b/pkg/cloudprovider/loadbalancerlistener.go index a697a4d5da..5e11342796 100644 --- a/pkg/cloudprovider/loadbalancerlistener.go +++ b/pkg/cloudprovider/loadbalancerlistener.go @@ -31,6 +31,11 @@ type SLoadbalancerListener struct { Description string EstablishedTimeout int + ClientRequestTimeout int + ClientIdleTimeout int + BackendConnectTimeout int + BackendIdleTimeout int + HealthCheckReq string HealthCheckExp string diff --git a/pkg/compute/guestdrivers/aws.go b/pkg/compute/guestdrivers/aws.go index b55d465c61..83b30e4d31 100644 --- a/pkg/compute/guestdrivers/aws.go +++ b/pkg/compute/guestdrivers/aws.go @@ -19,9 +19,13 @@ import ( "fmt" "strings" + "yunion.io/x/jsonutils" + "yunion.io/x/log" + "yunion.io/x/pkg/errors" "yunion.io/x/pkg/utils" api "yunion.io/x/onecloud/pkg/apis/compute" + "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/mcclient" @@ -165,3 +169,45 @@ func (self *SAwsGuestDriver) GetGuestInitialStateAfterRebuild() string { func (self *SAwsGuestDriver) IsSupportedBillingCycle(bc billing.SBillingCycle) bool { return false } + +func (self *SAwsGuestDriver) RequestSyncstatusOnHost(ctx context.Context, guest *models.SGuest, host *models.SHost, userCred mcclient.TokenCredential) (jsonutils.JSONObject, error) { + ihost, err := host.GetIHost() + if err != nil { + return nil, err + } + ivm, err := ihost.GetIVMById(guest.ExternalId) + if err != nil { + log.Errorf("fail to find ivm by id %s", err) + return nil, err + } + + err = guest.SyncAllWithCloudVM(ctx, userCred, host, ivm) + if err != nil { + return nil, err + } + + ieip, err := ivm.GetIEIP() + if err != nil { + return nil, errors.Wrap(err, "AwsGuestDriver.GetIEIP") + } + + // 如果aws已经绑定了EIP,则要把多余的公有IP删除 + if ieip.GetMode() == api.EIP_MODE_STANDALONE_EIP { + publicIP, err := guest.GetPublicIp() + if err != nil { + return nil, errors.Wrap(err, "AwsGuestDriver.GetPublicIp") + } + + if publicIP != nil { + err = db.DeleteModel(ctx, userCred, publicIP) + if err != nil { + return nil, errors.Wrap(err, "AwsGuestDriver.DeletePublicIp") + } + } + } + + status := GetCloudVMStatus(ivm) + body := jsonutils.NewDict() + body.Add(jsonutils.NewString(status), "status") + return body, nil +} diff --git a/pkg/compute/guestdrivers/managedvirtual.go b/pkg/compute/guestdrivers/managedvirtual.go index a4324ec8f2..5a89dfdb59 100644 --- a/pkg/compute/guestdrivers/managedvirtual.go +++ b/pkg/compute/guestdrivers/managedvirtual.go @@ -657,24 +657,7 @@ func (self *SManagedVirtualizedGuestDriver) RequestSyncstatusOnHost(ctx context. return nil, err } - status := ivm.GetStatus() - switch status { - case api.VM_RUNNING: - status = cloudprovider.CloudVMStatusRunning - case api.VM_READY: - status = cloudprovider.CloudVMStatusStopped - case api.VM_STARTING: - status = cloudprovider.CloudVMStatusStopped - case api.VM_STOPPING: - status = cloudprovider.CloudVMStatusStopping - case api.VM_CHANGE_FLAVOR: - status = cloudprovider.CloudVMStatusChangeFlavor - case api.VM_DEPLOYING: - status = cloudprovider.CloudVMStatusDeploying - default: - status = cloudprovider.CloudVMStatusOther - } - + status := GetCloudVMStatus(ivm) body := jsonutils.NewDict() body.Add(jsonutils.NewString(status), "status") return body, nil @@ -986,3 +969,25 @@ func (self *SManagedVirtualizedGuestDriver) chooseHostStorage( func (self *SManagedVirtualizedGuestDriver) IsSupportCdrom(guest *models.SGuest) (bool, error) { return false, nil } + +func GetCloudVMStatus(vm cloudprovider.ICloudVM) string { + status := vm.GetStatus() + switch status { + case api.VM_RUNNING: + status = cloudprovider.CloudVMStatusRunning + case api.VM_READY: + status = cloudprovider.CloudVMStatusStopped + case api.VM_STARTING: + status = cloudprovider.CloudVMStatusStopped + case api.VM_STOPPING: + status = cloudprovider.CloudVMStatusStopping + case api.VM_CHANGE_FLAVOR: + status = cloudprovider.CloudVMStatusChangeFlavor + case api.VM_DEPLOYING: + status = cloudprovider.CloudVMStatusDeploying + default: + status = cloudprovider.CloudVMStatusOther + } + + return status +} diff --git a/pkg/compute/models/cloudsyncelb.go b/pkg/compute/models/cloudsyncelb.go index 32056e5c58..f4694a81aa 100644 --- a/pkg/compute/models/cloudsyncelb.go +++ b/pkg/compute/models/cloudsyncelb.go @@ -81,7 +81,9 @@ func syncRegionLoadbalancers(ctx context.Context, userCred mcclient.TokenCredent } db.OpsLog.LogEvent(provider, db.ACT_SYNC_LB_COMPLETE, msg, userCred) // 同步未关联负载均衡的后端服务器组 - syncAwsLoadbalancerBackendgroups(ctx, userCred, syncResults, provider, localRegion, remoteRegion, syncRange) + if provider.Provider == compute.CLOUD_PROVIDER_AWS { + syncAwsLoadbalancerBackendgroups(ctx, userCred, syncResults, provider, localRegion, remoteRegion, syncRange) + } for i := 0; i < len(localLbs); i++ { func() { diff --git a/pkg/compute/models/elasticips.go b/pkg/compute/models/elasticips.go index f87f205bcf..4466727047 100644 --- a/pkg/compute/models/elasticips.go +++ b/pkg/compute/models/elasticips.go @@ -456,6 +456,30 @@ func (manager *SElasticipManager) getEipForInstance(instanceType string, instanc return &eip, nil } +func (manager *SElasticipManager) getPublicIpForInstance(instanceType string, instanceId string) (*SElasticip, error) { + eip := SElasticip{} + + q := manager.Query() + q = q.Equals("associate_type", instanceType) + q = q.Equals("associate_id", instanceId) + q = q.Equals("mode", api.EIP_MODE_INSTANCE_PUBLICIP) + + err := q.First(&eip) + + if err != nil { + if err != sql.ErrNoRows { + log.Errorf("getEipForInstance query fail %s", err) + return nil, err + } else { + return nil, nil + } + } + + eip.SetModelManager(manager, &eip) + + return &eip, nil +} + func (self *SElasticip) IsAssociated() bool { if len(self.AssociateId) == 0 { return false diff --git a/pkg/compute/models/guests.go b/pkg/compute/models/guests.go index 36fd15db51..092d8fa55b 100644 --- a/pkg/compute/models/guests.go +++ b/pkg/compute/models/guests.go @@ -4085,6 +4085,10 @@ func (self *SGuest) GetEip() (*SElasticip, error) { return ElasticipManager.getEipForInstance("server", self.Id) } +func (self *SGuest) GetPublicIp() (*SElasticip, error) { + return ElasticipManager.getPublicIpForInstance("server", self.Id) +} + func (self *SGuest) SyncVMEip(ctx context.Context, userCred mcclient.TokenCredential, provider *SCloudprovider, extEip cloudprovider.ICloudEIP, syncOwnerId mcclient.IIdentityProvider) compare.SyncResult { result := compare.SyncResult{} diff --git a/pkg/compute/models/loadbalancerawscachedlbb.go b/pkg/compute/models/loadbalancerawscachedlbb.go index 094eeba1c2..d7eda141ee 100644 --- a/pkg/compute/models/loadbalancerawscachedlbb.go +++ b/pkg/compute/models/loadbalancerawscachedlbb.go @@ -19,6 +19,7 @@ import ( "fmt" "yunion.io/x/jsonutils" + "yunion.io/x/pkg/errors" "yunion.io/x/pkg/util/compare" api "yunion.io/x/onecloud/pkg/apis/compute" @@ -92,7 +93,7 @@ func (man *SAwsCachedLbManager) CreateAwsCachedLb(ctx context.Context, userCred return nil, err } - err = man.TableSpec().Insert(lbb) + err = man.TableSpec().Insert(cachedlbb) if err != nil { return nil, err @@ -209,7 +210,26 @@ func (lbb *SAwsCachedLb) constructFieldsFromCloudLoadbalancerBackend(extLoadbala func (lbb *SAwsCachedLb) SyncWithCloudLoadbalancerBackend(ctx context.Context, userCred mcclient.TokenCredential, extLoadbalancerBackend cloudprovider.ICloudLoadbalancerBackend, syncOwnerId mcclient.IIdentityProvider) error { lbb.SetModelManager(AwsCachedLbManager, lbb) + cacheLbbg, err := lbb.GetCachedBackendGroup() + if err != nil { + return errors.Wrap(err, "AwsCachedLb.SyncWithCloudLoadbalancerBackend.GetCachedBackendGroup") + } + + localLbbg, err := cacheLbbg.GetLocalBackendGroup(ctx, userCred) + if err != nil { + return errors.Wrap(err, "AwsCachedLb.SyncWithCloudLoadbalancerBackend.GetLocalBackendGroup") + } + + locallbb, err := newLocalBackendFromCloudLoadbalancerBackend(ctx, userCred, localLbbg, extLoadbalancerBackend, syncOwnerId) + if err != nil { + return errors.Wrap(err, "AwsCachedLb.SyncWithCloudLoadbalancerBackend.newLocalBackendFromCloudLoadbalancerBackend") + } + diff, err := db.UpdateWithLock(ctx, lbb, func() error { + if locallbb != nil { + lbb.BackendId = locallbb.GetId() + } + return lbb.constructFieldsFromCloudLoadbalancerBackend(extLoadbalancerBackend) }) if err != nil { diff --git a/pkg/compute/models/loadbalancerawscachedlbbg.go b/pkg/compute/models/loadbalancerawscachedlbbg.go index 309a5c10ca..e866bcaef7 100644 --- a/pkg/compute/models/loadbalancerawscachedlbbg.go +++ b/pkg/compute/models/loadbalancerawscachedlbbg.go @@ -17,9 +17,12 @@ package models import ( "context" "fmt" + "strconv" + "strings" "yunion.io/x/jsonutils" "yunion.io/x/log" + "yunion.io/x/pkg/errors" "yunion.io/x/pkg/util/compare" "yunion.io/x/pkg/utils" "yunion.io/x/sqlchemy" @@ -314,8 +317,61 @@ func (lbbg *SAwsCachedLbbg) syncRemoveCloudLoadbalancerBackendgroup(ctx context. return err } +func (lbbg *SAwsCachedLbbg) isBackendsMatch(backends []SLoadbalancerBackend, ibackends []cloudprovider.ICloudLoadbalancerBackend) bool { + if len(ibackends) != len(backends) { + return false + } + + locals := []string{} + remotes := []string{} + + for i := range backends { + guest := backends[i].GetGuest() + seg := strings.Join([]string{guest.ExternalId, strconv.Itoa(backends[i].Port)}, "/") + locals = append(locals, seg) + } + + for i := range ibackends { + ibackend := ibackends[i] + seg := strings.Join([]string{ibackend.GetBackendId(), strconv.Itoa(ibackend.GetPort())}, "/") + remotes = append(remotes, seg) + } + + for i := range remotes { + if !utils.IsInStringArray(remotes[i], locals) { + return false + } + } + + return true +} + func (lbbg *SAwsCachedLbbg) SyncWithCloudLoadbalancerBackendgroup(ctx context.Context, userCred mcclient.TokenCredential, lb *SLoadbalancer, extLoadbalancerBackendgroup cloudprovider.ICloudLoadbalancerBackendGroup, syncOwnerId mcclient.IIdentityProvider) error { lbbg.SetModelManager(AwsCachedLbbgManager, lbbg) + + ibackends, err := extLoadbalancerBackendgroup.GetILoadbalancerBackends() + if err != nil { + return errors.Wrap(err, "AwsCachedLbbg.SyncWithCloudLoadbalancerBackendgroup.GetILoadbalancerBackends") + } + + localLbbg, err := lbbg.GetLocalBackendGroup(ctx, userCred) + if err != nil { + return errors.Wrap(err, "AwsCachedLbbg.SyncWithCloudLoadbalancerBackendgroup.GetLocalBackendGroup") + } + + backends, err := localLbbg.GetBackends() + if err != nil { + return errors.Wrap(err, "AwsCachedLbbg.SyncWithCloudLoadbalancerBackendgroup.GetBackends") + } + + var newLocalLbbg *SLoadbalancerBackendGroup + if !lbbg.isBackendsMatch(backends, ibackends) { + newLocalLbbg, err = newLocalBackendgroupFromCloudLoadbalancerBackendgroup(ctx, userCred, lb, extLoadbalancerBackendgroup, syncOwnerId) + if err != nil { + return errors.Wrap(err, "HuaweiCachedLbbg.SyncWithCloudLoadbalancerBackendgroup.newLocalBackendgroupFromCloudLoadbalancerBackendgroup") + } + } + diff, err := db.UpdateWithLock(ctx, lbbg, func() error { lbbg.Status = extLoadbalancerBackendgroup.GetStatus() metadata := extLoadbalancerBackendgroup.GetMetadata() @@ -328,6 +384,9 @@ func (lbbg *SAwsCachedLbbg) SyncWithCloudLoadbalancerBackendgroup(ctx context.Co if interval, _ := metadata.Int("health_check_interval"); interval > 0 { lbbg.HealthCheckInterval = int(interval) } + if newLocalLbbg != nil { + lbbg.BackendGroupId = newLocalLbbg.GetId() + } return nil }) if err != nil { diff --git a/pkg/compute/models/loadbalancerbackendgroups.go b/pkg/compute/models/loadbalancerbackendgroups.go index cb9d140318..72e784bf5c 100644 --- a/pkg/compute/models/loadbalancerbackendgroups.go +++ b/pkg/compute/models/loadbalancerbackendgroups.go @@ -283,7 +283,7 @@ func (lbbg *SLoadbalancerBackendGroup) RefCount() (int, error) { func (lbbg *SLoadbalancerBackendGroup) refCount(man db.IModelManager) (int, error) { t := man.TableSpec().Instance() pdF := t.Field("pending_deleted") - return t.Query(). + return t.Query().IsFalse("deleted"). Equals("backend_group_id", lbbg.Id). Filter(sqlchemy.OR(sqlchemy.IsNull(pdF), sqlchemy.IsFalse(pdF))). CountWithError() diff --git a/pkg/compute/models/loadbalancercachedcertificates.go b/pkg/compute/models/loadbalancercachedcertificates.go index 9379ed5b39..5358b1e8c8 100644 --- a/pkg/compute/models/loadbalancercachedcertificates.go +++ b/pkg/compute/models/loadbalancercachedcertificates.go @@ -337,7 +337,12 @@ func (lbcert *SCachedLoadbalancerCertificate) syncRemoveCloudLoadbalancerCertifi func (man *SCachedLoadbalancerCertificateManager) getLoadbalancerCertificateByRegion(provider *SCloudprovider, regionId string, localCertificateId string) (SCachedLoadbalancerCertificate, error) { certificates := []SCachedLoadbalancerCertificate{} - q := man.Query().Equals("cloudregion_id", regionId).Equals("manager_id", provider.Id).Equals("certificate_id", localCertificateId).IsFalse("pending_deleted") + q := man.Query().Equals("manager_id", provider.Id).Equals("certificate_id", localCertificateId).IsFalse("pending_deleted") + // aws 所有region共用一份证书.与region无关 + if provider.GetName() != api.CLOUD_PROVIDER_AWS { + q = q.Equals("cloudregion_id", regionId) + } + if err := db.FetchModelObjects(man, q, &certificates); err != nil { log.Errorf("failed to get lb certificate for region: %v provider: %v error: %v", regionId, provider, err) return SCachedLoadbalancerCertificate{}, err diff --git a/pkg/compute/models/loadbalancerhuaweicachedlbb.go b/pkg/compute/models/loadbalancerhuaweicachedlbb.go index 3040847817..af0c168ebb 100644 --- a/pkg/compute/models/loadbalancerhuaweicachedlbb.go +++ b/pkg/compute/models/loadbalancerhuaweicachedlbb.go @@ -16,9 +16,12 @@ package models import ( "context" + "database/sql" "fmt" "yunion.io/x/jsonutils" + "yunion.io/x/log" + "yunion.io/x/pkg/errors" "yunion.io/x/pkg/util/compare" api "yunion.io/x/onecloud/pkg/apis/compute" @@ -65,7 +68,7 @@ func (lbb *SHuaweiCachedLb) GetCustomizeColumns(context.Context, mcclient.TokenC func (man *SHuaweiCachedLbManager) GetBackendsByLocalBackendId(backendId string) ([]SHuaweiCachedLb, error) { loadbalancerBackends := []SHuaweiCachedLb{} - q := man.Query().Equals("backend_id", backendId) + q := man.Query().IsFalse("pending_deleted").Equals("backend_id", backendId) if err := db.FetchModelObjects(man, q, &loadbalancerBackends); err != nil { return nil, err } @@ -92,7 +95,7 @@ func (man *SHuaweiCachedLbManager) CreateHuaweiCachedLb(ctx context.Context, use return nil, err } - err = man.TableSpec().Insert(lbb) + err = man.TableSpec().Insert(cachedlbb) if err != nil { return nil, err @@ -209,7 +212,26 @@ func (lbb *SHuaweiCachedLb) constructFieldsFromCloudLoadbalancerBackend(extLoadb func (lbb *SHuaweiCachedLb) SyncWithCloudLoadbalancerBackend(ctx context.Context, userCred mcclient.TokenCredential, extLoadbalancerBackend cloudprovider.ICloudLoadbalancerBackend, syncOwnerId mcclient.IIdentityProvider) error { lbb.SetModelManager(HuaweiCachedLbManager, lbb) + cacheLbbg, err := lbb.GetCachedBackendGroup() + if err != nil { + return errors.Wrap(err, "HuaweiCachedLb.SyncWithCloudLoadbalancerBackend.GetCachedBackendGroup") + } + + localLbbg, err := cacheLbbg.GetLocalBackendGroup(ctx, userCred) + if err != nil { + return errors.Wrap(err, "HuaweiCachedLb.SyncWithCloudLoadbalancerBackend.GetLocalBackendGroup") + } + + locallbb, err := newLocalBackendFromCloudLoadbalancerBackend(ctx, userCred, localLbbg, extLoadbalancerBackend, syncOwnerId) + if err != nil { + return errors.Wrap(err, "HuaweiCachedLb.SyncWithCloudLoadbalancerBackend.newLocalBackendFromCloudLoadbalancerBackend") + } + diff, err := db.UpdateWithLock(ctx, lbb, func() error { + if locallbb != nil { + lbb.BackendId = locallbb.GetId() + } + return lbb.constructFieldsFromCloudLoadbalancerBackend(extLoadbalancerBackend) }) if err != nil { @@ -267,35 +289,69 @@ 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 + //} + man := LoadbalancerBackendManager - lbb := &SLoadbalancerBackend{} - lbb.SetModelManager(man, lbb) - - lbb.BackendGroupId = loadbalancerBackendgroup.Id - lbb.ExternalId = "" - - lbb.CloudregionId = loadbalancerBackendgroup.CloudregionId - lbb.ManagerId = loadbalancerBackendgroup.ManagerId - - newName, err := db.GenerateName(man, syncOwnerId, extLoadbalancerBackend.GetName()) + q := man.Query().IsFalse("pending_deleted").Equals("backend_group_id", loadbalancerBackendgroup.Id).Equals("cloudregion_id", loadbalancerBackendgroup.CloudregionId) + q = q.Equals("manager_id", loadbalancerBackendgroup.ManagerId).Equals("weight", extLoadbalancerBackend.GetWeight()).Equals("port", extLoadbalancerBackend.GetPort()) + q = q.Equals("backend_id", guest.Id) + //q = q.Equals("address", address) + lbbs := []SLoadbalancerBackend{} + err = db.FetchModelObjects(man, q, &lbbs) if err != nil { - return nil, err - } - lbb.Name = newName - - if err := lbb.constructFieldsFromCloudLoadbalancerBackend(extLoadbalancerBackend); err != nil { - return nil, err + if err != sql.ErrNoRows { + return nil, err + } } - err = man.TableSpec().Insert(lbb) + if err == sql.ErrNoRows || len(lbbs) == 0 { + lbb := &SLoadbalancerBackend{} + lbb.SetModelManager(man, lbb) - if err != nil { - return nil, err + lbb.BackendGroupId = loadbalancerBackendgroup.Id + lbb.ExternalId = "" + + lbb.CloudregionId = loadbalancerBackendgroup.CloudregionId + lbb.ManagerId = loadbalancerBackendgroup.ManagerId + + baseName := extLoadbalancerBackend.GetName() + if len(baseName) == 0 { + baseName = "backend" + } + + newName, err := db.GenerateName(man, syncOwnerId, extLoadbalancerBackend.GetName()) + if err != nil { + return nil, err + } + lbb.Name = newName + + if err := lbb.constructFieldsFromCloudLoadbalancerBackend(extLoadbalancerBackend); err != nil { + return nil, err + } + + err = man.TableSpec().Insert(lbb) + + if err != nil { + return nil, err + } + + SyncCloudProject(userCred, lbb, syncOwnerId, extLoadbalancerBackend, loadbalancerBackendgroup.ManagerId) + + db.OpsLog.LogEvent(lbb, db.ACT_CREATE, lbb.GetShortDesc(ctx), userCred) + return lbb, nil + } else if len(lbbs) == 1 { + return &lbbs[0], nil + } else { + log.Errorf("duplicate lbb found %#v", lbbs) + return nil, errors.Wrap(fmt.Errorf("duplicate lbb found"), "newLocalBackendFromCloudLoadbalancerBackend") } - - SyncCloudProject(userCred, lbb, syncOwnerId, extLoadbalancerBackend, loadbalancerBackendgroup.ManagerId) - - db.OpsLog.LogEvent(lbb, db.ACT_CREATE, lbb.GetShortDesc(ctx), userCred) - - return lbb, nil } diff --git a/pkg/compute/models/loadbalancerhuaweicachedlbbg.go b/pkg/compute/models/loadbalancerhuaweicachedlbbg.go index 3c990d93cd..a3c1bbf10b 100644 --- a/pkg/compute/models/loadbalancerhuaweicachedlbbg.go +++ b/pkg/compute/models/loadbalancerhuaweicachedlbbg.go @@ -17,10 +17,14 @@ package models import ( "context" "fmt" + "strconv" + "strings" "yunion.io/x/jsonutils" "yunion.io/x/log" + "yunion.io/x/pkg/errors" "yunion.io/x/pkg/util/compare" + "yunion.io/x/pkg/utils" api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db" @@ -129,7 +133,7 @@ func (lbbg *SHuaweiCachedLbbg) GetICloudLoadbalancerBackendGroup() (cloudprovide func (man *SHuaweiCachedLbbgManager) GetUsableCachedBackendGroups(backendGroupId string, protocolType string) ([]SHuaweiCachedLbbg, error) { ret := []SHuaweiCachedLbbg{} - err := man.TableSpec().Query().Equals("backend_group_id", backendGroupId).Equals("protocol_type", protocolType).IsNullOrEmpty("associated_id").IsNotEmpty("external_id").All(&ret) + err := man.Query().IsFalse("pending_deleted").Equals("backend_group_id", backendGroupId).Equals("protocol_type", protocolType).IsNullOrEmpty("associated_id").IsNotEmpty("external_id").All(&ret) if err != nil { return ret, err } @@ -152,17 +156,18 @@ func (man *SHuaweiCachedLbbgManager) GetUsableCachedBackendGroup(backendGroupId func (man *SHuaweiCachedLbbgManager) GetCachedBackendGroupByAssociateId(associateId string) (*SHuaweiCachedLbbg, error) { ret := &SHuaweiCachedLbbg{} - err := man.TableSpec().Query().Equals("associated_id", associateId).First(ret) + err := man.Query().IsFalse("pending_deleted").Equals("associated_id", associateId).First(ret) if err != nil { return nil, err } + ret.SetModelManager(man, ret) return ret, nil } func (man *SHuaweiCachedLbbgManager) GetCachedBackendGroups(backendGroupId string) ([]SHuaweiCachedLbbg, error) { ret := []SHuaweiCachedLbbg{} - err := man.TableSpec().Query().Equals("backend_group_id", backendGroupId).All(&ret) + err := man.Query().IsFalse("pending_deleted").Equals("backend_group_id", backendGroupId).All(&ret) if err != nil { return nil, err } @@ -172,7 +177,7 @@ func (man *SHuaweiCachedLbbgManager) GetCachedBackendGroups(backendGroupId strin func (man *SHuaweiCachedLbbgManager) getLoadbalancerBackendgroupsByLoadbalancer(lb *SLoadbalancer) ([]SHuaweiCachedLbbg, error) { lbbgs := []SHuaweiCachedLbbg{} - q := man.Query().Equals("loadbalancer_id", lb.Id) + q := man.Query().IsFalse("pending_deleted").Equals("loadbalancer_id", lb.Id) if err := db.FetchModelObjects(man, q, &lbbgs); err != nil { log.Errorf("failed to get lbbgs for lb: %s error: %v", lb.Name, err) return nil, err @@ -257,11 +262,66 @@ func (lbbg *SHuaweiCachedLbbg) syncRemoveCloudLoadbalancerBackendgroup(ctx conte return err } +func (lbbg *SHuaweiCachedLbbg) isBackendsMatch(backends []SLoadbalancerBackend, ibackends []cloudprovider.ICloudLoadbalancerBackend) bool { + if len(ibackends) != len(backends) { + return false + } + + locals := []string{} + remotes := []string{} + + for i := range backends { + guest := backends[i].GetGuest() + seg := strings.Join([]string{guest.ExternalId, strconv.Itoa(backends[i].Weight), strconv.Itoa(backends[i].Port)}, "/") + locals = append(locals, seg) + } + + for i := range ibackends { + ibackend := ibackends[i] + seg := strings.Join([]string{ibackend.GetBackendId(), strconv.Itoa(ibackend.GetWeight()), strconv.Itoa(ibackend.GetPort())}, "/") + remotes = append(remotes, seg) + } + + for i := range remotes { + if !utils.IsInStringArray(remotes[i], locals) { + return false + } + } + + return true +} + func (lbbg *SHuaweiCachedLbbg) SyncWithCloudLoadbalancerBackendgroup(ctx context.Context, userCred mcclient.TokenCredential, lb *SLoadbalancer, extLoadbalancerBackendgroup cloudprovider.ICloudLoadbalancerBackendGroup, syncOwnerId mcclient.IIdentityProvider) error { lbbg.SetModelManager(HuaweiCachedLbbgManager, lbbg) + + ibackends, err := extLoadbalancerBackendgroup.GetILoadbalancerBackends() + if err != nil { + return errors.Wrap(err, "HuaweiCachedLbbg.SyncWithCloudLoadbalancerBackendgroup.GetILoadbalancerBackends") + } + + localLbbg, err := lbbg.GetLocalBackendGroup(ctx, userCred) + if err != nil { + return errors.Wrap(err, "HuaweiCachedLbbg.SyncWithCloudLoadbalancerBackendgroup.GetLocalBackendGroup") + } + + backends, err := localLbbg.GetBackends() + if err != nil { + return errors.Wrap(err, "HuaweiCachedLbbg.SyncWithCloudLoadbalancerBackendgroup.GetBackends") + } + + var newLocalLbbg *SLoadbalancerBackendGroup + if !lbbg.isBackendsMatch(backends, ibackends) { + newLocalLbbg, err = newLocalBackendgroupFromCloudLoadbalancerBackendgroup(ctx, userCred, lb, extLoadbalancerBackendgroup, syncOwnerId) + if err != nil { + return errors.Wrap(err, "HuaweiCachedLbbg.SyncWithCloudLoadbalancerBackendgroup.newLocalBackendgroupFromCloudLoadbalancerBackendgroup") + } + } + diff, err := db.UpdateWithLock(ctx, lbbg, func() error { lbbg.Status = extLoadbalancerBackendgroup.GetStatus() - + if newLocalLbbg != nil { + lbbg.BackendGroupId = newLocalLbbg.GetId() + } return nil }) if err != nil { diff --git a/pkg/compute/models/loadbalancerlistenerrules.go b/pkg/compute/models/loadbalancerlistenerrules.go index 3c8b085eec..781a5b260f 100644 --- a/pkg/compute/models/loadbalancerlistenerrules.go +++ b/pkg/compute/models/loadbalancerlistenerrules.go @@ -20,6 +20,7 @@ import ( "yunion.io/x/jsonutils" "yunion.io/x/log" + "yunion.io/x/pkg/errors" "yunion.io/x/pkg/util/compare" "yunion.io/x/sqlchemy" @@ -647,6 +648,39 @@ func (lbr *SLoadbalancerListenerRule) constructFieldsFromCloudListenerRule(userC } } +func (lbr *SLoadbalancerListenerRule) updateCachedLoadbalancerBackendGroupAssociate(ctx context.Context, extRule cloudprovider.ICloudLoadbalancerListenerRule) error { + exteralLbbgId := extRule.GetBackendGroupId() + if len(exteralLbbgId) == 0 { + return nil + } + + switch lbr.GetProviderName() { + case api.CLOUD_PROVIDER_HUAWEI: + _group, err := db.FetchByExternalId(HuaweiCachedLbbgManager, exteralLbbgId) + if err != nil { + return fmt.Errorf("Fetch huawei loadbalancer backendgroup by external id %s failed: %s", exteralLbbgId, err) + } + + if _group != nil { + group := _group.(*SHuaweiCachedLbbg) + if group.AssociatedId != lbr.Id { + _, err := db.UpdateWithLock(ctx, group, func() error { + group.AssociatedId = lbr.Id + group.AssociatedType = api.LB_ASSOCIATE_TYPE_RULE + return nil + }) + if err != nil { + return errors.Wrap(err, "LoadbalancerListener.updateCachedLoadbalancerBackendGroupAssociate") + } + } + } + default: + return nil + } + + return nil +} + func (man *SLoadbalancerListenerRuleManager) newFromCloudLoadbalancerListenerRule(ctx context.Context, userCred mcclient.TokenCredential, listener *SLoadbalancerListener, extRule cloudprovider.ICloudLoadbalancerListenerRule, syncOwnerId mcclient.IIdentityProvider) (*SLoadbalancerListenerRule, error) { lbr := &SLoadbalancerListenerRule{} lbr.SetModelManager(man, lbr) @@ -670,22 +704,9 @@ func (man *SLoadbalancerListenerRuleManager) newFromCloudLoadbalancerListenerRul return nil, err } - groupId := extRule.GetBackendGroupId() - if lbr.GetProviderName() == api.CLOUD_PROVIDER_HUAWEI && len(groupId) > 0 { - group, err := db.FetchByExternalId(HuaweiCachedLbbgManager, groupId) - if err != nil { - log.Errorf("Fetch huawei loadbalancer backendgroup by external id %s failed: %s", groupId, err) - } - - cachedGroup := group.(*SHuaweiCachedLbbg) - _, err = db.UpdateWithLock(context.Background(), cachedGroup, func() error { - cachedGroup.AssociatedId = lbr.GetId() - cachedGroup.AssociatedType = api.LB_ASSOCIATE_TYPE_RULE - return nil - }) - if err != nil { - log.Errorf("Update huawei loadbalancer backendgroup cache %s failed: %s", groupId, err) - } + err = lbr.updateCachedLoadbalancerBackendGroupAssociate(ctx, extRule) + if err != nil { + return nil, errors.Wrap(err, "LoadbalancerListenerRuleManager.newFromCloudLoadbalancerListenerRule") } SyncCloudProject(userCred, lbr, syncOwnerId, extRule, listener.ManagerId) @@ -719,6 +740,12 @@ func (lbr *SLoadbalancerListenerRule) SyncWithCloudLoadbalancerListenerRule(ctx if err != nil { return err } + + err = lbr.updateCachedLoadbalancerBackendGroupAssociate(ctx, extRule) + if err != nil { + return errors.Wrap(err, "LoadbalancerListenerRule.SyncWithCloudLoadbalancerListenerRule") + } + db.OpsLog.LogSyncUpdate(lbr, diff, userCred) SyncCloudProject(userCred, lbr, syncOwnerId, extRule, listener.ManagerId) diff --git a/pkg/compute/models/loadbalancerlisteners.go b/pkg/compute/models/loadbalancerlisteners.go index b17b08c561..19170036d9 100644 --- a/pkg/compute/models/loadbalancerlisteners.go +++ b/pkg/compute/models/loadbalancerlisteners.go @@ -56,8 +56,8 @@ func init() { } type SLoadbalancerHTTPRateLimiter struct { - HTTPRequestRate int `nullable:"true" list:"user" create:"optional" update:"user"` - HTTPRequestRatePerSrc int `nullable:"true" list:"user" create:"optional" update:"user"` + HTTPRequestRate int `nullable:"true" list:"user" create:"optional" update:"user"` // 限定监听接收请示速率 + HTTPRequestRatePerSrc int `nullable:"true" list:"user" create:"optional" update:"user"` // 源IP监听请求最大速率 } type SLoadbalancerRateLimiter struct { @@ -65,20 +65,20 @@ type SLoadbalancerRateLimiter struct { } type SLoadbalancerHealthCheck struct { - HealthCheck string `width:"16" charset:"ascii" nullable:"true" list:"user" create:"optional" update:"user"` - HealthCheckType string `width:"16" charset:"ascii" nullable:"true" list:"user" create:"optional" update:"user"` + HealthCheck string `width:"16" charset:"ascii" nullable:"true" list:"user" create:"optional" update:"user"` // 健康检查开启状态 on|off + HealthCheckType string `width:"16" charset:"ascii" nullable:"true" list:"user" create:"optional" update:"user"` // 健康检查协议 HTTP|TCP - HealthCheckDomain string `charset:"ascii" nullable:"true" list:"user" create:"optional" update:"user"` - HealthCheckURI string `charset:"ascii" nullable:"true" list:"user" create:"optional" update:"user"` - HealthCheckHttpCode string `charset:"ascii" nullable:"true" list:"user" create:"optional" update:"user"` + HealthCheckDomain string `charset:"ascii" nullable:"true" list:"user" create:"optional" update:"user"` // 健康检查域名 yunion.cn + HealthCheckURI string `charset:"ascii" nullable:"true" list:"user" create:"optional" update:"user"` // 健康检查路径 / + HealthCheckHttpCode string `charset:"ascii" nullable:"true" list:"user" create:"optional" update:"user"` // HTTP正常状态码 http_2xx,http_3xx - HealthCheckRise int `nullable:"true" list:"user" create:"optional" update:"user"` - HealthCheckFall int `nullable:"true" list:"user" create:"optional" update:"user"` - HealthCheckTimeout int `nullable:"true" list:"user" create:"optional" update:"user"` - HealthCheckInterval int `nullable:"true" list:"user" create:"optional" update:"user"` + HealthCheckRise int `nullable:"true" list:"user" create:"optional" update:"user"` // 健康检查健康阈值 3秒 + HealthCheckFall int `nullable:"true" list:"user" create:"optional" update:"user"` // 健康检查不健康阈值 15秒 + HealthCheckTimeout int `nullable:"true" list:"user" create:"optional" update:"user"` // 健康检查超时时间 10秒 + HealthCheckInterval int `nullable:"true" list:"user" create:"optional" update:"user"` // 健康检查间隔时间 5秒 - HealthCheckReq string `list:"user" create:"optional" update:"user"` - HealthCheckExp string `list:"user" create:"optional" update:"user"` + HealthCheckReq string `list:"user" create:"optional" update:"user"` // UDP监听健康检查的请求串 + HealthCheckExp string `list:"user" create:"optional" update:"user"` // UDP监听健康检查的响应串 } type SLoadbalancerTCPListener struct{} @@ -86,13 +86,13 @@ type SLoadbalancerUDPListener struct{} // TODO sensible default for knobs type SLoadbalancerHTTPListener struct { - StickySession string `width:"16" charset:"ascii" nullable:"true" list:"user" create:"optional" update:"user"` - StickySessionType string `width:"16" charset:"ascii" nullable:"true" list:"user" create:"optional" update:"user"` - StickySessionCookie string `width:"128" charset:"ascii" nullable:"true" list:"user" create:"optional" update:"user"` - StickySessionCookieTimeout int `nullable:"true" list:"user" create:"optional" update:"user"` + StickySession string `width:"16" charset:"ascii" nullable:"true" list:"user" create:"optional" update:"user"` // 会话保持开启状态 on|off + StickySessionType string `width:"16" charset:"ascii" nullable:"true" list:"user" create:"optional" update:"user"` // Cookie处理方式 insert(植入cookie)|server(重写cookie) + StickySessionCookie string `width:"128" charset:"ascii" nullable:"true" list:"user" create:"optional" update:"user"` // Cookie名称 + StickySessionCookieTimeout int `nullable:"true" list:"user" create:"optional" update:"user"` // 会话超时时间 - XForwardedFor bool `nullable:"true" list:"user" create:"optional" update:"user"` - Gzip bool `nullable:"true" list:"user" create:"optional" update:"user"` + XForwardedFor bool `nullable:"true" list:"user" create:"optional" update:"user"` // 获取客户端真实IP + Gzip bool `nullable:"true" list:"user" create:"optional" update:"user"` // Gzip数据压缩 } // TODO @@ -125,10 +125,10 @@ type SLoadbalancerListener struct { SendProxy string `width:"16" charset:"ascii" nullable:"false" list:"user" create:"optional" update:"user" default:"off"` - ClientRequestTimeout int `nullable:"true" list:"user" create:"optional" update:"user"` - ClientIdleTimeout int `nullable:"true" list:"user" create:"optional" update:"user"` - BackendConnectTimeout int `nullable:"true" list:"user" create:"optional" update:"user"` - BackendIdleTimeout int `nullable:"true" list:"user" create:"optional" update:"user"` + ClientRequestTimeout int `nullable:"true" list:"user" create:"optional" update:"user"` // 连接请求超时时间 + ClientIdleTimeout int `nullable:"true" list:"user" create:"optional" update:"user"` // 连接空闲超时时间 + BackendConnectTimeout int `nullable:"true" list:"user" create:"optional" update:"user"` // 后端连接超时时间 + BackendIdleTimeout int `nullable:"true" list:"user" create:"optional" update:"user"` // 后端连接空闲时间 AclStatus string `width:"16" charset:"ascii" nullable:"true" list:"user" create:"optional" update:"user"` AclType string `width:"16" charset:"ascii" nullable:"true" list:"user" create:"optional" update:"user"` @@ -502,6 +502,11 @@ func (lblis *SLoadbalancerListener) GetLoadbalancerListenerParams() (*cloudprovi EstablishedTimeout: lblis.BackendConnectTimeout, AccessControlListStatus: lblis.AclStatus, + ClientRequestTimeout: lblis.ClientRequestTimeout, + ClientIdleTimeout: lblis.ClientIdleTimeout, + BackendIdleTimeout: lblis.BackendIdleTimeout, + BackendConnectTimeout: lblis.BackendConnectTimeout, + HealthCheckReq: lblis.HealthCheckReq, HealthCheckExp: lblis.HealthCheckExp, @@ -845,18 +850,17 @@ func (lblis *SLoadbalancerListener) constructFieldsFromCloudListener(userCred mc } } - if len(lblis.BackendGroupId) == 0 && len(groupId) == 0 { - lblis.BackendGroupId = lb.BackendGroupId - } else if lblis.GetProviderName() == api.CLOUD_PROVIDER_HUAWEI { + switch lblis.GetProviderName() { + case api.CLOUD_PROVIDER_HUAWEI: if len(groupId) > 0 { group, err := db.FetchByExternalId(HuaweiCachedLbbgManager, groupId) if err != nil { log.Errorf("Fetch huawei loadbalancer backendgroup by external id %s failed: %s", groupId, err) + } else { + lblis.BackendGroupId = group.(*SHuaweiCachedLbbg).BackendGroupId } - - lblis.BackendGroupId = group.(*SHuaweiCachedLbbg).BackendGroupId } - } else if lblis.GetProviderName() == api.CLOUD_PROVIDER_AWS { + case api.CLOUD_PROVIDER_AWS: if len(groupId) > 0 { group, err := db.FetchByExternalId(AwsCachedLbbgManager, groupId) if err != nil { @@ -865,11 +869,48 @@ func (lblis *SLoadbalancerListener) constructFieldsFromCloudListener(userCred mc lblis.BackendGroupId = group.(*SAwsCachedLbbg).BackendGroupId } } - } else if group, err := db.FetchByExternalId(LoadbalancerBackendGroupManager, groupId); err == nil { - lblis.BackendGroupId = group.GetId() + default: + if len(lblis.BackendGroupId) == 0 && len(groupId) == 0 { + lblis.BackendGroupId = lb.BackendGroupId + } else if group, err := db.FetchByExternalId(LoadbalancerBackendGroupManager, groupId); err == nil { + lblis.BackendGroupId = group.GetId() + } } } +func (lblis *SLoadbalancerListener) updateCachedLoadbalancerBackendGroupAssociate(ctx context.Context, extListener cloudprovider.ICloudLoadbalancerListener) error { + exteralLbbgId := extListener.GetBackendGroupId() + if len(exteralLbbgId) == 0 { + return nil + } + + switch lblis.GetProviderName() { + case api.CLOUD_PROVIDER_HUAWEI: + _group, err := db.FetchByExternalId(HuaweiCachedLbbgManager, exteralLbbgId) + if err != nil { + return fmt.Errorf("Fetch huawei loadbalancer backendgroup by external id %s failed: %s", exteralLbbgId, err) + } + + if _group != nil { + group := _group.(*SHuaweiCachedLbbg) + if group.AssociatedId != lblis.Id { + _, err := db.UpdateWithLock(ctx, group, func() error { + group.AssociatedId = lblis.Id + group.AssociatedType = api.LB_ASSOCIATE_TYPE_LISTENER + return nil + }) + if err != nil { + return errors.Wrap(err, "LoadbalancerListener.updateCachedLoadbalancerBackendGroupAssociate") + } + } + } + default: + return nil + } + + return nil +} + func (lblis *SLoadbalancerListener) syncRemoveCloudLoadbalancerListener(ctx context.Context, userCred mcclient.TokenCredential) error { lockman.LockObject(ctx, lblis) defer lockman.ReleaseObject(ctx, lblis) @@ -892,6 +933,11 @@ func (lblis *SLoadbalancerListener) SyncWithCloudLoadbalancerListener(ctx contex return err } + err = lblis.updateCachedLoadbalancerBackendGroupAssociate(ctx, extListener) + if err != nil { + return errors.Wrap(err, "LoadbalancerListener.SyncWithCloudLoadbalancerListener") + } + db.OpsLog.LogSyncUpdate(lblis, diff, userCred) SyncCloudProject(userCred, lblis, syncOwnerId, extListener, lblis.ManagerId) @@ -919,22 +965,9 @@ func (man *SLoadbalancerListenerManager) newFromCloudLoadbalancerListener(ctx co return nil, err } - groupId := extListener.GetBackendGroupId() - if lblis.GetProviderName() == api.CLOUD_PROVIDER_HUAWEI && len(groupId) > 0 { - group, err := db.FetchByExternalId(HuaweiCachedLbbgManager, groupId) - if err != nil { - log.Errorf("Fetch huawei loadbalancer backendgroup by external id %s failed: %s", groupId, err) - } - - cachedGroup := group.(*SHuaweiCachedLbbg) - _, err = db.UpdateWithLock(context.Background(), cachedGroup, func() error { - cachedGroup.AssociatedId = lblis.GetId() - cachedGroup.AssociatedType = api.LB_ASSOCIATE_TYPE_LISTENER - return nil - }) - if err != nil { - log.Errorf("Update huawei loadbalancer backendgroup cache %s failed: %s", groupId, err) - } + err = lblis.updateCachedLoadbalancerBackendGroupAssociate(ctx, extListener) + if err != nil { + return nil, errors.Wrap(err, "LoadbalancerListener.newFromCloudLoadbalancerListener") } SyncCloudProject(userCred, lblis, syncOwnerId, extListener, lblis.ManagerId) diff --git a/pkg/compute/regiondrivers/aws.go b/pkg/compute/regiondrivers/aws.go index 0faebbf621..366a2d72ec 100644 --- a/pkg/compute/regiondrivers/aws.go +++ b/pkg/compute/regiondrivers/aws.go @@ -16,6 +16,7 @@ package regiondrivers import ( "context" + "database/sql" "fmt" "regexp" "strings" @@ -963,18 +964,18 @@ func (self *SAwsRegionDriver) RequestCreateLoadbalancerListener(ctx context.Cont cert, err := models.LoadbalancerCertificateManager.FetchById(certId) if err != nil { - return nil, err + return nil, errors.Wrap(err, "awsRegionDriver.RequestCreateLoadbalancerListener.FetchById") } lbcert, err := models.CachedLoadbalancerCertificateManager.GetOrCreateCachedCertificate(ctx, userCred, provider, lblis, cert.(*models.SLoadbalancerCertificate)) if err != nil { - return nil, err + return nil, errors.Wrap(err, "awsRegionDriver.RequestCreateLoadbalancerListener.GetOrCreateCachedCertificate") } if len(lbcert.ExternalId) == 0 { _, err = self.createLoadbalancerCertificate(ctx, userCred, lbcert) if err != nil { - return nil, err + return nil, errors.Wrap(err, "awsRegionDriver.RequestCreateLoadbalancerListener.createLoadbalancerCertificate") } } @@ -996,7 +997,7 @@ func (self *SAwsRegionDriver) RequestCreateLoadbalancerListener(ctx context.Cont params, err := lbbg.GetAwsBackendGroupParams(lblis, nil) if err != nil { - return nil, err + return nil, errors.Wrap(err, "awsRegionDriver.RequestCreateLoadbalancerListener.GetAwsBackendGroupParams") } group, _ := models.AwsCachedLbbgManager.GetUsableCachedBackendGroup(lblis.LoadbalancerId, lblis.BackendGroupId, lblis.ListenerType, lblis.HealthCheckType, lblis.HealthCheckInterval) @@ -1004,46 +1005,46 @@ func (self *SAwsRegionDriver) RequestCreateLoadbalancerListener(ctx context.Cont // 服务器组存在 ilbbg, err := group.GetICloudLoadbalancerBackendGroup() if err != nil { - return nil, err + return nil, errors.Wrap(err, "awsRegionDriver.RequestCreateLoadbalancerListener.GetICloudLoadbalancerBackendGroup") } // 服务器组已经存在,直接同步即可 if err := ilbbg.Sync(params); err != nil { - return nil, err + return nil, errors.Wrap(err, "awsRegionDriver.RequestCreateLoadbalancerListener.Sync") } } else { backends, err := lbbg.GetBackendsParams() if err != nil { - return nil, err + return nil, errors.Wrap(err, "awsRegionDriver.RequestCreateLoadbalancerListener.GetBackendsParams") } // 服务器组不存在 _, err = self.createLoadbalancerBackendGroup(ctx, userCred, loadbalancer, lblis, lbbg, backends) if err != nil { - return nil, err + return nil, errors.Wrap(err, "awsRegionDriver.RequestCreateLoadbalancerListener.createLoadbalancerBackendGroup") } } } params, err := lblis.GetAwsLoadbalancerListenerParams() if err != nil { - return nil, err + return nil, errors.Wrap(err, "awsRegionDriver.RequestCreateLoadbalancerListener.GetAwsLoadbalancerListenerParams") } iRegion, err := loadbalancer.GetIRegion() if err != nil { - return nil, err + return nil, errors.Wrap(err, "awsRegionDriver.RequestCreateLoadbalancerListener.GetIRegion") } iLoadbalancer, err := iRegion.GetILoadBalancerById(loadbalancer.ExternalId) if err != nil { - return nil, err + return nil, errors.Wrap(err, "awsRegionDriver.RequestCreateLoadbalancerListener.GetILoadBalancerById") } iListener, err := iLoadbalancer.CreateILoadBalancerListener(params) if err != nil { - return nil, err + return nil, errors.Wrap(err, "awsRegionDriver.RequestCreateLoadbalancerListener.CreateILoadBalancerListener") } lblis.SetModelManager(models.LoadbalancerListenerManager, lblis) if err := db.SetExternalId(lblis, userCred, iListener.GetGlobalId()); err != nil { - return nil, err + return nil, errors.Wrap(err, "awsRegionDriver.RequestCreateLoadbalancerListener.SetExternalId") } return nil, lblis.SyncWithCloudLoadbalancerListener(ctx, userCred, loadbalancer, iListener, loadbalancer.GetOwnerId()) @@ -1078,7 +1079,7 @@ func (self *SAwsRegionDriver) RequestCreateLoadbalancerListenerRule(ctx context. Condition: lbr.Condition, } - group, err := models.AwsCachedLbbgManager.GetUsableCachedBackendGroup(listener.LoadbalancerId, listener.BackendGroupId, listener.ListenerType, listener.HealthCheckType, listener.HealthCheckInterval) + group, err := models.AwsCachedLbbgManager.GetUsableCachedBackendGroup(listener.LoadbalancerId, lbr.BackendGroupId, listener.ListenerType, listener.HealthCheckType, listener.HealthCheckInterval) if err != nil { return nil, err } @@ -1098,6 +1099,63 @@ func (self *SAwsRegionDriver) RequestCreateLoadbalancerListenerRule(ctx context. return nil } +func (self *SAwsRegionDriver) RequestDeleteLoadbalancerBackendGroup(ctx context.Context, userCred mcclient.TokenCredential, lbbg *models.SLoadbalancerBackendGroup, task taskman.ITask) error { + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + if jsonutils.QueryBoolean(task.GetParams(), "purge", false) { + return nil, nil + } + + iRegion, err := lbbg.GetIRegion() + if err != nil { + return nil, err + } + loadbalancer := lbbg.GetLoadbalancer() + if loadbalancer == nil { + return nil, fmt.Errorf("failed to find loadbalancer for backendgroup %s", lbbg.Name) + } + iLoadbalancer, err := iRegion.GetILoadBalancerById(loadbalancer.ExternalId) + if err != nil { + if err == cloudprovider.ErrNotFound { + return nil, nil + } + return nil, err + } + + cachedLbbgs, err := models.AwsCachedLbbgManager.GetCachedBackendGroups(lbbg.GetId()) + if err != nil { + return nil, err + } + + for i := range cachedLbbgs { + cachedLbbg := cachedLbbgs[i] + if len(cachedLbbg.ExternalId) == 0 { + continue + } + + iLoadbalancerBackendGroup, err := iLoadbalancer.GetILoadBalancerBackendGroupById(cachedLbbg.ExternalId) + if err != nil { + if err == cloudprovider.ErrNotFound { + return nil, nil + } + return nil, err + } + + err = iLoadbalancerBackendGroup.Delete() + if err != nil { + return nil, err + } + + err = cachedLbbg.Delete(ctx, userCred) + if err != nil { + return nil, err + } + } + + return nil, nil + }) + return nil +} + func (self *SAwsRegionDriver) RequestDeleteLoadbalancer(ctx context.Context, userCred mcclient.TokenCredential, lb *models.SLoadbalancer, task taskman.ITask) error { taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { if jsonutils.QueryBoolean(task.GetParams(), "purge", false) { @@ -1154,18 +1212,18 @@ func (self *SAwsRegionDriver) RequestSyncLoadbalancerListener(ctx context.Contex cert, err := models.LoadbalancerCertificateManager.FetchById(certId) if err != nil { - return nil, err + return nil, errors.Wrap(err, "awsRegionDriver.RequestSyncLoadbalancerListener.FetchById") } lbcert, err := models.CachedLoadbalancerCertificateManager.GetOrCreateCachedCertificate(ctx, userCred, provider, lblis, cert.(*models.SLoadbalancerCertificate)) if err != nil { - return nil, err + return nil, errors.Wrap(err, "awsRegionDriver.RequestSyncLoadbalancerListener.GetOrCreateCachedCertificate") } if len(lbcert.ExternalId) == 0 { _, err = self.createLoadbalancerCertificate(ctx, userCred, lbcert) if err != nil { - return nil, err + return nil, errors.Wrap(err, "awsRegionDriver.RequestSyncLoadbalancerListener.createLoadbalancerCertificate") } } @@ -1181,7 +1239,7 @@ func (self *SAwsRegionDriver) RequestSyncLoadbalancerListener(ctx context.Contex params, err := lblis.GetAwsLoadbalancerListenerParams() if err != nil { - return nil, err + return nil, errors.Wrap(err, "awsRegionDriver.RequestSyncLoadbalancerListener.GetAwsLoadbalancerListenerParams") } loadbalancer := lblis.GetLoadbalancer() if loadbalancer == nil { @@ -1189,27 +1247,57 @@ func (self *SAwsRegionDriver) RequestSyncLoadbalancerListener(ctx context.Contex } iRegion, err := loadbalancer.GetIRegion() if err != nil { - return nil, err + return nil, errors.Wrap(err, "awsRegionDriver.RequestSyncLoadbalancerListener.GetIRegion") } iLoadbalancer, err := iRegion.GetILoadBalancerById(loadbalancer.ExternalId) if err != nil { - return nil, err + return nil, errors.Wrap(err, "awsRegionDriver.RequestSyncLoadbalancerListener.GetILoadBalancerById") } iListener, err := iLoadbalancer.GetILoadBalancerListenerById(lblis.ExternalId) if err != nil { - return nil, err + return nil, errors.Wrap(err, "awsRegionDriver.RequestSyncLoadbalancerListener.GetILoadBalancerListenerById") } if err := iListener.Sync(params); err != nil { - return nil, err + return nil, errors.Wrap(err, "awsRegionDriver.RequestSyncLoadbalancerListener.Sync") } if err := iListener.Refresh(); err != nil { - return nil, err + return nil, errors.Wrap(err, "awsRegionDriver.RequestSyncLoadbalancerListener.Refresh") } return nil, lblis.SyncWithCloudLoadbalancerListener(ctx, userCred, loadbalancer, iListener, nil) }) return nil } +func (self *SAwsRegionDriver) RequestSyncLoadbalancerBackendGroup(ctx context.Context, userCred mcclient.TokenCredential, lblis *models.SLoadbalancerListener, lbbg *models.SLoadbalancerBackendGroup, task taskman.ITask) error { + lb := lblis.GetLoadbalancer() + if lb == nil { + return errors.Wrap(fmt.Errorf("listener %s related loadbalancer not found", lblis.GetId()), "AwsRegionDriver.RequestSyncLoadbalancerBackendGroup.GetLoadbalancer") + } + + cachedLbbg, err := models.AwsCachedLbbgManager.GetUsableCachedBackendGroup(lb.GetId(), lblis.BackendGroupId, lblis.ListenerType, lblis.HealthCheckType, lblis.HealthCheckInterval) + if err != nil { + if err != sql.ErrNoRows { + return errors.Wrap(err, "AwsRegionDriver.RequestSyncLoadbalancerBackendGroup.GetLoadbalancer") + } + } + + if cachedLbbg == nil { + backends, err := lbbg.GetBackendsParams() + if err != nil { + return errors.Wrap(err, "AwsRegionDriver.RequestSyncLoadbalancerBackendGroup.GetBackendsParams") + } + + taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { + return self.createLoadbalancerBackendGroup(ctx, userCred, lb, lblis, lbbg, backends) + }) + + return nil + } + + task.ScheduleRun(nil) + return nil +} + func (self *SAwsRegionDriver) IsSecurityGroupBelongVpc() bool { return true } diff --git a/pkg/compute/regiondrivers/huawei.go b/pkg/compute/regiondrivers/huawei.go index dfbb2f85a7..c58b493fb1 100644 --- a/pkg/compute/regiondrivers/huawei.go +++ b/pkg/compute/regiondrivers/huawei.go @@ -352,7 +352,7 @@ func (self *SHuaWeiRegionDriver) ValidateCreateLoadbalancerListenerData(ctx cont "health_check_path": validators.NewURLPathValidator("health_check_path").Default(""), "health_check_http_code": validators.NewStringMultiChoicesValidator("health_check_http_code", api.LB_HEALTH_CHECK_HTTP_CODES).Sep(",").Default(api.LB_HEALTH_CHECK_HTTP_CODE_DEFAULT), - "health_check_fall": validators.NewRangeValidator("health_check_fall", 1, 10).Default(3), + "health_check_rise": validators.NewRangeValidator("health_check_rise", 1, 10).Default(3), "health_check_timeout": validators.NewRangeValidator("health_check_timeout", 1, 50).Default(10), "health_check_interval": validators.NewRangeValidator("health_check_interval", 1, 50).Default(5), } @@ -361,8 +361,8 @@ func (self *SHuaWeiRegionDriver) ValidateCreateLoadbalancerListenerData(ctx cont return nil, err } - if t, _ := data.Int("health_check_fall"); t > 0 { - data.Set("health_check_rise", jsonutils.NewInt(t)) + if t, _ := data.Int("health_check_rise"); t > 0 { + data.Set("health_check_fall", jsonutils.NewInt(t)) } // acl check @@ -504,7 +504,9 @@ func (self *SHuaWeiRegionDriver) ValidateUpdateLoadbalancerListenerData(ctx cont aclTypeV.Default(lblis.AclType) } aclV := validators.NewModelIdOrNameValidator("acl", "loadbalanceracl", ownerId) - aclV.Default(lblis.AclId) + if len(lblis.AclId) > 0 { + aclV.Default(lblis.AclId) + } certV := validators.NewModelIdOrNameValidator("certificate", "loadbalancercertificate", ownerId) tlsCipherPolicyV := validators.NewStringChoicesValidator("tls_cipher_policy", api.LB_TLS_CIPHER_POLICIES).Default(api.LB_TLS_CIPHER_POLICY_1_2) @@ -528,7 +530,7 @@ func (self *SHuaWeiRegionDriver) ValidateUpdateLoadbalancerListenerData(ctx cont "health_check_path": validators.NewURLPathValidator("health_check_path").Default(""), "health_check_http_code": validators.NewStringMultiChoicesValidator("health_check_http_code", api.LB_HEALTH_CHECK_HTTP_CODES).Sep(",").Default(api.LB_HEALTH_CHECK_HTTP_CODE_DEFAULT), - "health_check_fall": validators.NewRangeValidator("health_check_fall", 1, 10).Default(3), + "health_check_rise": validators.NewRangeValidator("health_check_rise", 1, 10).Default(3), "health_check_timeout": validators.NewRangeValidator("health_check_timeout", 1, 50).Default(10), "health_check_interval": validators.NewRangeValidator("health_check_interval", 1, 50).Default(5), @@ -544,8 +546,8 @@ func (self *SHuaWeiRegionDriver) ValidateUpdateLoadbalancerListenerData(ctx cont return nil, err } - if t, _ := data.Int("health_check_fall"); t > 0 { - data.Set("health_check_rise", jsonutils.NewInt(t)) + if t, _ := data.Int("health_check_rise"); t > 0 { + data.Set("health_check_fall", jsonutils.NewInt(t)) } { @@ -558,6 +560,77 @@ func (self *SHuaWeiRegionDriver) ValidateUpdateLoadbalancerListenerData(ctx cont return self.SManagedVirtualizationRegionDriver.ValidateUpdateLoadbalancerListenerData(ctx, userCred, data, lblis, backendGroup) } +func (self *SHuaWeiRegionDriver) createCachedLbbg(lb *models.SLoadbalancer, lblis *models.SLoadbalancerListener, lbr *models.SLoadbalancerListenerRule, lbbg *models.SLoadbalancerBackendGroup) (*models.SHuaweiCachedLbbg, error) { + // create loadbalancer backendgroup cache + cachedLbbg := &models.SHuaweiCachedLbbg{} + cachedLbbg.ManagerId = lb.ManagerId + cachedLbbg.CloudregionId = lb.CloudregionId + cachedLbbg.LoadbalancerId = lb.GetId() + cachedLbbg.BackendGroupId = lbbg.GetId() + if lbr != nil { + cachedLbbg.AssociatedType = api.LB_ASSOCIATE_TYPE_RULE + cachedLbbg.AssociatedId = lbr.GetId() + cachedLbbg.ProtocolType = lblis.ListenerType + } else { + cachedLbbg.AssociatedType = api.LB_ASSOCIATE_TYPE_LISTENER + cachedLbbg.AssociatedId = lblis.GetId() + cachedLbbg.ProtocolType = lblis.ListenerType + } + + err := models.HuaweiCachedLbbgManager.TableSpec().Insert(cachedLbbg) + if err != nil { + return nil, err + } + + cachedLbbg.SetModelManager(models.HuaweiCachedLbbgManager, cachedLbbg) + return cachedLbbg, nil +} + +func (self *SHuaWeiRegionDriver) syncCloudlbbs(ctx context.Context, userCred mcclient.TokenCredential, lb *models.SLoadbalancer, cachedLbbg *models.SHuaweiCachedLbbg, extlbbg cloudprovider.ICloudLoadbalancerBackendGroup, backends []cloudprovider.SLoadbalancerBackend) error { + ibackends, err := extlbbg.GetILoadbalancerBackends() + if err != nil { + return errors.Wrap(err, "HuaWeiRegionDriver.syncCloudLoadbalancerBackends.GetILoadbalancerBackends") + } + + for i := range ibackends { + ibackend := ibackends[i] + err = extlbbg.RemoveBackendServer(ibackend.GetId(), ibackend.GetWeight(), ibackend.GetPort()) + if err != nil { + return errors.Wrap(err, "HuaWeiRegionDriver.syncCloudLoadbalancerBackends.RemoveBackendServer") + } + } + + for _, backend := range backends { + _, err = extlbbg.AddBackendServer(backend.ExternalID, backend.Weight, backend.Port) + if err != nil { + return errors.Wrap(err, "HuaWeiRegionDriver.syncCloudLoadbalancerBackends.AddBackendServer") + } + } + + return nil +} + +func (self *SHuaWeiRegionDriver) syncCachedLbbs(ctx context.Context, userCred mcclient.TokenCredential, lb *models.SLoadbalancer, lbbg *models.SHuaweiCachedLbbg, extlbbg cloudprovider.ICloudLoadbalancerBackendGroup) error { + iBackends, err := extlbbg.GetILoadbalancerBackends() + if err != nil { + return errors.Wrap(err, "HuaWeiRegionDriver.syncLoadbalancerBackendCaches.GetILoadbalancerBackends") + } + + if len(iBackends) > 0 { + provider := lb.GetCloudprovider() + if provider == nil { + return fmt.Errorf("failed to find cloudprovider for lb %s", lb.Name) + } + + result := models.HuaweiCachedLbManager.SyncLoadbalancerBackends(ctx, userCred, provider, lbbg, iBackends, &models.SSyncRange{}) + if result.IsError() { + return errors.Wrap(result.AllError(), "HuaWeiRegionDriver.syncLoadbalancerBackendCaches.SyncLoadbalancerBackends") + } + } + + return nil +} + func (self *SHuaWeiRegionDriver) createLoadbalancerBackendGroup(ctx context.Context, userCred mcclient.TokenCredential, lblis *models.SLoadbalancerListener, lbr *models.SLoadbalancerListenerRule, lbbg *models.SLoadbalancerBackendGroup, backends []cloudprovider.SLoadbalancerBackend) (jsonutils.JSONObject, error) { if len(lblis.ListenerType) == 0 { return nil, fmt.Errorf("loadbalancer backendgroup missing protocol type") @@ -576,25 +649,9 @@ func (self *SHuaWeiRegionDriver) createLoadbalancerBackendGroup(ctx context.Cont return nil, err } - // create loadbalancer backendgroup cache - cachedLbbg := &models.SHuaweiCachedLbbg{} - cachedLbbg.ManagerId = lb.ManagerId - cachedLbbg.CloudregionId = lb.CloudregionId - cachedLbbg.LoadbalancerId = lb.GetId() - cachedLbbg.BackendGroupId = lbbg.GetId() - if lbr != nil { - cachedLbbg.AssociatedType = api.LB_ASSOCIATE_TYPE_RULE - cachedLbbg.AssociatedId = lbr.GetId() - cachedLbbg.ProtocolType = lblis.ListenerType - } else { - cachedLbbg.AssociatedType = api.LB_ASSOCIATE_TYPE_LISTENER - cachedLbbg.AssociatedId = lblis.GetId() - cachedLbbg.ProtocolType = lblis.ListenerType - } - - err = models.HuaweiCachedLbbgManager.TableSpec().Insert(cachedLbbg) + cachedLbbg, err := self.createCachedLbbg(lb, lblis, lbr, lbbg) if err != nil { - return nil, err + return nil, errors.Wrap(err, "HuaWeiRegionDriver.createLoadbalancerBackendGroupCache") } group, err := lbbg.GetHuaweiBackendGroupParams(lblis, lbr) @@ -607,45 +664,20 @@ func (self *SHuaWeiRegionDriver) createLoadbalancerBackendGroup(ctx context.Cont return nil, err } - cachedLbbg.SetModelManager(models.HuaweiCachedLbbgManager, cachedLbbg) if err := db.SetExternalId(cachedLbbg, userCred, iLoadbalancerBackendGroup.GetGlobalId()); err != nil { return nil, err } - for _, backend := range backends { - cachedlbb := &models.SHuaweiCachedLb{} - cachedlbb.ManagerId = lb.ManagerId - cachedlbb.CloudregionId = lb.CloudregionId - cachedlbb.CachedBackendGroupId = cachedLbbg.GetId() - cachedlbb.BackendId = backend.ID - err = models.HuaweiCachedLbManager.TableSpec().Insert(cachedlbb) - if err != nil { - return nil, err - } - - ibackend, err := iLoadbalancerBackendGroup.AddBackendServer(backend.ExternalID, backend.Weight, backend.Port) - if err != nil { - return nil, err - } - - cachedlbb.SetModelManager(models.HuaweiCachedLbManager, cachedlbb) - err = db.SetExternalId(cachedlbb, userCred, ibackend.GetGlobalId()) - if err != nil { - return nil, err - } - } - - iBackends, err := iLoadbalancerBackendGroup.GetILoadbalancerBackends() + err = self.syncCloudlbbs(ctx, userCred, lb, cachedLbbg, iLoadbalancerBackendGroup, backends) if err != nil { - return nil, err + return nil, errors.Wrap(err, "HuaWeiRegionDriver.createLoadbalancerBackendGroup.syncCloudLoadbalancerBackends") } - if len(iBackends) > 0 { - provider := lb.GetCloudprovider() - if provider == nil { - return nil, fmt.Errorf("failed to find cloudprovider for lb %s", lb.Name) - } - models.HuaweiCachedLbManager.SyncLoadbalancerBackends(ctx, userCred, provider, cachedLbbg, iBackends, &models.SSyncRange{}) + + err = self.syncCachedLbbs(ctx, userCred, lb, cachedLbbg, iLoadbalancerBackendGroup) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.createLoadbalancerBackendGroup.syncLoadbalancerBackendCaches") } + return nil, nil } @@ -700,7 +732,7 @@ func (self *SHuaWeiRegionDriver) createLoadbalancerAcl(ctx context.Context, user ListenerId: listener.GetExternalId(), Name: lbacl.Name, Entrys: []cloudprovider.SLoadbalancerAccessControlListEntry{}, - AccessControlEnable: (listener.AclStatus == api.LB_BOOL_ON), + AccessControlEnable: listener.AclStatus == api.LB_BOOL_ON, } _originAcl, err := db.FetchById(models.LoadbalancerAclManager, lbacl.AclId) @@ -733,6 +765,103 @@ func (self *SHuaWeiRegionDriver) RequestCreateLoadbalancerAcl(ctx context.Contex return nil } +func (self *SHuaWeiRegionDriver) updateCachedLbbg(ctx context.Context, cachedLbbg *models.SHuaweiCachedLbbg, backendGroupId string, externalBackendGroupId string, asscoicateId string, asscoicateType string) error { + _, err := db.UpdateWithLock(ctx, cachedLbbg, func() error { + if len(backendGroupId) > 0 { + cachedLbbg.BackendGroupId = backendGroupId + } + + if len(externalBackendGroupId) > 0 { + cachedLbbg.ExternalId = externalBackendGroupId + } + + if len(asscoicateId) > 0 { + cachedLbbg.AssociatedId = asscoicateId + } + + if len(asscoicateType) > 0 { + cachedLbbg.AssociatedType = asscoicateType + } + + return nil + }) + if err != nil { + return err + } + + return nil +} + +func (self *SHuaWeiRegionDriver) validatorCachedLbbgConflict(ilb cloudprovider.ICloudLoadbalancer, lblis *models.SLoadbalancerListener, olbbg *models.SHuaweiCachedLbbg, nlbbg *models.SHuaweiCachedLbbg, rlbbg *models.SHuaweiCachedLbbg) error { + if rlbbg == nil { + return nil + } + + if len(rlbbg.AssociatedId) > 0 && lblis.GetId() != rlbbg.AssociatedId { + return fmt.Errorf("sync required, backend group aready assoicate with %s %s", rlbbg.AssociatedType, rlbbg.AssociatedId) + } + + if olbbg != nil && len(rlbbg.AssociatedId) > 0 && rlbbg.AssociatedId != olbbg.AssociatedId { + // 与本地关联状态不一致 + return fmt.Errorf("sync required, backend group aready assoicate with %s %s", rlbbg.AssociatedType, rlbbg.AssociatedId) + } + + if nlbbg != nil && len(nlbbg.ExternalId) > 0 { + // 需绑定的服务器组,被人为删除了 + _, err := ilb.GetILoadBalancerBackendGroupById(nlbbg.GetExternalId()) + if err != nil { + return fmt.Errorf("sync required, validatorCachedLbbgConflict.GetILoadBalancerBackendGroupById %s", err) + } + + // 需绑定的服务器组,被人为关联了其他监听或者监听规则 + listeners, err := ilb.GetILoadBalancerListeners() + if err != nil { + return fmt.Errorf("sync required, validatorCachedLbbgConflict.GetILoadBalancerListeners %s", err) + } + + for i := range listeners { + listener := listeners[i] + if bgId := listener.GetBackendGroupId(); len(bgId) > 0 && bgId == nlbbg.ExternalId && listener.GetGlobalId() != lblis.GetExternalId() { + return fmt.Errorf("sync required, backend group aready assoicate with listener %s(external id )", listener.GetGlobalId()) + } + + rules, err := listener.GetILoadbalancerListenerRules() + if err != nil { + return fmt.Errorf("sync required, validatorCachedLbbgConflict.GetILoadbalancerListenerRules %s", err) + } + + for j := range rules { + rule := rules[j] + if bgId := rule.GetBackendGroupId(); len(bgId) > 0 && bgId == nlbbg.ExternalId { + return fmt.Errorf("sync required, backend group aready assoicate with rule %s(external id )", rule.GetGlobalId()) + } + } + } + } + + return nil +} + +func (self *SHuaWeiRegionDriver) removeCachedLbbg(ctx context.Context, userCred mcclient.TokenCredential, lbbg *models.SHuaweiCachedLbbg) error { + backends, err := lbbg.GetCachedBackends() + if err != nil && err != sql.ErrNoRows { + return errors.Wrap(err, "huaweiRegionDriver.GetCachedBackends") + } + for i := range backends { + err = db.DeleteModel(ctx, userCred, &backends[i]) + if err != nil { + return errors.Wrap(err, "huaweiRegionDriver.DeleteModel") + } + } + + err = db.DeleteModel(ctx, userCred, lbbg) + if err != nil { + return errors.Wrap(err, "huaweiRegionDriver.DeleteModel") + } + + return nil +} + func (self *SHuaWeiRegionDriver) RequestSyncLoadbalancerBackendGroup(ctx context.Context, userCred mcclient.TokenCredential, lblis *models.SLoadbalancerListener, lbbg *models.SLoadbalancerBackendGroup, task taskman.ITask) error { taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { iRegion, err := lbbg.GetIRegion() @@ -746,72 +875,256 @@ func (self *SHuaWeiRegionDriver) RequestSyncLoadbalancerBackendGroup(ctx context return nil, err } + provider := lb.GetCloudprovider() + if provider == nil { + return nil, errors.Wrap(fmt.Errorf("provider is nil"), "HuaWeiRegionDriver.RequestSyncLoadbalancerBackendGroup.GetCloudprovider") + } + ilisten, err := ilb.GetILoadBalancerListenerById(lblis.GetExternalId()) if err != nil { return nil, err } - olbbg := ilisten.GetBackendGroupId() - if len(olbbg) > 0 { - p, err := lblis.GetHuaweiLoadbalancerListenerParams() - if err != nil { - return nil, err - } - - err = ilisten.Sync(p) - if err != nil { - return nil, err - } - - _t, err := db.FetchByExternalId(models.HuaweiCachedLbbgManager, olbbg) - if err != nil { - return nil, err - } - - tmpLbbg := _t.(*models.SHuaweiCachedLbbg) - _, err = db.UpdateWithLock(ctx, tmpLbbg, func() error { - tmpLbbg.AssociatedId = "" - tmpLbbg.AssociatedType = "" - return nil - }) - - if err != nil { - return nil, err - } + // new group params + groupInput, err := lbbg.GetHuaweiBackendGroupParams(lblis, nil) + if err != nil { + return nil, errors.Wrap(err, "GetHuaweiBackendGroupParams") } - backends, err := lbbg.GetBackendsParams() + // new backends params + backendsInput, err := lbbg.GetBackendsParams() if err != nil { return nil, err } - nlbbg, err := models.HuaweiCachedLbbgManager.GetUsableCachedBackendGroup(lbbg.GetId(), lblis.ListenerType) - if nlbbg == nil { - if _, err := self.createLoadbalancerBackendGroup(ctx, task.GetUserCred(), lblis, nil, lbbg, backends); err != nil { - return nil, err - } - } else { - ilbbg, err := ilb.GetILoadBalancerBackendGroupById(nlbbg.GetExternalId()) + { + var olbbg, nlbbg, rlbbg *models.SHuaweiCachedLbbg + // current related cachedLbbg + olbbg, err := models.HuaweiCachedLbbgManager.GetCachedBackendGroupByAssociateId(lblis.GetId()) if err != nil { - return nil, err + if err != sql.ErrNoRows { + return nil, errors.Wrap(fmt.Errorf("provider is nil"), "HuaWeiRegionDriver.RequestSyncLoadbalancerBackendGroup.GetCachedBackendGroupByAssociateId") + } } - group, err := lbbg.GetHuaweiBackendGroupParams(lblis, nil) + // new related backendgroup + if olbbg == nil || olbbg.BackendGroupId != lbbg.GetId() { + nlbbg, err = models.HuaweiCachedLbbgManager.GetUsableCachedBackendGroup(lbbg.GetId(), lblis.ListenerType) + if err != nil { + if err != sql.ErrNoRows { + return nil, errors.Wrap(fmt.Errorf("provider is nil"), "HuaWeiRegionDriver.RequestSyncLoadbalancerBackendGroup.GetUsableCachedBackendGroup") + } + } + } + + // remote relasted backendgroup + rlbbgId := ilisten.GetBackendGroupId() + if len(rlbbgId) > 0 { + _rlbbg, err := db.FetchByExternalId(models.HuaweiCachedLbbgManager, rlbbgId) + if err != nil { + if err != sql.ErrNoRows { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.RequestSyncLoadbalancerBackendGroup.FetchByExternalId") + } + } + + rlbbg = _rlbbg.(*models.SHuaweiCachedLbbg) + } + + // validator confilct + err = self.validatorCachedLbbgConflict(ilb, lblis, olbbg, nlbbg, rlbbg) if err != nil { - return nil, err + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.validatorCachedLbbgConflict") } - if err := ilbbg.Sync(group); err != nil { - return nil, err + // case 1:新绑定,创建本地缓存,并创建远端服务器组 + if olbbg == nil && nlbbg == nil && rlbbg == nil { + if _, err := self.createLoadbalancerBackendGroup(ctx, task.GetUserCred(), lblis, nil, lbbg, backendsInput); err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case1.createLoadbalancerBackendGroup") + } } - nlbbg.SetModelManager(models.HuaweiCachedLbbgManager, nlbbg) - _, err = db.UpdateWithLock(ctx, nlbbg, func() error { - nlbbg.AssociatedId = lblis.GetId() - nlbbg.AssociatedType = api.LB_ASSOCIATE_TYPE_LISTENER - return nil - }) + // case 2:新绑定,创建本地缓存 + if olbbg == nil && nlbbg == nil && rlbbg != nil { + cachedLbbg, err := self.createCachedLbbg(lb, lblis, nil, lbbg) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case2.createCachedLbbg") + } + + err = db.SetExternalId(cachedLbbg, userCred, rlbbgId) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case2.SetExternalId") + } + + ilbbg, err := ilb.GetILoadBalancerBackendGroupById(rlbbgId) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case2.GetILoadBalancerBackendGroupById") + } + + err = self.syncCloudlbbs(ctx, userCred, lb, cachedLbbg, ilbbg, backendsInput) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case2.syncCloudlbbs") + } + + err = self.syncCachedLbbs(ctx, userCred, lb, cachedLbbg, ilbbg) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case2.syncCloudlbbs") + } + } + + // case 3:新绑定, 关联并同步 + if olbbg == nil && nlbbg != nil && rlbbg == nil { + ilbbg, err := ilb.GetILoadBalancerBackendGroupById(nlbbg.GetExternalId()) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case4.GetILoadBalancerBackendGroupById") + } + + err = db.SetExternalId(nlbbg, userCred, ilbbg.GetGlobalId()) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case3.SetExternalId") + } + + err = self.syncCloudlbbs(ctx, userCred, lb, nlbbg, ilbbg, backendsInput) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case3.syncCloudlbbs") + } + + err = self.syncCachedLbbs(ctx, userCred, lb, nlbbg, ilbbg) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case3.syncCloudlbbs") + } + } + + // case 4:新绑定, 关联rlbbg并同步 + if olbbg == nil && nlbbg != nil && rlbbg != nil { + ilbbg, err := ilb.GetILoadBalancerBackendGroupById(rlbbgId) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case4.GetILoadBalancerBackendGroupById") + } + + err = db.SetExternalId(nlbbg, userCred, ilbbg.GetGlobalId()) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case4.SetExternalId") + } + + err = self.syncCloudlbbs(ctx, userCred, lb, nlbbg, ilbbg, backendsInput) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case4.syncCloudlbbs") + } + + err = self.syncCachedLbbs(ctx, userCred, lb, nlbbg, ilbbg) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case4.syncCloudlbbs") + } + } + + // case 5:已绑定,更新原有缓存 + if olbbg != nil && nlbbg == nil && rlbbg == nil { + ilbbg, err := ilb.CreateILoadBalancerBackendGroup(groupInput) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case5.CreateILoadBalancerBackendGroup") + } + + err = self.updateCachedLbbg(ctx, olbbg, lbbg.GetId(), ilbbg.GetGlobalId(), "", "") + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case5.updateCachedLbbg") + } + + err = self.syncCloudlbbs(ctx, userCred, lb, olbbg, ilbbg, backendsInput) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case4.syncCloudlbbs") + } + + err = self.syncCachedLbbs(ctx, userCred, lb, olbbg, ilbbg) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case4.syncCloudlbbs") + } + } + + // case 6:已绑定, 更新olbbg并同步 + if olbbg != nil && nlbbg == nil && rlbbg != nil { + ilbbg, err := ilb.GetILoadBalancerBackendGroupById(rlbbgId) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case6.GetILoadBalancerBackendGroupById") + } + + err = self.updateCachedLbbg(ctx, olbbg, lbbg.GetId(), ilbbg.GetGlobalId(), "", "") + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case6.updateCachedLbbg") + } + + err = self.syncCloudlbbs(ctx, userCred, lb, olbbg, ilbbg, backendsInput) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case6.syncCloudlbbs") + } + + err = self.syncCachedLbbs(ctx, userCred, lb, olbbg, ilbbg) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case6.syncCloudlbbs") + } + } + + // case 7:已绑定, 更新olbbg并同步 + if olbbg != nil && nlbbg != nil && rlbbg == nil { + err := self.removeCachedLbbg(ctx, userCred, olbbg) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case7.removeCachedLbbg") + } + + err = self.updateCachedLbbg(ctx, nlbbg, lbbg.GetId(), "", lblis.GetId(), api.LB_ASSOCIATE_TYPE_LISTENER) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case7.updateCachedLbbg") + } + + ilbbg, err := ilb.GetILoadBalancerBackendGroupById(nlbbg.GetExternalId()) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case7.GetILoadBalancerBackendGroupById") + } + + err = self.syncCloudlbbs(ctx, userCred, lb, nlbbg, ilbbg, backendsInput) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case6.syncCloudlbbs") + } + + err = self.syncCachedLbbs(ctx, userCred, lb, nlbbg, ilbbg) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case6.syncCloudlbbs") + } + } + + // case 8:已绑定, 更新olbbg并同步 + if olbbg != nil && nlbbg != nil && rlbbg != nil { + ilbbg, err := ilb.GetILoadBalancerBackendGroupById(rlbbgId) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case8.GetILoadBalancerBackendGroupById") + } + + err = deleteHuaweiLoadbalancerBackendGroup(ctx, userCred, ilb, ilbbg) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case8.deleteHuaweiLoadbalancerBackendGroup") + } + + err = deleteHuaweiCachedLbbg(ctx, userCred, lblis.GetId()) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case8.deleteHuaweiCachedLbbg") + } + + err = self.updateCachedLbbg(ctx, nlbbg, lbbg.GetId(), "", lblis.GetId(), api.LB_ASSOCIATE_TYPE_LISTENER) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case8.updateCachedLbbg") + } + + err = self.syncCloudlbbs(ctx, userCred, lb, nlbbg, ilbbg, backendsInput) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case6.syncCloudlbbs") + } + + err = self.syncCachedLbbs(ctx, userCred, lb, nlbbg, ilbbg) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Sync.Case6.syncCloudlbbs") + } + } } + // continue here cachedLbbg, err := models.HuaweiCachedLbbgManager.GetCachedBackendGroupByAssociateId(lblis.GetId()) if err != nil { @@ -1093,6 +1406,169 @@ func (self *SHuaWeiRegionDriver) RequestSyncLoadbalancerListener(ctx context.Con return nil } +func deleteHuaweiLoadbalancerListenerRule(ctx context.Context, userCred mcclient.TokenCredential, ilb cloudprovider.ICloudLoadbalancer, irule cloudprovider.ICloudLoadbalancerListenerRule) error { + err := irule.Refresh() + if err != nil { + if err == cloudprovider.ErrNotFound { + return nil + } + + return errors.Wrap(err, "HuaWeiRegionDriver.Rule.Refresh") + } + + lbbgId := irule.GetBackendGroupId() + + err = irule.Delete() + if err != nil && err != cloudprovider.ErrNotFound { + return errors.Wrap(err, "HuaWeiRegionDriver.Rule.Delete") + } + + // delete backendgroup + if len(lbbgId) > 0 { + ilbbg, err := ilb.GetILoadBalancerBackendGroupById(lbbgId) + if err != nil { + if err == cloudprovider.ErrNotFound { + return nil + } + + return errors.Wrap(err, "HuaWeiRegionDriver.Rule.GetILoadBalancerBackendGroupById") + } + + err = deleteHuaweiLoadbalancerBackendGroup(ctx, userCred, ilb, ilbbg) + if err != nil { + return errors.Wrap(err, "HuaWeiRegionDriver.Rule.deleteHuaweiLoadbalancerBackendGroup") + } + } + + return nil +} + +func deleteHuaweiLoadbalancerBackendGroup(ctx context.Context, userCred mcclient.TokenCredential, ilb cloudprovider.ICloudLoadbalancer, ilbbg cloudprovider.ICloudLoadbalancerBackendGroup) error { + err := ilbbg.Refresh() + if err != nil { + if err == cloudprovider.ErrNotFound { + return nil + } + + return errors.Wrap(err, "HuaWeiRegionDriver.BackendGroup.Refresh") + } + + ibackends, err := ilbbg.GetILoadbalancerBackends() + if err != nil { + return errors.Wrap(err, "HuaWeiRegionDriver.BackendGroup.GetILoadbalancerBackends") + } + + for i := range ibackends { + ilbb := ibackends[i] + err = deleteHuaweiLoadbalancerBackend(ctx, userCred, ilb, ilbbg, ilbb) + if err != nil { + return errors.Wrap(err, "HuaWeiRegionDriver.BackendGroup.deleteHuaweiLoadbalancerBackend") + } + } + + err = ilbbg.Delete() + if err != nil && err != cloudprovider.ErrNotFound { + return errors.Wrap(err, "HuaWeiRegionDriver.BackendGroup.Delete") + } + + return nil +} + +func deleteHuaweiLoadbalancerBackend(ctx context.Context, userCred mcclient.TokenCredential, ilb cloudprovider.ICloudLoadbalancer, ilbbg cloudprovider.ICloudLoadbalancerBackendGroup, ilbb cloudprovider.ICloudLoadbalancerBackend) error { + err := ilbbg.RemoveBackendServer(ilbb.GetId(), ilbb.GetWeight(), ilbb.GetPort()) + if err != nil && err != cloudprovider.ErrNotFound { + return errors.Wrap(err, "HuaWeiRegionDriver.Backend.Delete") + } + + return nil +} + +func deleteHuaweiLblisRule(ctx context.Context, userCred mcclient.TokenCredential, ruleId string) error { + rule, err := db.FetchById(models.LoadbalancerListenerRuleManager, ruleId) + if err != nil { + if err == sql.ErrNoRows { + return nil + } + + return errors.Wrap(err, "deleteHuaweiLblisRule.FetchById") + } + + err = deleteHuaweiCachedLbbg(ctx, userCred, ruleId) + if err != nil { + return errors.Wrap(err, "deleteHuaweiLblisRule.deleteHuaweiCachedLbbg") + } + + err = db.DeleteModel(ctx, userCred, rule) + if err != nil { + return errors.Wrap(err, "deleteHuaweiLblisRule.DeleteModel") + } + + return nil +} + +func deleteHuaweiCachedLbbg(ctx context.Context, userCred mcclient.TokenCredential, associatedId string) error { + lbbg, err := models.HuaweiCachedLbbgManager.GetCachedBackendGroupByAssociateId(associatedId) + if err != nil { + if err == sql.ErrNoRows { + return nil + } + + return errors.Wrap(err, "deleteHuaweiCachedLbbg.GetCachedBackendGroupByAssociateId") + } + + if err := deleteHuaweiCachedLbbsByLbbg(ctx, userCred, lbbg.GetId()); err != nil { + return errors.Wrap(err, "deleteHuaweiCachedLbbg.deleteHuaweiCachedLbbsByLbbg") + } + + err = db.DeleteModel(ctx, userCred, lbbg) + if err != nil { + return errors.Wrap(err, "deleteHuaweiCachedLbbg.DeleteModel") + } + + return nil +} + +func deleteHuaweiCachedLbbsByLbbg(ctx context.Context, userCred mcclient.TokenCredential, cachedLbbgId string) error { + cachedLbbs := []models.SHuaweiCachedLb{} + q := models.HuaweiCachedLbManager.Query().IsFalse("pending_deleted").Equals("cached_backend_group_id", cachedLbbgId) + err := db.FetchModelObjects(models.HuaweiCachedLbManager, q, &cachedLbbs) + if err != nil { + if err == sql.ErrNoRows { + return nil + } + + return errors.Wrap(err, "deleteHuaweiCachedLbbsByLbbg.FetchModelObjects") + } + + for i := range cachedLbbs { + cachedLbb := cachedLbbs[i] + err = db.DeleteModel(ctx, userCred, &cachedLbb) + if err != nil { + return errors.Wrap(err, "deleteHuaweiCachedLbbsByLbbg.DeleteModel") + } + } + + return nil +} + +func deleteHuaweiCachedLbb(ctx context.Context, userCred mcclient.TokenCredential, cachedLbbId string) error { + lbb, err := db.FetchById(models.HuaweiCachedLbManager, cachedLbbId) + if err != nil { + if err == sql.ErrNoRows { + return nil + } + + return errors.Wrap(err, "deleteHuaweiCachedLbb.FetchById") + } + + err = db.DeleteModel(ctx, userCred, lbb) + if err != nil { + return errors.Wrap(err, "deleteHuaweiCachedLbb.DeleteModel") + } + + return nil +} + func (self *SHuaWeiRegionDriver) RequestDeleteLoadbalancerListener(ctx context.Context, userCred mcclient.TokenCredential, lblis *models.SLoadbalancerListener, task taskman.ITask) error { taskman.LocalTaskRun(task, func() (jsonutils.JSONObject, error) { if jsonutils.QueryBoolean(task.GetParams(), "purge", false) { @@ -1143,19 +1619,25 @@ func (self *SHuaWeiRegionDriver) RequestDeleteLoadbalancerListener(ctx context.C err = iListener.Sync(params) if err != nil { return nil, err + } else { + iListener.Refresh() } - _cachedLbbg, err := db.FetchByExternalId(models.HuaweiCachedLbbgManager, backendgroupId) + // 删除后端服务器组 + ilbbg, err := iLoadbalancer.GetILoadBalancerBackendGroupById(backendgroupId) if err != nil { - return nil, err + return nil, errors.Wrap(err, "HuaWeiRegionDriver.RequestDeleteLoadbalancerListener.GetILoadBalancerBackendGroup") } - cachedLbbg := _cachedLbbg.(*models.SHuaweiCachedLbbg) - _, err = db.UpdateWithLock(ctx, cachedLbbg, func() error { - cachedLbbg.AssociatedId = "" - cachedLbbg.AssociatedType = "" - return nil - }) + err = deleteHuaweiLoadbalancerBackendGroup(ctx, userCred, iLoadbalancer, ilbbg) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.RequestDeleteLoadbalancerListener.DeleteBackendGroup") + } + } + + err = deleteHuaweiCachedLbbg(ctx, userCred, lblis.GetId()) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.RequestDeleteLoadbalancerListener.deleteHuaweiCachedLbbg") } // 删除访问控制 @@ -1180,6 +1662,32 @@ func (self *SHuaWeiRegionDriver) RequestDeleteLoadbalancerListener(ctx context.C } } + // remove rules + irules, err := iListener.GetILoadbalancerListenerRules() + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.GetILoadbalancerListenerRules") + } + + for i := range irules { + irule := irules[i] + err = deleteHuaweiLoadbalancerListenerRule(ctx, userCred, iLoadbalancer, irule) + if err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.deleteHuaweiLoadbalancerListenerRule") + } + } + + rules, err := lblis.GetLoadbalancerListenerRules() + if err != nil && err != sql.ErrNoRows { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.GetLoadbalancerListenerRules") + } + + for i := range rules { + rule := rules[i] + if err := deleteHuaweiLblisRule(ctx, userCred, rule.GetId()); err != nil { + return nil, errors.Wrap(err, "HuaWeiRegionDriver.Rule.deleteHuaweiCachedLbbg") + } + } + return nil, iListener.Delete() }) return nil @@ -1220,6 +1728,10 @@ func (self *SHuaWeiRegionDriver) RequestDeleteLoadbalancerBackendGroup(ctx conte iLoadbalancerBackendGroup, err := iLoadbalancer.GetILoadBalancerBackendGroupById(cachedLbbg.ExternalId) if err != nil { if err == cloudprovider.ErrNotFound { + if err := deleteHuaweiCachedLbbg(ctx, userCred, cachedLbbg.AssociatedId); err != nil { + return nil, errors.Wrap(err, "huaweiRegionDriver.RequestDeleteLoadbalancerBackendGroup.deleteHuaweiCachedLbbg") + } + continue } return nil, errors.Wrap(err, "huaweiRegionDriver.RequestDeleteLoadbalancerBackendGroup.GetILoadBalancerBackendGroupById") @@ -1279,7 +1791,8 @@ func (self *SHuaWeiRegionDriver) RequestDeleteLoadbalancerBackend(ctx context.Co for _, cachedlbb := range cachedlbbs { cachedlbbg, _ := cachedlbb.GetCachedBackendGroup() if cachedlbbg == nil { - return nil, fmt.Errorf("failed to find lbbg for backend %s", cachedlbb.Name) + log.Warningf("failed to find lbbg for backend %s", cachedlbb.Name) + continue } lb := cachedlbbg.GetLoadbalancer() if lb == nil { @@ -1480,6 +1993,14 @@ func (self *SHuaWeiRegionDriver) RequestCreateLoadbalancerBackend(ctx context.Co for _, cachedLbbg := range cachedlbbgs { iLoadbalancerBackendGroup, err := cachedLbbg.GetICloudLoadbalancerBackendGroup() if err != nil { + if err == cloudprovider.ErrNotFound { + if err := deleteHuaweiCachedLbbg(ctx, userCred, cachedLbbg.AssociatedId); err != nil { + return nil, errors.Wrap(err, "huaweiRegionDriver.RequestCreateLoadbalancerBackend.deleteHuaweiCachedLbbg") + } + + continue + } + return nil, errors.Wrap(err, "huaweiRegionDriver.RequestCreateLoadbalancerBackend.GetICloudLoadbalancerBackendGroup") } diff --git a/pkg/compute/regiondrivers/kvm.go b/pkg/compute/regiondrivers/kvm.go index fa1055cd47..28b7dec0fd 100644 --- a/pkg/compute/regiondrivers/kvm.go +++ b/pkg/compute/regiondrivers/kvm.go @@ -591,6 +591,11 @@ func (self *SKVMRegionDriver) RequestDeleteLoadbalancerAcl(ctx context.Context, return nil } +func (self *SKVMRegionDriver) RequestSyncLoadbalancerBackendGroup(ctx context.Context, userCred mcclient.TokenCredential, lblis *models.SLoadbalancerListener, lbbg *models.SLoadbalancerBackendGroup, task taskman.ITask) error { + task.ScheduleRun(nil) + return nil +} + func (self *SKVMRegionDriver) RequestCreateLoadbalancerCertificate(ctx context.Context, userCred mcclient.TokenCredential, lbcert *models.SCachedLoadbalancerCertificate, task taskman.ITask) error { task.ScheduleRun(nil) return nil diff --git a/pkg/compute/regiondrivers/managedvirtual.go b/pkg/compute/regiondrivers/managedvirtual.go index 17f6f84d9e..2e3b04e405 100644 --- a/pkg/compute/regiondrivers/managedvirtual.go +++ b/pkg/compute/regiondrivers/managedvirtual.go @@ -418,6 +418,18 @@ func (self *SManagedVirtualizationRegionDriver) createLoadbalancerCertificate(ct if err := db.SetExternalId(lbcert, userCred, iLoadbalancerCert.GetGlobalId()); err != nil { return nil, errors.Wrap(err, "db.SetExternalId") } + + err = cloudprovider.WaitCreated(3*time.Second, 30*time.Second, func() bool { + err := iLoadbalancerCert.Refresh() + if err == nil { + return true + } + return false + }) + if err != nil { + return nil, errors.Wrap(err, "iRegion.createLoadbalancerCertificate.WaitCreated") + } + return nil, lbcert.SyncWithCloudLoadbalancerCertificate(ctx, userCred, iLoadbalancerCert, lbcert.GetOwnerId()) } @@ -529,18 +541,6 @@ func (self *SManagedVirtualizationRegionDriver) RequestDeleteLoadbalancerBackend return nil, err } - cachedLbbgs, err := models.AwsCachedLbbgManager.GetCachedBackendGroups(lbbg.GetId()) - if err != nil { - return nil, err - } - - for i := range cachedLbbgs { - err = cachedLbbgs[i].Delete(ctx, userCred) - if err != nil { - return nil, err - } - } - return nil, nil }) return nil diff --git a/pkg/compute/tasks/loadbalancer_listener_create_task.go b/pkg/compute/tasks/loadbalancer_listener_create_task.go index 81beb673b0..2c4fc2226b 100644 --- a/pkg/compute/tasks/loadbalancer_listener_create_task.go +++ b/pkg/compute/tasks/loadbalancer_listener_create_task.go @@ -16,6 +16,7 @@ package tasks import ( "context" + "database/sql" "fmt" "yunion.io/x/jsonutils" @@ -65,7 +66,12 @@ func onHuaweiLoadbalancerListenerCreateComplete(ctx context.Context, lblis *mode params := jsonutils.NewDict() params.Set("listenerId", jsonutils.NewString(lblis.GetId())) - group, _ := models.HuaweiCachedLbbgManager.GetUsableCachedBackendGroup(lbbg.GetId(), lblis.ListenerType) + group, err := models.HuaweiCachedLbbgManager.GetCachedBackendGroupByAssociateId(lblis.GetId()) + if err != nil && err != sql.ErrNoRows { + self.taskFail(ctx, lblis, err.Error()) + return + } + if group != nil { // 服务器组存在 ilbbg, err := group.GetICloudLoadbalancerBackendGroup() diff --git a/pkg/compute/tasks/loadbalancer_listener_rule_create_task.go b/pkg/compute/tasks/loadbalancer_listener_rule_create_task.go index 26e098a4e9..54e3bb5e02 100644 --- a/pkg/compute/tasks/loadbalancer_listener_rule_create_task.go +++ b/pkg/compute/tasks/loadbalancer_listener_rule_create_task.go @@ -67,7 +67,7 @@ func onHuaiweiPrepareLoadbalancerBackendgroup(ctx context.Context, region *model params := jsonutils.NewDict() params.Set("ruleId", jsonutils.NewString(lbr.GetId())) - group, _ := models.HuaweiCachedLbbgManager.GetUsableCachedBackendGroup(lbbg.GetId(), lblis.ListenerType) + group, _ := models.HuaweiCachedLbbgManager.GetCachedBackendGroupByAssociateId(lbr.GetId()) if group != nil { ilbbg, err := group.GetICloudLoadbalancerBackendGroup() if err != nil { @@ -124,7 +124,7 @@ func onAwsPrepareLoadbalancerBackendgroup(ctx context.Context, region *models.SC return } - group, _ := models.AwsCachedLbbgManager.GetUsableCachedBackendGroup(lblis.LoadbalancerId, lblis.BackendGroupId, lblis.ListenerType, lblis.HealthCheckType, lblis.HealthCheckInterval) + group, _ := models.AwsCachedLbbgManager.GetUsableCachedBackendGroup(lblis.LoadbalancerId, lbr.BackendGroupId, lblis.ListenerType, lblis.HealthCheckType, lblis.HealthCheckInterval) if group != nil { ilbbg, err := group.GetICloudLoadbalancerBackendGroup() if err != nil { diff --git a/pkg/compute/tasks/loadbalancer_listener_sync_task.go b/pkg/compute/tasks/loadbalancer_listener_sync_task.go index 95aacab7c3..14682071e2 100644 --- a/pkg/compute/tasks/loadbalancer_listener_sync_task.go +++ b/pkg/compute/tasks/loadbalancer_listener_sync_task.go @@ -52,12 +52,6 @@ func (self *LoadbalancerListenerSyncTask) OnInit(ctx context.Context, obj db.ISt return } - // // todo: 这个if应该可以删除 - // if lblis.GetProviderName() != api.CLOUD_PROVIDER_HUAWEI || lblis.GetProviderName() != api.CLOUD_PROVIDER_AWS { - // self.OnLoadbalancerBackendgroupSyncComplete(ctx, lblis, data) - // return - // } - lbbg := lblis.GetLoadbalancerBackendGroup() if lbbg == nil { self.taskFail(ctx, lblis, fmt.Sprintf("failed to find lbbg for lblis %s", lblis.Name)) @@ -65,7 +59,10 @@ func (self *LoadbalancerListenerSyncTask) OnInit(ctx context.Context, obj db.ISt } self.SetStage("OnLoadbalancerBackendgroupSyncComplete", nil) - if err := region.GetDriver().RequestSyncLoadbalancerBackendGroup(ctx, self.GetUserCred(), lblis, lbbg, self); err != nil { + driver := region.GetDriver() + userCred := self.GetUserCred() + err := driver.RequestSyncLoadbalancerBackendGroup(ctx, userCred, lblis, lbbg, self) + if err != nil { self.taskFail(ctx, lblis, err.Error()) } } @@ -82,9 +79,9 @@ func (self *LoadbalancerListenerSyncTask) OnLoadbalancerBackendgroupSyncComplete } } -func (self *LoadbalancerListenerSyncTask) OnLoadbalancerBackendgroupSyncCompleteFail(ctx context.Context, lblis *models.SLoadbalancerListener, reason jsonutils.JSONObject) { +func (self *LoadbalancerListenerSyncTask) OnLoadbalancerBackendgroupSyncCompleteFailed(ctx context.Context, lblis *models.SLoadbalancerListener, reason jsonutils.JSONObject) { lblis.SetStatus(self.GetUserCred(), api.LB_SYNC_CONF_FAILED, reason.String()) - self.SetStageFailed(ctx, reason.String()) + self.taskFail(ctx, lblis, reason.String()) } func (self *LoadbalancerListenerSyncTask) OnLoadbalancerListenerSyncComplete(ctx context.Context, lblis *models.SLoadbalancerListener, data jsonutils.JSONObject) { diff --git a/pkg/multicloud/aliyun/loadbalancerbackendgroup.go b/pkg/multicloud/aliyun/loadbalancerbackendgroup.go index bdf23cf295..2f74917c9f 100644 --- a/pkg/multicloud/aliyun/loadbalancerbackendgroup.go +++ b/pkg/multicloud/aliyun/loadbalancerbackendgroup.go @@ -56,6 +56,10 @@ type SLoadbalancerBackendGroup struct { AssociatedObjects AssociatedObjects } +func (backendgroup *SLoadbalancerBackendGroup) GetILoadbalancer() cloudprovider.ICloudLoadbalancer { + return backendgroup.lb +} + func (backendgroup *SLoadbalancerBackendGroup) GetLoadbalancerId() string { return backendgroup.lb.GetId() } diff --git a/pkg/multicloud/aliyun/loadbalancerdefaultbackendgroup.go b/pkg/multicloud/aliyun/loadbalancerdefaultbackendgroup.go index 2b7147c266..dc1a99dce6 100644 --- a/pkg/multicloud/aliyun/loadbalancerdefaultbackendgroup.go +++ b/pkg/multicloud/aliyun/loadbalancerdefaultbackendgroup.go @@ -27,6 +27,10 @@ type SLoadbalancerDefaultBackendGroup struct { lb *SLoadbalancer } +func (backendgroup *SLoadbalancerDefaultBackendGroup) GetILoadbalancer() cloudprovider.ICloudLoadbalancer { + return backendgroup.lb +} + func (backendgroup *SLoadbalancerDefaultBackendGroup) GetLoadbalancerId() string { return backendgroup.lb.GetId() } diff --git a/pkg/multicloud/aliyun/loadbalancerhttpslistener.go b/pkg/multicloud/aliyun/loadbalancerhttpslistener.go index 5a129f6129..4e3c687ab5 100644 --- a/pkg/multicloud/aliyun/loadbalancerhttpslistener.go +++ b/pkg/multicloud/aliyun/loadbalancerhttpslistener.go @@ -281,6 +281,9 @@ func (region *SRegion) constructHTTPCreateListenerParams(params map[string]strin } params["HealthCheckTimeout"] = fmt.Sprintf("%d", listener.HealthCheckTimeout) } + params["RequestTimeout"] = fmt.Sprintf("%d", listener.ClientRequestTimeout) + params["IdleTimeout"] = fmt.Sprintf("%d", listener.ClientIdleTimeout) + params["StickySession"] = listener.StickySession params["StickySessionType"] = listener.StickySessionType params["Cookie"] = listener.StickySessionCookie @@ -304,6 +307,12 @@ func (region *SRegion) CreateLoadbalancerHTTPSListener(lb *SLoadbalancer, listen params := region.constructBaseCreateListenerParams(lb, listener) params = region.constructHTTPCreateListenerParams(params, listener) params["ServerCertificateId"] = listener.CertificateID + if listener.EnableHTTP2 { + params["EnableHttp2"] = "on" + } else { + params["EnableHttp2"] = "off" + } + if len(listener.TLSCipherPolicy) > 0 { params["TLSCipherPolicy"] = listener.TLSCipherPolicy } @@ -361,6 +370,12 @@ func (region *SRegion) SyncLoadbalancerHTTPSListener(lb *SLoadbalancer, listener params := region.constructBaseCreateListenerParams(lb, listener) params = region.constructHTTPCreateListenerParams(params, listener) params["ServerCertificateId"] = listener.CertificateID + if listener.EnableHTTP2 { + params["EnableHttp2"] = "on" + } else { + params["EnableHttp2"] = "off" + } + if len(lb.LoadBalancerSpec) > 0 && len(listener.TLSCipherPolicy) > 0 { params["TLSCipherPolicy"] = listener.TLSCipherPolicy } diff --git a/pkg/multicloud/aws/aws.go b/pkg/multicloud/aws/aws.go index a08452ff8f..5473c0570a 100644 --- a/pkg/multicloud/aws/aws.go +++ b/pkg/multicloud/aws/aws.go @@ -224,6 +224,12 @@ func (client *SAwsClient) fetchBuckets() error { log.Errorf("s3cli.GetBucketLocation error %s", err) continue } + + if output == nil { + log.Errorf("s3cli.GetBucketLocation nil output") + continue + } + location := *output.LocationConstraint region, err := client.getIRegionByRegionId(location) if err != nil { diff --git a/pkg/multicloud/aws/image.go b/pkg/multicloud/aws/image.go index 12d77b3e62..846c98336f 100644 --- a/pkg/multicloud/aws/image.go +++ b/pkg/multicloud/aws/image.go @@ -35,6 +35,10 @@ const ( ImageStatusCreating ImageStatusType = "pending" ImageStatusAvailable ImageStatusType = "available" ImageStatusCreateFailed ImageStatusType = "failed" + + ImageImportStatusCompleted = "completed" + ImageImportStatusUncompleted = "uncompleted" + ImageImportStatusError = "error" ) type TImageOwnerType string @@ -53,6 +57,8 @@ var ( ) type ImageImportTask struct { + region *SRegion + ImageId string RegionId string TaskId string @@ -95,6 +101,56 @@ type SImage struct { OSBuildId string } +func (self *ImageImportTask) GetId() string { + return self.TaskId +} + +func (self *ImageImportTask) GetName() string { + return self.GetId() +} + +func (self *ImageImportTask) GetGlobalId() string { + return self.GetId() +} + +func (self *ImageImportTask) Refresh() error { + return nil +} + +func (self *ImageImportTask) IsEmulated() bool { + return true +} + +func (self *ImageImportTask) GetMetadata() *jsonutils.JSONDict { + return nil +} + +func (self *ImageImportTask) GetStatus() string { + ret, err := self.region.ec2Client.DescribeImportImageTasks(&ec2.DescribeImportImageTasksInput{ImportTaskIds: []*string{&self.TaskId}}) + if err != nil { + log.Errorf("DescribeImportImageTasks %s", err) + return ImageImportStatusError + } + + err = FillZero(ret) + if err != nil { + log.Errorf("DescribeImportImageTask.FillZero %s", err) + return ImageImportStatusError + } + + // 打印上传进度 + log.Debugf("DescribeImportImage Task %s", ret.String()) + for _, item := range ret.ImportImageTasks { + if *item.Status == "completed" { + return ImageImportStatusCompleted + } else { + return ImageImportStatusUncompleted + } + } + + return ImageImportStatusUncompleted +} + func (self *SImage) GetMinRamSizeMb() int { return 0 } diff --git a/pkg/multicloud/aws/loadbalancerbackendgroup.go b/pkg/multicloud/aws/loadbalancerbackendgroup.go index a844d98460..c5d9e01a04 100644 --- a/pkg/multicloud/aws/loadbalancerbackendgroup.go +++ b/pkg/multicloud/aws/loadbalancerbackendgroup.go @@ -57,6 +57,10 @@ func (self *SElbBackendGroup) GetLoadbalancerId() string { return "" } +func (self *SElbBackendGroup) GetILoadbalancer() cloudprovider.ICloudLoadbalancer { + return self.lb +} + type Matcher struct { HTTPCode string `json:"HttpCode"` } @@ -246,7 +250,7 @@ func (self *SRegion) GetELbBackends(backendgroupId string) ([]SElbBackend, error ret := []SElbBackend{} for i := range backends { - if !utils.IsInStringArray(backends[i].TargetHealth.Reason, []string{"Target.InvalidState", "Target.DeregistrationInProgress"}) { + if !utils.IsInStringArray(backends[i].TargetHealth.Reason, []string{"Target.InvalidState", "Target.NotInUse", "Target.DeregistrationInProgress"}) { backends[i].region = self backends[i].group = group ret = append(ret, backends[i]) diff --git a/pkg/multicloud/aws/loadbalancerlistener.go b/pkg/multicloud/aws/loadbalancerlistener.go index 1b30d17a9c..94f5f3bffb 100644 --- a/pkg/multicloud/aws/loadbalancerlistener.go +++ b/pkg/multicloud/aws/loadbalancerlistener.go @@ -19,6 +19,7 @@ import ( "fmt" "strconv" "strings" + "time" "github.com/aws/aws-sdk-go/service/elbv2" "github.com/pkg/errors" @@ -551,6 +552,18 @@ func (self *SRegion) CreateElbListener(listener *cloudprovider.SLoadbalancerList action.SetTargetGroupArn(listener.BackendGroupID) params.SetDefaultActions([]*elbv2.Action{action}) if listenerType == "HTTPS" { + err = cloudprovider.WaitCreated(3*time.Second, 30*time.Second, func() bool { + _, err := self.GetILoadBalancerCertificateById(listener.CertificateID) + if err == nil { + return true + } + + return false + }) + if err != nil { + return nil, errors.Wrap(err, "Region.GetILoadBalancerCertificateById") + } + cert := &elbv2.Certificate{ CertificateArn: &listener.CertificateID, } diff --git a/pkg/multicloud/aws/loadbalancerlistenerrule.go b/pkg/multicloud/aws/loadbalancerlistenerrule.go index e15a047442..4a5d083cf0 100644 --- a/pkg/multicloud/aws/loadbalancerlistenerrule.go +++ b/pkg/multicloud/aws/loadbalancerlistenerrule.go @@ -23,6 +23,7 @@ import ( "yunion.io/x/jsonutils" "yunion.io/x/log" + "yunion.io/x/pkg/errors" api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudprovider" @@ -202,8 +203,12 @@ func (self *SRegion) CreateElbListenerRule(listenerId string, config *cloudprovi return nil, err } + if len(ret.Rules) == 0 { + return nil, errors.Wrap(fmt.Errorf("empty rules"), "Region.CreateElbListenerRule.len") + } + rule := SElbListenerRule{} - err = unmarshalAwsOutput(ret, "", &rule) + err = unmarshalAwsOutput(ret.Rules[0], "", &rule) if err != nil { return nil, err } diff --git a/pkg/multicloud/aws/storagecache.go b/pkg/multicloud/aws/storagecache.go index 5b61cdc9fe..b1a9e5f802 100644 --- a/pkg/multicloud/aws/storagecache.go +++ b/pkg/multicloud/aws/storagecache.go @@ -225,30 +225,14 @@ func (self *SStoragecache) uploadImage(ctx context.Context, userCred mcclient.To } // todo:// 等待镜像导入完成 - for i := 1; i < 120; i++ { - time.Sleep(2 * time.Minute) - ret, err := self.region.ec2Client.DescribeImportImageTasks(&ec2.DescribeImportImageTasksInput{ImportTaskIds: []*string{&task.TaskId}}) - if err != nil { - return "", err - } - - err = FillZero(ret) - if err != nil { - return "", err - } - - log.Debugf("DescribeImportImage Task %s", ret.String()) - for _, item := range ret.ImportImageTasks { - if *item.Status == "completed" { - // add name tag - self.region.addTags(*item.ImageId, "Name", image.ImageId) - return *item.ImageId, nil - } - } + err = cloudprovider.WaitStatusWithDelay(task, ImageImportStatusCompleted, 2*time.Minute, 2*time.Minute, 4*time.Hour) + if err != nil { + return "", errors.Wrap(err, "SStoragecache.WaitStatusWithDelay") } - return task.ImageId, fmt.Errorf("uploadImage uncompleted: %s", task) - + // add name tag + self.region.addTags(task.ImageId, "Name", image.ImageId) + return task.ImageId, nil } func (self *SStoragecache) downloadImage(userCred mcclient.TokenCredential, imageId string, extId string) (jsonutils.JSONObject, error) { diff --git a/pkg/multicloud/huawei/eip.go b/pkg/multicloud/huawei/eip.go index 24a68b2a9a..876b2f0256 100644 --- a/pkg/multicloud/huawei/eip.go +++ b/pkg/multicloud/huawei/eip.go @@ -248,7 +248,7 @@ func (self *SEipAddress) Associate(instanceId string) error { return err } - err = cloudprovider.WaitStatus(self, api.EIP_STATUS_READY, 10*time.Second, 180*time.Second) + err = cloudprovider.WaitStatusWithDelay(self, api.EIP_STATUS_READY, 10*time.Second, 10*time.Second, 180*time.Second) return err } diff --git a/pkg/multicloud/huawei/loadbalancer_backendgroup.go b/pkg/multicloud/huawei/loadbalancer_backendgroup.go index 66c67560fd..8f11425b9d 100644 --- a/pkg/multicloud/huawei/loadbalancer_backendgroup.go +++ b/pkg/multicloud/huawei/loadbalancer_backendgroup.go @@ -17,8 +17,10 @@ package huawei import ( "fmt" "strings" + "time" "yunion.io/x/jsonutils" + "yunion.io/x/pkg/errors" api "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudprovider" @@ -46,6 +48,10 @@ func (self *SElbBackendGroup) GetLoadbalancerId() string { return self.lb.GetId() } +func (self *SElbBackendGroup) GetILoadbalancer() cloudprovider.ICloudLoadbalancer { + return self.lb +} + type StickySession struct { Type string `json:"type"` CookieName string `json:"cookie_name"` @@ -289,18 +295,37 @@ func (self *SElbBackendGroup) AddBackendServer(serverId string, weight int, port } func (self *SElbBackendGroup) RemoveBackendServer(backendId string, weight int, port int) error { - return self.region.RemoveLoadBalancerBackend(self.GetId(), backendId) + ibackend, err := self.GetILoadbalancerBackendById(backendId) + if err != nil { + if err == cloudprovider.ErrNotFound { + return nil + } + + return errors.Wrap(err, "ElbBackendGroup.GetILoadbalancerBackendById") + } + + err = self.region.RemoveLoadBalancerBackend(self.GetId(), backendId) + if err != nil { + return errors.Wrap(err, "ElbBackendGroup.RemoveBackendServer") + } + + return cloudprovider.WaitDeleted(ibackend, 2*time.Second, 30*time.Second) } func (self *SElbBackendGroup) Delete() error { if len(self.HealthMonitorID) > 0 { err := self.region.DeleteLoadbalancerHealthCheck(self.HealthMonitorID) if err != nil { - return err + return errors.Wrap(err, "ElbBackendGroup.Delete.DeleteLoadbalancerHealthCheck") } } - return self.region.DeleteLoadBalancerBackendGroup(self.GetId()) + err := self.region.DeleteLoadBalancerBackendGroup(self.GetId()) + if err != nil { + return errors.Wrap(err, "ElbBackendGroup.Delete.DeleteLoadBalancerBackendGroup") + } + + return cloudprovider.WaitDeleted(self, 2*time.Second, 30*time.Second) } func (self *SElbBackendGroup) Sync(group *cloudprovider.SLoadbalancerBackendGroup) error { diff --git a/pkg/multicloud/huawei/loadbalancer_listener.go b/pkg/multicloud/huawei/loadbalancer_listener.go index 440cc41e6e..f22cb3e327 100644 --- a/pkg/multicloud/huawei/loadbalancer_listener.go +++ b/pkg/multicloud/huawei/loadbalancer_listener.go @@ -418,7 +418,7 @@ func (self *SElbListener) GetILoadbalancerListenerRules() ([]cloudprovider.IClou rule := ret[i] rule.listener = self rule.lb = self.lb - iret = append(iret, &ret[i]) + iret = append(iret, &rule) } return iret, nil diff --git a/pkg/multicloud/qcloud/loadbalancer.go b/pkg/multicloud/qcloud/loadbalancer.go index e37a5884f6..44b5bfc5f3 100644 --- a/pkg/multicloud/qcloud/loadbalancer.go +++ b/pkg/multicloud/qcloud/loadbalancer.go @@ -404,7 +404,7 @@ func (self *SRegion) GetLoadbalancers(ids []string) ([]SLoadbalancer, error) { func (self *SRegion) GetLoadbalancer(id string) (*SLoadbalancer, error) { if len(id) == 0 { - return nil, fmt.Errorf("GetLoadbalancer id should not empty") + return nil, fmt.Errorf("GetILoadbalancer id should not empty") } lbs, err := self.GetLoadbalancers([]string{id}) @@ -418,7 +418,7 @@ func (self *SRegion) GetLoadbalancer(id string) (*SLoadbalancer, error) { case 1: return &lbs[0], nil default: - return nil, fmt.Errorf("GetLoadbalancer %s found %d", id, len(lbs)) + return nil, fmt.Errorf("GetILoadbalancer %s found %d", id, len(lbs)) } } diff --git a/pkg/multicloud/qcloud/loadbalancer_backendgroup.go b/pkg/multicloud/qcloud/loadbalancer_backendgroup.go index a7914b6de8..ae76b473d5 100644 --- a/pkg/multicloud/qcloud/loadbalancer_backendgroup.go +++ b/pkg/multicloud/qcloud/loadbalancer_backendgroup.go @@ -37,6 +37,10 @@ func (self *SLBBackendGroup) GetLoadbalancerId() string { return self.lb.GetId() } +func (self *SLBBackendGroup) GetILoadbalancer() cloudprovider.ICloudLoadbalancer { + return self.lb +} + func (self *SLBBackendGroup) GetProtocolType() string { return "" }