From cea9cf95b5cc952292e21f293f8b4c21ac1ceb4d Mon Sep 17 00:00:00 2001 From: Qu Xuan Date: Thu, 23 Apr 2020 17:49:09 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20cloudevent=E4=BD=BF=E7=94=A8=E4=BB=A3?= =?UTF-8?q?=E7=90=86=E4=BC=98=E5=8C=96?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pkg/cloudevent/models/cloudproviders.go | 167 +++++++++++++++--------- 1 file changed, 106 insertions(+), 61 deletions(-) diff --git a/pkg/cloudevent/models/cloudproviders.go b/pkg/cloudevent/models/cloudproviders.go index 448b046a80..fb0295c8cf 100644 --- a/pkg/cloudevent/models/cloudproviders.go +++ b/pkg/cloudevent/models/cloudproviders.go @@ -36,8 +36,10 @@ import ( "yunion.io/x/onecloud/pkg/cloudevent/options" "yunion.io/x/onecloud/pkg/cloudprovider" "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/mcclient/auth" "yunion.io/x/onecloud/pkg/mcclient/modules" "yunion.io/x/onecloud/pkg/s3gateway/session" + "yunion.io/x/onecloud/pkg/util/httputils" ) type SCloudproviderManager struct { @@ -61,39 +63,27 @@ func init() { type SCloudprovider struct { db.SEnabledStatusStandaloneResourceBase - HealthStatus string `width:"16" charset:"ascii" default:"normal" nullable:"false" list:"domain"` SyncStatus string LastSync time.Time LastSyncEndAt time.Time LastSyncTimeAt time.Time - AccessUrl string `width:"64" charset:"ascii" nullable:"true" list:"domain" update:"domain"` - Account string `width:"128" charset:"ascii" nullable:"false" list:"domain"` - Secret string `length:"0" charset:"ascii" nullable:"false" list:"domain"` - Provider string `width:"64" charset:"ascii" list:"domain"` Brand string `width:"64" charset:"ascii" list:"domain"` - - ProxySetting *proxyapi.SProxySetting } func (manager *SCloudproviderManager) GetRegionCloudproviders(ctx context.Context, userCred mcclient.TokenCredential) ([]SCloudprovider, error) { s := session.GetSession(ctx, userCred) params := jsonutils.NewDict() - params.Add(jsonutils.JSONTrue, "details") result, err := modules.Cloudproviders.List(s, params) if err != nil { return nil, errors.Wrap(err, "modules.Cloudproviders.List") } + providers := []SCloudprovider{} - for _, _provider := range result.Data { - provider := SCloudprovider{} - provider.SetModelManager(manager, &provider) - err = _provider.Unmarshal(&provider) - if err != nil { - return nil, errors.Wrap(err, "_provider.Unmarshal") - } - providers = append(providers, provider) + err = jsonutils.Update(&providers, result.Data) + if err != nil { + return nil, errors.Wrap(err, "jsonutils.Update") } return providers, nil } @@ -108,17 +98,18 @@ func (manager *SCloudproviderManager) GetLocalCloudproviders() ([]SCloudprovider return dbProviders, nil } -func (manager *SCloudproviderManager) SyncCloudproviders(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) { +func (manager *SCloudproviderManager) syncCloudproviders(ctx context.Context, userCred mcclient.TokenCredential) compare.SyncResult { + result := compare.SyncResult{} providers, err := manager.GetRegionCloudproviders(ctx, userCred) if err != nil { - log.Errorf("failed to get region cloudproviders: %v", err) - return + result.Error(errors.Wrap(err, "GetRegionCloudproviders")) + return result } dbProviders, err := manager.GetLocalCloudproviders() if err != nil { - log.Errorf("failed to get local cloudproviders: %v", err) - return + result.Error(errors.Wrap(err, "GetLocalCloudproviders")) + return result } removed := make([]SCloudprovider, 0) @@ -128,36 +119,48 @@ func (manager *SCloudproviderManager) SyncCloudproviders(ctx context.Context, us err = compare.CompareSets(dbProviders, providers, &removed, &commondb, &commonext, &added) if err != nil { - log.Errorf("compare.CompareSets: %v", err) - return + result.Error(errors.Wrap(err, "CompareSets")) + return result } for i := 0; i < len(removed); i++ { err = removed[i].Delete(ctx, userCred) if err != nil { - log.Errorf("failed to remove cloudprovider %s(%s) error: %v", removed[i].Name, removed[i].Id, err) + result.DeleteError(err) + continue } + result.Delete() } for i := 0; i < len(commondb); i++ { err = commondb[i].syncWithRegionProvider(ctx, userCred, commonext[i]) if err != nil { - log.Errorf("failed to sync cloudprovider %s(%s) error: %v", commondb[i].Name, commondb[i].Id, err) + result.UpdateError(err) + continue } + result.Update() } for i := 0; i < len(added); i++ { err = manager.newFromRegionProvider(ctx, userCred, added[i]) if err != nil { - log.Errorf("failed to add cloudprovider %s(%s) error: %v", added[i].Name, added[i].Id, err) + result.AddError(err) + continue } + result.Add() } + return result +} + +func (manager *SCloudproviderManager) SyncCloudproviders(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) { + result := manager.syncCloudproviders(ctx, userCred) + info := result.Result() + log.Infof("sync cloudproviders result: %s", info) } func (provider *SCloudprovider) syncWithRegionProvider(ctx context.Context, userCred mcclient.TokenCredential, cloudprovider SCloudprovider) error { _, err := db.Update(provider, func() error { provider.Status = cloudprovider.Status - provider.Secret = cloudprovider.Secret provider.Enabled = cloudprovider.Enabled provider.Brand = cloudprovider.Brand return nil @@ -173,8 +176,7 @@ func (self *SCloudprovider) MarkSyncing(userCred mcclient.TokenCredential) error return nil }) if err != nil { - log.Errorf("Failed to MarkSyncing error: %v", err) - return err + return errors.Wrap(err, "db.Update") } return nil } @@ -186,8 +188,7 @@ func (self *SCloudprovider) MarkEndSync(userCred mcclient.TokenCredential) error return nil }) if err != nil { - log.Errorf("Failed to markEndSync error: %v", err) - return err + return errors.Wrap(err, "db.Update") } return nil } @@ -198,8 +199,7 @@ func (self *SCloudprovider) SetLastSyncTimeAt(userCred mcclient.TokenCredential, return nil }) if err != nil { - log.Errorf("Failed to SetLastSyncTimeAt error: %v", err) - return err + return errors.Wrap(err, "db.Update") } return nil } @@ -209,30 +209,36 @@ func (manager *SCloudproviderManager) newFromRegionProvider(ctx context.Context, return manager.TableSpec().Insert(&cloudprovider) } -func (manager *SCloudproviderManager) SyncCloudeventTask(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) { +func (manager *SCloudproviderManager) syncCloudeventTask(ctx context.Context, userCred mcclient.TokenCredential) error { cloudproviders := []SCloudprovider{} q := manager.Query().IsTrue("enabled").Equals("status", api.CLOUD_PROVIDER_CONNECTED).Equals("sync_status", api.CLOUD_PROVIDER_SYNC_STATUS_IDLE) err := db.FetchModelObjects(manager, q, &cloudproviders) if err != nil { - log.Errorf("failed to fetch cloudproviders") - return + return errors.Wrap(err, "db.FetchModelObjects") } - for _, provider := range cloudproviders { - err = provider.StartCloudeventSyncTask(ctx, userCred) + for i := range cloudproviders { + err = cloudproviders[i].StartCloudeventSyncTask(ctx, userCred) if err != nil { - log.Errorf("Failed start cloudevent sync task error: %v", err) + log.Errorf("Failed start cloudevent sync task for cloudprovider %s (%s) error: %v", cloudproviders[i].Name, cloudproviders[i].Id, err) } } - return + return nil +} + +func (manager *SCloudproviderManager) SyncCloudeventTask(ctx context.Context, userCred mcclient.TokenCredential, isStart bool) { + err := manager.syncCloudeventTask(ctx, userCred) + if err != nil { + log.Errorf("syncCloudeventTask error: %v", err) + } } func (provider *SCloudprovider) StartCloudeventSyncTask(ctx context.Context, userCred mcclient.TokenCredential) error { params := jsonutils.NewDict() - provider.MarkSyncing(userCred) task, err := taskman.TaskManager.NewTask(ctx, "CloudeventSyncTask", provider, userCred, params, "", "", nil) if err != nil { return errors.Wrap(err, "NewTask") } + provider.MarkSyncing(userCred) task.ScheduleRun(nil) return nil } @@ -277,11 +283,43 @@ func (self *SCloudprovider) GetNextTimeRange() (time.Time, time.Time, error) { return start, end, nil } -func (provider *SCloudprovider) getPassword() (string, error) { +type SCloudproviderDelegate struct { + Id string + Name string + + Enabled bool + Status string + SyncStatus string + + AccessUrl string + Account string + Secret string + + Provider string + Brand string + + ProxySetting proxyapi.SProxySetting +} + +func (self *SCloudprovider) GetDelegate() (*SCloudproviderDelegate, error) { + s := auth.GetAdminSession(context.Background(), options.Options.Region, "") + result, err := modules.Cloudproviders.Get(s, self.Id, nil) + if err != nil { + return nil, errors.Wrap(err, "modules.Cloudproviders.Get") + } + provider := &SCloudproviderDelegate{} + err = result.Unmarshal(provider) + if err != nil { + return nil, errors.Wrap(err, "result.Unmarshal") + } + return provider, nil +} + +func (provider *SCloudproviderDelegate) getPassword() (string, error) { return utils.DescryptAESBase64(provider.Id, provider.Secret) } -func (provider *SCloudprovider) getAccessUrl() string { +func (provider *SCloudproviderDelegate) getAccessUrl() string { return provider.AccessUrl } @@ -297,33 +335,40 @@ func (provider SCloudprovider) GetExternalId() string { return provider.Id } -func (provider *SCloudprovider) GetProvider() (cloudprovider.ICloudProvider, error) { - if !provider.GetEnabled() { +func (self *SCloudprovider) GetProvider() (cloudprovider.ICloudProvider, error) { + if !self.GetEnabled() { return nil, errors.Error("Cloud provider is not enabled") } - accessUrl := provider.getAccessUrl() - passwd, err := provider.getPassword() + delegate, err := self.GetDelegate() if err != nil { - return nil, err + return nil, errors.Wrap(err, "GetDelegate") + } + accessUrl := delegate.getAccessUrl() + passwd, err := delegate.getPassword() + if err != nil { + return nil, errors.Wrap(err, "delegate.getPassword") } - ps := provider.ProxySetting - cfg := &httpproxy.Config{ - HTTPProxy: ps.HTTPProxy, - HTTPSProxy: ps.HTTPSProxy, - NoProxy: ps.NoProxy, - } - cfgProxyFunc := cfg.ProxyFunc() - proxyFunc := func(req *http.Request) (*url.URL, error) { - return cfgProxyFunc(req.URL) + var proxyFunc httputils.TransportProxyFunc + { + cfg := &httpproxy.Config{ + HTTPProxy: delegate.ProxySetting.HTTPProxy, + HTTPSProxy: delegate.ProxySetting.HTTPSProxy, + NoProxy: delegate.ProxySetting.NoProxy, + } + cfgProxyFunc := cfg.ProxyFunc() + proxyFunc = func(req *http.Request) (*url.URL, error) { + return cfgProxyFunc(req.URL) + } } + return cloudprovider.GetProvider( cloudprovider.ProviderConfig{ - Id: provider.Id, - Name: provider.Name, - Vendor: provider.Provider, + Id: self.Id, + Name: self.Name, + Vendor: self.Provider, URL: accessUrl, - Account: provider.Account, + Account: delegate.Account, Secret: passwd, ProxyFunc: proxyFunc,