From fff91b8eca5c7d6d8fb53732798ffb8882362746 Mon Sep 17 00:00:00 2001 From: Zexi Li Date: Fri, 15 Dec 2023 17:10:23 +0800 Subject: [PATCH] fix(scheduler): reload cloudprovider's hosts when syncing finished --- pkg/cloudcommon/db/fetch.go | 13 ++++- pkg/compute/models/hosts.go | 2 + pkg/scheduler/cache/candidate/builder.go | 2 +- .../data_manager/candidate_manager.go | 16 ++++-- .../cloudprovider/cloudprovider.go | 51 +++++++++++++++++-- .../data_manager/common/cache_manager.go | 15 ++++++ pkg/scheduler/data_manager/common/common.go | 36 +++++++++++-- pkg/scheduler/manager/manager.go | 7 ++- pkg/scheduler/service/service.go | 2 +- 9 files changed, 127 insertions(+), 17 deletions(-) create mode 100644 pkg/scheduler/data_manager/common/cache_manager.go diff --git a/pkg/cloudcommon/db/fetch.go b/pkg/cloudcommon/db/fetch.go index 48d8817768..c54ff5b6d6 100644 --- a/pkg/cloudcommon/db/fetch.go +++ b/pkg/cloudcommon/db/fetch.go @@ -569,8 +569,11 @@ func FetchStandaloneObjectsByIds(modelManager IModelManager, ids []string, targe return FetchModelObjectsByIds(modelManager, "id", ids, targets) } -func FetchDistinctField(modelManager IModelManager, field string) ([]string, error) { - q := modelManager.Query(field).Distinct() +func FetchField(modelMan IModelManager, field string, qCallback func(q *sqlchemy.SQuery) *sqlchemy.SQuery) ([]string, error) { + q := modelMan.Query(field) + if qCallback != nil { + q = qCallback(q) + } rows, err := q.Rows() if err != nil { if errors.Cause(err) == sql.ErrNoRows { @@ -593,3 +596,9 @@ func FetchDistinctField(modelManager IModelManager, field string) ([]string, err } return values, nil } + +func FetchDistinctField(modelManager IModelManager, field string) ([]string, error) { + return FetchField(modelManager, field, func(q *sqlchemy.SQuery) *sqlchemy.SQuery { + return q.Distinct() + }) +} diff --git a/pkg/compute/models/hosts.go b/pkg/compute/models/hosts.go index 3cfa9dca42..6de48fa26f 100644 --- a/pkg/compute/models/hosts.go +++ b/pkg/compute/models/hosts.go @@ -4664,6 +4664,7 @@ func (h *SHost) addNetif(ctx context.Context, userCred mcclient.TokenCredential, } // else not found netif = &SNetInterface{} + netif.SetModelManager(NetInterfaceManager, netif) netif.Mac = mac netif.VlanId = vlanId } @@ -5003,6 +5004,7 @@ func (hh *SHost) Attach2Network( } bn := &SHostnetwork{} bn.BaremetalId = hh.Id + bn.SetModelManager(HostnetworkManager, bn) bn.NetworkId = net.Id bn.IpAddr = freeIp bn.MacAddr = netif.Mac diff --git a/pkg/scheduler/cache/candidate/builder.go b/pkg/scheduler/cache/candidate/builder.go index 3565e0b8c4..0536c5aabc 100644 --- a/pkg/scheduler/cache/candidate/builder.go +++ b/pkg/scheduler/cache/candidate/builder.go @@ -220,7 +220,7 @@ func (b *baseBuilder) setCloudproviderAccounts(hosts []computemodels.SHost, errC } providerObjs := make([]computemodels.SCloudprovider, 0) for _, pId := range providerSets.List() { - pObj, ok := cloudprovider.Manager.GetResource(pId) + pObj, ok := cloudprovider.GetManager().GetResource(pId) if !ok { errCh <- errors.Errorf("Not found cloudprovider by id: %q", pId) return diff --git a/pkg/scheduler/data_manager/candidate_manager.go b/pkg/scheduler/data_manager/candidate_manager.go index eed2d36e35..e0cf6d1ace 100644 --- a/pkg/scheduler/data_manager/candidate_manager.go +++ b/pkg/scheduler/data_manager/candidate_manager.go @@ -322,6 +322,11 @@ func (cm *CandidateManager) AddImpl(name string, impl *CandidateManagerImpl) { cm.impls[name] = impl } +const ( + CANDIDATE_MANAGER_IMPL_HOST = "host" + CANDIDATE_MANAGER_IMPL_BAREMETAL = "baremetal" +) + func NewCandidateManager(dataManager *DataManager, stopCh <-chan struct{}) *CandidateManager { candidateManager := &CandidateManager{ @@ -331,10 +336,10 @@ func NewCandidateManager(dataManager *DataManager, stopCh <-chan struct{}) *Cand //dirtyPool: ttlpool.NewCountPool(), } - candidateManager.AddImpl("host", NewCandidateManagerImpl( + candidateManager.AddImpl(CANDIDATE_MANAGER_IMPL_HOST, NewCandidateManagerImpl( &HostCandidateManagerImplProvider{dataManager: dataManager}, stopCh)) - candidateManager.AddImpl("baremetal", NewCandidateManagerImpl( + candidateManager.AddImpl(CANDIDATE_MANAGER_IMPL_BAREMETAL, NewCandidateManagerImpl( &BaremetalCandidateManagerImplProvider{dataManager: dataManager}, stopCh)) return candidateManager @@ -347,8 +352,11 @@ func (cm *CandidateManager) Run() { } } -func (cm *CandidateManager) Reload(resType string, candidateIds []string) ( - []interface{}, error) { +func (cm *CandidateManager) ReloadHosts(ids []string) ([]interface{}, error) { + return cm.Reload(CANDIDATE_MANAGER_IMPL_HOST, ids) +} + +func (cm *CandidateManager) Reload(resType string, candidateIds []string) ([]interface{}, error) { if len(candidateIds) == 0 { return []interface{}{}, nil diff --git a/pkg/scheduler/data_manager/cloudprovider/cloudprovider.go b/pkg/scheduler/data_manager/cloudprovider/cloudprovider.go index 560e43b1f6..5d9161f350 100644 --- a/pkg/scheduler/data_manager/cloudprovider/cloudprovider.go +++ b/pkg/scheduler/data_manager/cloudprovider/cloudprovider.go @@ -15,17 +15,29 @@ package cloudprovider import ( + "fmt" "time" + "yunion.io/x/jsonutils" + "yunion.io/x/log" + "yunion.io/x/pkg/errors" + "yunion.io/x/sqlchemy" + + computeapi "yunion.io/x/onecloud/pkg/apis/compute" + "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/mcclient/modules/compute" "yunion.io/x/onecloud/pkg/scheduler/data_manager/common" ) -var Manager common.IResourceManager[models.SCloudprovider] +var manager common.IResourceManager[models.SCloudprovider] -func init() { - Manager = NewResourceManager() +func GetManager() common.IResourceManager[models.SCloudprovider] { + if manager != nil { + return manager + } + manager = NewResourceManager() + return manager } func NewResourceManager() common.IResourceManager[models.SCloudprovider] { @@ -41,5 +53,36 @@ func NewResourceStore() common.IResourceStore[models.SCloudprovider] { return common.NewResourceStore[models.SCloudprovider]( models.CloudproviderManager, compute.Cloudproviders, - ) + ).WithOnUpdate(onCloudproviderUpdate) +} + +func onCloudproviderUpdate(oldObj *jsonutils.JSONDict, newObj db.IModel) { + syncStatusKey := "sync_status" + if !oldObj.Contains(syncStatusKey) { + return + } + // process cloudaccount syncing finished status + cp := newObj.(*models.SCloudprovider) + prevStatus, _ := oldObj.GetString(syncStatusKey) + curStatus := cp.SyncStatus + if prevStatus == computeapi.CLOUD_PROVIDER_SYNC_STATUS_SYNCING && curStatus == computeapi.CLOUD_PROVIDER_SYNC_STATUS_IDLE { + if err := onCloudproviderSyncFinished(cp); err != nil { + log.Infof("onCloudproviderSyncFinished error: %v", err) + } + } +} + +func onCloudproviderSyncFinished(cp *models.SCloudprovider) error { + cpHint := fmt.Sprintf("%s/%s", cp.GetId(), cp.GetName()) + hostdIds, err := db.FetchField(models.HostManager, "id", func(q *sqlchemy.SQuery) *sqlchemy.SQuery { + return q.Equals("manager_id", cp.GetId()) + }) + if err != nil { + return errors.Wrapf(err, "get all hostdIds from cloudprovider %s", cpHint) + } + log.Infof("Start reload cloudprovider %s hosts: %v", cpHint, hostdIds) + if _, err := common.GetCacheManager().ReloadHosts(hostdIds); err != nil { + return errors.Wrapf(err, "Reload cache hosts of cloudprovider %s", cpHint) + } + return nil } diff --git a/pkg/scheduler/data_manager/common/cache_manager.go b/pkg/scheduler/data_manager/common/cache_manager.go new file mode 100644 index 0000000000..b842ee60ce --- /dev/null +++ b/pkg/scheduler/data_manager/common/cache_manager.go @@ -0,0 +1,15 @@ +package common + +var cacheManager CacheManager + +type CacheManager interface { + ReloadHosts(ids []string) ([]interface{}, error) +} + +func RegisterCacheManager(man CacheManager) { + cacheManager = man +} + +func GetCacheManager() CacheManager { + return cacheManager +} diff --git a/pkg/scheduler/data_manager/common/common.go b/pkg/scheduler/data_manager/common/common.go index 83c56d95c8..bd9bb81cb5 100644 --- a/pkg/scheduler/data_manager/common/common.go +++ b/pkg/scheduler/data_manager/common/common.go @@ -114,12 +114,15 @@ type ResourceStore[O lockman.ILockedObject] struct { getId func(O) string getWatchId func(*jsonutils.JSONDict) string getDBObject FGetDBObject + onAdd func(obj db.IModel) + onUpdate func(oldObj *jsonutils.JSONDict, newObj db.IModel) + onDelete func(obj *jsonutils.JSONDict) } func NewResourceStore[O lockman.ILockedObject]( modelMan db.IModelManager, res informer.IResourceManager, -) IResourceStore[O] { +) *ResourceStore[O] { return newResourceStore[O](modelMan, res, nil, nil, nil) } @@ -129,7 +132,7 @@ func NewJointResourceStore[O lockman.ILockedObject]( getId func(O) string, getWatchId func(*jsonutils.JSONDict) string, getDBObject FGetDBObject, -) IResourceStore[O] { +) *ResourceStore[O] { return newResourceStore(modelMan, res, getId, getWatchId, getDBObject) } @@ -139,7 +142,7 @@ func newResourceStore[O lockman.ILockedObject]( getId func(O) string, getWatchId func(*jsonutils.JSONDict) string, getDBObject FGetDBObject, -) IResourceStore[O] { +) *ResourceStore[O] { if getId == nil { getId = func(o O) string { return o.GetId() @@ -163,9 +166,27 @@ func newResourceStore[O lockman.ILockedObject]( getId: getId, getWatchId: getWatchId, getDBObject: getDBObject, + onAdd: nil, + onUpdate: nil, + onDelete: nil, } } +func (s *ResourceStore[O]) WithOnAdd(onAdd func(db.IModel)) *ResourceStore[O] { + s.onAdd = onAdd + return s +} + +func (s *ResourceStore[O]) WithOnUpdate(onUpdate func(old *jsonutils.JSONDict, newObj db.IModel)) *ResourceStore[O] { + s.onUpdate = onUpdate + return s +} + +func (s *ResourceStore[O]) WithOnDelete(onDelete func(*jsonutils.JSONDict)) *ResourceStore[O] { + s.onDelete = onDelete + return s +} + func (s *ResourceStore[O]) GetInformerResourceManager() informer.IResourceManager { return s.res } @@ -220,6 +241,9 @@ func (s *ResourceStore[O]) Add(obj *jsonutils.JSONDict) { tmpObj := v.Elem().Interface() s.dataMap.Store(id, tmpObj) log.Infof("Add %s %s", s.modelMan.Keyword(), obj.String()) + if s.onAdd != nil { + s.onAdd(dbObj) + } } else { log.Errorf("Fetch %s by id %s error when created: %v", s.modelMan.Keyword(), id, err) } @@ -250,6 +274,9 @@ func (s *ResourceStore[O]) Update(oldObj, newObj *jsonutils.JSONDict) { tmpObj := v.Elem().Interface() s.dataMap.Store(id, tmpObj) log.Infof("Update %s %s", s.modelMan.Keyword(), newObj.String()) + if s.onUpdate != nil { + s.onUpdate(oldObj, dbObj) + } } else { log.Errorf("Fetch %s by id %s error when updated: %v", s.modelMan.Keyword(), id, err) } @@ -261,6 +288,9 @@ func (s *ResourceStore[O]) Delete(obj *jsonutils.JSONDict) { if id != "" { s.dataMap.Delete(id) log.Infof("Delete %s %s", s.modelMan.Keyword(), obj.String()) + if s.onDelete != nil { + s.onDelete(obj) + } } } diff --git a/pkg/scheduler/manager/manager.go b/pkg/scheduler/manager/manager.go index 4051305a5c..0adcca6352 100644 --- a/pkg/scheduler/manager/manager.go +++ b/pkg/scheduler/manager/manager.go @@ -28,6 +28,7 @@ import ( candidatecache "yunion.io/x/onecloud/pkg/scheduler/cache/candidate" "yunion.io/x/onecloud/pkg/scheduler/core" "yunion.io/x/onecloud/pkg/scheduler/data_manager" + "yunion.io/x/onecloud/pkg/scheduler/data_manager/common" schedmodels "yunion.io/x/onecloud/pkg/scheduler/models" o "yunion.io/x/onecloud/pkg/scheduler/options" "yunion.io/x/onecloud/pkg/util/k8s" @@ -48,7 +49,7 @@ type SchedulerManager struct { KubeClusterManager *k8s.SKubeClusterManager } -func NewSchedulerManager(stopCh <-chan struct{}) *SchedulerManager { +func newSchedulerManager(stopCh <-chan struct{}) *SchedulerManager { sm := &SchedulerManager{} sm.DataManager = data_manager.NewDataManager(stopCh) sm.CandidateManager = data_manager.NewCandidateManager(sm.DataManager, stopCh) @@ -58,6 +59,8 @@ func NewSchedulerManager(stopCh <-chan struct{}) *SchedulerManager { sm.TaskManager = NewTaskManager(stopCh) sm.KubeClusterManager = k8s.NewKubeClusterManager(o.Options.Region, 30*time.Second) + common.RegisterCacheManager(sm.CandidateManager) + return sm } @@ -74,7 +77,7 @@ func InitAndStart(stopCh <-chan struct{}) { log.Warningf("Global scheduler already init.") return } - schedManager = NewSchedulerManager(stopCh) + schedManager = newSchedulerManager(stopCh) go schedManager.start() log.Infof("InitAndStart ok") } diff --git a/pkg/scheduler/service/service.go b/pkg/scheduler/service/service.go index 855a514265..80c3c5403f 100644 --- a/pkg/scheduler/service/service.go +++ b/pkg/scheduler/service/service.go @@ -125,7 +125,7 @@ func StartService() error { for _, f := range []func(ctx context.Context){ cloudregion.Manager.Start, zone.Manager.Start, - cloudprovider.Manager.Start, + cloudprovider.GetManager().Start, cloudaccount.Manager.Start, wire.Manager.Start, network.Manager.Start,