fix(scheduler): reload cloudprovider's hosts when syncing finished

This commit is contained in:
Zexi Li
2023-12-15 20:19:32 +08:00
parent 644b2a20ab
commit fff91b8eca
9 changed files with 127 additions and 17 deletions
+11 -2
View File
@@ -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()
})
}
+2
View File
@@ -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
+1 -1
View File
@@ -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
@@ -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
@@ -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
}
@@ -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
}
+33 -3
View File
@@ -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)
}
}
}
+5 -2
View File
@@ -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")
}
+1 -1
View File
@@ -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,