From b30d21b110e5db22effcaa3b30f9f4d485c65837 Mon Sep 17 00:00:00 2001 From: Zexi Li Date: Wed, 27 Sep 2023 16:39:59 +0800 Subject: [PATCH] fix(scheduler): some joint resources are not synced --- pkg/scheduler/data_manager/common/common.go | 37 ++++++++++++------- .../data_manager/netinterface/netinterface.go | 23 ++++++++++++ 2 files changed, 47 insertions(+), 13 deletions(-) diff --git a/pkg/scheduler/data_manager/common/common.go b/pkg/scheduler/data_manager/common/common.go index ac8f59d8dd..aaa14a4b7d 100644 --- a/pkg/scheduler/data_manager/common/common.go +++ b/pkg/scheduler/data_manager/common/common.go @@ -91,19 +91,22 @@ func (m *CommonResourceManager[O]) SyncOnce() error { return m.GetStore().Init() } +type FGetDBObject func(man db.IModelManager, id string, obj *jsonutils.JSONDict) (db.IModel, error) + type ResourceStore[O lockman.ILockedObject] struct { - dataMap *sync.Map - modelMan db.IModelManager - res informer.IResourceManager - getId func(O) string - getWatchId func(*jsonutils.JSONDict) string + dataMap *sync.Map + modelMan db.IModelManager + res informer.IResourceManager + getId func(O) string + getWatchId func(*jsonutils.JSONDict) string + getDBObject FGetDBObject } func NewResourceStore[O lockman.ILockedObject]( modelMan db.IModelManager, res informer.IResourceManager, ) IResourceStore[O] { - return newResourceStore[O](modelMan, res, nil, nil) + return newResourceStore[O](modelMan, res, nil, nil, nil) } func NewJointResourceStore[O lockman.ILockedObject]( @@ -111,8 +114,9 @@ func NewJointResourceStore[O lockman.ILockedObject]( res informer.IResourceManager, getId func(O) string, getWatchId func(*jsonutils.JSONDict) string, + getDBObject FGetDBObject, ) IResourceStore[O] { - return newResourceStore(modelMan, res, getId, getWatchId) + return newResourceStore(modelMan, res, getId, getWatchId, getDBObject) } func newResourceStore[O lockman.ILockedObject]( @@ -120,6 +124,7 @@ func newResourceStore[O lockman.ILockedObject]( res informer.IResourceManager, getId func(O) string, getWatchId func(*jsonutils.JSONDict) string, + getDBObject FGetDBObject, ) IResourceStore[O] { if getId == nil { getId = func(o O) string { @@ -132,12 +137,18 @@ func newResourceStore[O lockman.ILockedObject]( return id } } + if getDBObject == nil { + getDBObject = func(man db.IModelManager, id string, o *jsonutils.JSONDict) (db.IModel, error) { + return man.FetchById(id) + } + } return &ResourceStore[O]{ - dataMap: new(sync.Map), - modelMan: modelMan, - res: res, - getId: getId, - getWatchId: getWatchId, + dataMap: new(sync.Map), + modelMan: modelMan, + res: res, + getId: getId, + getWatchId: getWatchId, + getDBObject: getDBObject, } } @@ -189,7 +200,7 @@ func (s *ResourceStore[O]) GetAll() []O { func (s *ResourceStore[O]) Add(obj *jsonutils.JSONDict) { id := s.getWatchId(obj) if id != "" { - dbObj, err := s.modelMan.FetchById(id) + dbObj, err := s.getDBObject(s.modelMan, id, obj) if err == nil { v := reflect.ValueOf(dbObj) tmpObj := v.Elem().Interface() diff --git a/pkg/scheduler/data_manager/netinterface/netinterface.go b/pkg/scheduler/data_manager/netinterface/netinterface.go index 0e381d91ea..434a71561e 100644 --- a/pkg/scheduler/data_manager/netinterface/netinterface.go +++ b/pkg/scheduler/data_manager/netinterface/netinterface.go @@ -2,10 +2,13 @@ package netinterface import ( "fmt" + "strings" "time" "yunion.io/x/jsonutils" + "yunion.io/x/pkg/errors" + "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" @@ -48,6 +51,26 @@ func NewResourceStore() common.IResourceStore[models.SNetInterface] { vlan, _ := o.Int("vlan_id") return GetId(hostId, wireId, mac, int(vlan)) }, + func(man db.IModelManager, id string, obj *jsonutils.JSONDict) (db.IModel, error) { + ids := strings.Split(id, "/") + if len(ids) != 3 { + return nil, errors.Errorf("Invalid id: %q", id) + } + q := man.Query() + hostId := ids[0] + wireId := ids[1] + mac := ids[2] + q = q.Equals("host_id", hostId).Equals("wire_id", wireId).Equals("mac", mac) + objs, err := db.FetchIModelObjects(man, q) + errHint := fmt.Sprintf("hostId %q, wireId %q, mac %q", hostId, wireId, mac) + if err != nil { + return nil, errors.Wrapf(err, "db.FetchIModelObjects by %s", errHint) + } + if len(objs) != 1 { + return nil, errors.Errorf("Found %d objects by %s", len(objs), errHint) + } + return objs[0], nil + }, ) }