diff --git a/pkg/cloudcommon/db/modelbase.go b/pkg/cloudcommon/db/modelbase.go index 75c198bcea..fa7a772f3c 100644 --- a/pkg/cloudcommon/db/modelbase.go +++ b/pkg/cloudcommon/db/modelbase.go @@ -548,19 +548,19 @@ func (manager *SModelBaseManager) CustomizedTotalCount(ctx context.Context, user return ret.Count, nil, nil } -func (model *SModelBase) GetId() string { +func (model SModelBase) GetId() string { return "" } -func (model *SModelBase) Keyword() string { +func (model SModelBase) Keyword() string { return model.GetModelManager().Keyword() } -func (model *SModelBase) KeywordPlural() string { +func (model SModelBase) KeywordPlural() string { return model.GetModelManager().KeywordPlural() } -func (model *SModelBase) GetName() string { +func (model SModelBase) GetName() string { return "" } diff --git a/pkg/cloudcommon/db/standalone_anon.go b/pkg/cloudcommon/db/standalone_anon.go index 815dec3def..9b2e47e678 100644 --- a/pkg/cloudcommon/db/standalone_anon.go +++ b/pkg/cloudcommon/db/standalone_anon.go @@ -179,7 +179,7 @@ func (model *SStandaloneAnonResourceBase) StandaloneModelManager() IStandaloneMo return model.GetModelManager().(IStandaloneModelManager) } -func (model *SStandaloneAnonResourceBase) GetId() string { +func (model SStandaloneAnonResourceBase) GetId() string { return model.Id } diff --git a/pkg/compute/models/netinterfaces.go b/pkg/compute/models/netinterfaces.go index 78bf1bcf47..aead76cdd9 100644 --- a/pkg/compute/models/netinterfaces.go +++ b/pkg/compute/models/netinterfaces.go @@ -63,7 +63,7 @@ func init() { NetInterfaceManager.SetVirtualObject(NetInterfaceManager) } -func (netif *SNetInterface) GetId() string { +func (netif SNetInterface) GetId() string { return netif.Mac } diff --git a/pkg/scheduler/algorithm/predicates/guest/hypervisor_predicate.go b/pkg/scheduler/algorithm/predicates/guest/hypervisor_predicate.go index b762e714d2..95438121e4 100644 --- a/pkg/scheduler/algorithm/predicates/guest/hypervisor_predicate.go +++ b/pkg/scheduler/algorithm/predicates/guest/hypervisor_predicate.go @@ -22,6 +22,7 @@ import ( "yunion.io/x/onecloud/pkg/scheduler/algorithm/predicates" "yunion.io/x/onecloud/pkg/scheduler/api" "yunion.io/x/onecloud/pkg/scheduler/core" + "yunion.io/x/onecloud/pkg/scheduler/data_manager/schedtag" ) const ( @@ -43,9 +44,9 @@ func (f *HypervisorPredicate) Clone() core.FitPredicate { } func hostHasContainerTag(c core.Candidater) bool { - aggs := c.Getter().HostSchedtags() + aggs := schedtag.GetCandidateSchedtags("hosts", c.Getter().Host().GetId()) for _, agg := range aggs { - if agg.Name == CONTAINER_ALLOWED_TAG { + if agg.GetName() == CONTAINER_ALLOWED_TAG { return true } } diff --git a/pkg/scheduler/algorithm/predicates/host_schedtag_predicate.go b/pkg/scheduler/algorithm/predicates/host_schedtag_predicate.go index 537fa73e08..13ab573375 100644 --- a/pkg/scheduler/algorithm/predicates/host_schedtag_predicate.go +++ b/pkg/scheduler/algorithm/predicates/host_schedtag_predicate.go @@ -96,10 +96,6 @@ func (r hostSchedtagResW) Keyword() string { return r.Candidater.Getter().Host().Keyword() } -func (r hostSchedtagResW) GetSchedtags() []computemodels.SSchedtag { - return r.Candidater.Getter().HostSchedtags() -} - func (r hostSchedtagResW) GetSchedtagJointManager() computemodels.ISchedtagJointManager { return r.Candidater.Getter().Host().GetSchedtagJointManager() } diff --git a/pkg/scheduler/algorithm/predicates/zone_schedtag_predicate.go b/pkg/scheduler/algorithm/predicates/zone_schedtag_predicate.go index 32f27c21e4..808233f8c5 100644 --- a/pkg/scheduler/algorithm/predicates/zone_schedtag_predicate.go +++ b/pkg/scheduler/algorithm/predicates/zone_schedtag_predicate.go @@ -96,11 +96,6 @@ func (r zoneSchedtagResW) Keyword() string { return r.zone.Keyword() } -func (r zoneSchedtagResW) GetSchedtags() []computemodels.SSchedtag { - // TODO: not fetch schedtags from database, should get zone schedtags from r.c cache candidater - return r.zone.GetSchedtags() -} - func (r zoneSchedtagResW) GetSchedtagJointManager() computemodels.ISchedtagJointManager { return r.zone.GetSchedtagJointManager() } diff --git a/pkg/scheduler/api/types.go b/pkg/scheduler/api/types.go index 28c38f6fe5..8ca1fddc96 100644 --- a/pkg/scheduler/api/types.go +++ b/pkg/scheduler/api/types.go @@ -95,15 +95,13 @@ func SchedtagStrategyCheck(strategy string) (err error) { type CandidateStorage struct { *models.SStorage - Schedtags []models.SSchedtag `json:"schedtags"` - FreeCapacity int64 `json:"free_capacity"` - ActualFreeCapacity int64 `json:"actual_free_capacity"` + FreeCapacity int64 `json:"free_capacity"` + ActualFreeCapacity int64 `json:"actual_free_capacity"` } type CandidateNetwork struct { *models.SNetwork - Schedtags []models.SSchedtag `json:"schedtags"` - FreePort int `json:"free_port"` + FreePort int `json:"free_port"` Provider string VpcId string diff --git a/pkg/scheduler/cache/cache.go b/pkg/scheduler/cache/cache.go index bf3dcc5979..8fbebd19e0 100644 --- a/pkg/scheduler/cache/cache.go +++ b/pkg/scheduler/cache/cache.go @@ -18,6 +18,7 @@ import ( "fmt" "reflect" "sync" + "time" "yunion.io/x/log" expirationcache "yunion.io/x/pkg/util/cache" @@ -152,9 +153,10 @@ func (c *schedulerCache) updateAllObjects() { func (c *schedulerCache) loadObjects(ids []string) ([]interface{}, error) { log.Infof("Start load %s, period: %v, ttl: %v", c.Name(), c.item.Period(), c.item.TTL()) + startTime := time.Now() defer func() { - log.Infof("End load %s", c.Name()) + log.Infof("End load %s, elapsed %s", c.Name(), time.Since(startTime)) }() var ( diff --git a/pkg/scheduler/cache/candidate/baremetals.go b/pkg/scheduler/cache/candidate/baremetals.go index b547315654..3b8c1b1169 100644 --- a/pkg/scheduler/cache/candidate/baremetals.go +++ b/pkg/scheduler/cache/candidate/baremetals.go @@ -15,12 +15,7 @@ package candidate import ( - "fmt" - "strings" - "time" - "yunion.io/x/jsonutils" - "yunion.io/x/log" "yunion.io/x/sqlchemy" computeapi "yunion.io/x/onecloud/pkg/apis/compute" @@ -133,52 +128,17 @@ func (bd *BaremetalDesc) FreeStorageSize() int64 { } func newBaremetalBuilder() *BaremetalBuilder { - return &BaremetalBuilder{ - baseBuilder: newBaseBuilder(BaremetalDescBuilder), - } + builder := new(BaremetalBuilder) + builder.baseBuilder = newBaseBuilder(BaremetalDescBuilder, builder) + return builder } -func (bb *BaremetalBuilder) init(ids []string) error { - bms, err := FetchHostsByIds(ids) - if err != nil { - return err - } - - //bb.baremetalAgents = agents - bb.baremetals = bms - - wg := &WaitGroupWrapper{} - errMessageChannel := make(chan error, 2) - defer close(errMessageChannel) - - setFuncs := []func(){ - func() { bb.setIsolatedDevs(ids, errMessageChannel) }, - } - - for _, f := range setFuncs { - wg.Wrap(f) - } - - if ok := waitTimeOut(wg, time.Duration(20*time.Second)); !ok { - log.Errorln("BaremetalBuilder waitgroup timeout.") - } - - if len(errMessageChannel) != 0 { - errMessages := make([]string, 0) - lengthChan := len(errMessageChannel) - for ; lengthChan >= 0; lengthChan-- { - errMessages = append(errMessages, fmt.Sprintf("%s", <-errMessageChannel)) - } - return fmt.Errorf("%s\n", strings.Join(errMessages, ";")) - } - - return nil +func (bb *BaremetalBuilder) FetchHosts(ids []string) ([]computemodels.SHost, error) { + return FetchHostsByIds(ids) } func (bb *BaremetalBuilder) Clone() BuildActor { - return &BaremetalBuilder{ - baseBuilder: newBaseBuilder(BaremetalDescBuilder), - } + return newBaremetalBuilder() } func (bb *BaremetalBuilder) AllIDs() ([]string, error) { @@ -187,37 +147,11 @@ func (bb *BaremetalBuilder) AllIDs() ([]string, error) { return FetchModelIds(q) } -func (bb *BaremetalBuilder) Do(ids []string) ([]interface{}, error) { - err := bb.init(ids) - if err != nil { - return nil, err - } - netGetter := newNetworkGetter() - descs, err := bb.build(netGetter) - if err != nil { - return nil, err - } - return descs, nil +func (bb *BaremetalBuilder) InitFuncs() []InitFunc { + return nil } -func (bb *BaremetalBuilder) build(netGetter *networkGetter) ([]interface{}, error) { - schedDescs := []interface{}{} - for _, bm := range bb.baremetals { - desc, err := bb.buildOne(&bm, netGetter) - if err != nil { - log.Errorf("BaremetalBuilder error: %v", err) - continue - } - schedDescs = append(schedDescs, desc) - } - return schedDescs, nil -} - -func (bb *BaremetalBuilder) buildOne(hostObj *computemodels.SHost, netGetter *networkGetter) (interface{}, error) { - baseDesc, err := newBaseHostDesc(bb.baseBuilder, hostObj, netGetter) - if err != nil { - return nil, err - } +func (bb *BaremetalBuilder) BuildOne(hostObj *computemodels.SHost, netGetter *networkGetter, baseDesc *BaseHostDesc) (interface{}, error) { desc := &BaremetalDesc{ BaseHostDesc: baseDesc, } @@ -230,7 +164,7 @@ func (bb *BaremetalBuilder) buildOne(hostObj *computemodels.SHost, netGetter *ne desc.StorageInfo = baremetalStorages desc.Tenants = make(map[string]int64, 0) - err = bb.fillServerID(desc, hostObj) + err := bb.fillServerID(desc, hostObj) if err != nil { return nil, err } diff --git a/pkg/scheduler/cache/candidate/base.go b/pkg/scheduler/cache/candidate/base.go index c62ff70061..b88ca517e8 100644 --- a/pkg/scheduler/cache/candidate/base.go +++ b/pkg/scheduler/cache/candidate/base.go @@ -22,17 +22,21 @@ import ( "yunion.io/x/jsonutils" "yunion.io/x/log" "yunion.io/x/pkg/errors" + "yunion.io/x/pkg/util/sets" "yunion.io/x/pkg/utils" "yunion.io/x/sqlchemy" computeapi "yunion.io/x/onecloud/pkg/apis/compute" - computedb "yunion.io/x/onecloud/pkg/cloudcommon/db" "yunion.io/x/onecloud/pkg/cloudcommon/types" - "yunion.io/x/onecloud/pkg/compute/models" computemodels "yunion.io/x/onecloud/pkg/compute/models" "yunion.io/x/onecloud/pkg/scheduler/api" "yunion.io/x/onecloud/pkg/scheduler/core" + "yunion.io/x/onecloud/pkg/scheduler/data_manager/cloudregion" + "yunion.io/x/onecloud/pkg/scheduler/data_manager/hostwire" + "yunion.io/x/onecloud/pkg/scheduler/data_manager/netinterface" + "yunion.io/x/onecloud/pkg/scheduler/data_manager/network" "yunion.io/x/onecloud/pkg/scheduler/data_manager/sku" + "yunion.io/x/onecloud/pkg/scheduler/data_manager/zone" schedmodels "yunion.io/x/onecloud/pkg/scheduler/models" ) @@ -48,8 +52,7 @@ type BaseHostDesc struct { IsolatedDevices []*core.IsolatedDeviceDesc `json:"isolated_devices"` - Tenants map[string]int64 `json:"tenants"` - HostSchedtags []computemodels.SSchedtag `json:"schedtags"` + Tenants map[string]int64 `json:"tenants"` InstanceGroups map[string]*api.CandidateGroup `json:"instance_groups"` IpmiInfo types.SIPMIInfo `json:"ipmi_info"` @@ -135,10 +138,6 @@ func (b baseHostGetter) HostType() string { return b.h.HostType } -func (b baseHostGetter) HostSchedtags() []computemodels.SSchedtag { - return b.h.HostSchedtags -} - func (b baseHostGetter) Sku(instanceType string) *sku.ServerSku { zone := b.Zone() return sku.GetByZone(instanceType, zone.GetId()) @@ -360,7 +359,7 @@ func newBaseHostDesc(b *baseBuilder, host *computemodels.SHost, netGetter *netwo SHost: host, } - if err := desc.fillCloudProvider(host); err != nil { + if err := desc.fillCloudProvider(b, host); err != nil { return nil, fmt.Errorf("Fill cloudprovider info error: %v", err) } @@ -395,10 +394,6 @@ func newBaseHostDesc(b *baseBuilder, host *computemodels.SHost, netGetter *netwo return nil, fmt.Errorf("Fill storage error: %v", err) } - if err := desc.fillSchedtags(); err != nil { - return nil, fmt.Errorf("Fill schedtag error: %v", err) - } - if err := desc.fillInstanceGroups(host); err != nil { return nil, fmt.Errorf("Fill instance group error: %v", err) } @@ -569,26 +564,36 @@ func (h *BaseHostDesc) fillIsolatedDevices(b *baseBuilder, host *computemodels.S return nil } -func (b *BaseHostDesc) fillCloudProvider(host *computemodels.SHost) error { - b.Cloudprovider = host.GetCloudprovider() - if b.Cloudprovider != nil { - var err error - b.Cloudaccount, err = b.Cloudprovider.GetCloudaccount() - if err != nil { - return err - } +func (b *BaseHostDesc) fillCloudProvider(builder *baseBuilder, host *computemodels.SHost) error { + provider, ok := builder.hostCloudproviers[host.GetId()] + if !ok { + return nil } + b.Cloudprovider = provider + account, ok := builder.hostCloudaccounts[host.GetId()] + if !ok { + return nil + } + b.Cloudaccount = account return nil } func (b *BaseHostDesc) fillRegion(host *computemodels.SHost) error { - b.Region, _ = host.GetRegion() + regionId := b.Zone.GetCloudRegionId() + obj, ok := cloudregion.Manager.GetResource(regionId) + if !ok { + return errors.Errorf("Not found cloudregion by host %q with id %q", host.GetName(), regionId) + } + b.Region = &obj return nil } func (b *BaseHostDesc) fillZone(host *computemodels.SHost) error { - zone, _ := host.GetZone() - b.Zone = zone + obj, ok := zone.Manager.GetResource(host.ZoneId) + if !ok { + return errors.Errorf("Not found zone by host %q with id %q", host.GetName(), host.ZoneId) + } + b.Zone = &obj b.ZoneId = host.ZoneId return nil } @@ -604,11 +609,6 @@ func (b *BaseHostDesc) fillResidentTenants(host *computemodels.SHost) error { return nil } -func (b *BaseHostDesc) fillSchedtags() error { - b.HostSchedtags = b.SHost.GetSchedtags() - return nil -} - func (b *BaseHostDesc) fillSharedDomains() error { b.SharedDomains = b.SHost.GetSharedDomains() return nil @@ -616,15 +616,19 @@ func (b *BaseHostDesc) fillSharedDomains() error { func (b *BaseHostDesc) fillNetworks(host *computemodels.SHost, netGetter *networkGetter) error { hostId := host.Id - hostwires := computemodels.HostwireManager.Query().SubQuery() - sq := hostwires.Query(sqlchemy.DISTINCT("wire_id", hostwires.Field("wire_id"))).Equals("host_id", hostId) - networks := computemodels.NetworkManager.Query().SubQuery() - q := networks.Query().In("wire_id", sq) + + hostwires := hostwire.GetByHost(hostId) + wireIds := sets.NewString() + for _, hw := range hostwires { + wireIds.Insert(hw.WireId) + } nets := make([]computemodels.SNetwork, 0) - err := computedb.FetchModelObjects(computemodels.NetworkManager, q, &nets) - if err != nil { - return err + allNets := network.Manager.GetStore().GetAll() + for _, net := range allNets { + if wireIds.Has(net.WireId) { + nets = append(nets, net) + } } b.Networks = make([]*api.CandidateNetwork, len(nets)) for idx, n := range nets { @@ -633,26 +637,26 @@ func (b *BaseHostDesc) fillNetworks(host *computemodels.SHost, netGetter *networ return errors.Wrapf(err, "GetFreePort for network %s(%s)", n.GetName(), n.GetId()) } b.Networks[idx] = &api.CandidateNetwork{ - SNetwork: &nets[idx], - Schedtags: n.GetSchedtags(), - FreePort: freePort, + SNetwork: &nets[idx], + FreePort: freePort, } } - netifs := host.GetNetInterfaces() + // netifs := host.GetNetInterfaces() + netifs := netinterface.GetByHost(hostId) netifIndexs := make(map[string][]computemodels.SNetInterface, 0) for _, netif := range netifs { if !netif.IsUsableServernic() { continue } - wire := netif.GetWire() - if wire == nil { + wireId := netif.WireId + if wireId == "" { continue } - if _, exist := netifIndexs[wire.Id]; !exist { - netifIndexs[wire.Id] = make([]computemodels.SNetInterface, 0) + if _, exist := netifIndexs[wireId]; !exist { + netifIndexs[wireId] = make([]computemodels.SNetInterface, 0) } - netifIndexs[wire.Id] = append(netifIndexs[wire.Id], netif) + netifIndexs[wireId] = append(netifIndexs[wireId], netif) } b.NetInterfaces = netifIndexs @@ -751,12 +755,12 @@ func (b *BaseHostDesc) fillOnecloudVpcNetworks(netGetter *networkGetter) error { return nil } -func (b *BaseHostDesc) GetHypervisorDriver() models.IGuestDriver { +func (b *BaseHostDesc) GetHypervisorDriver() computemodels.IGuestDriver { hypervisor := computeapi.HOSTTYPE_HYPERVISOR[b.HostType] if hypervisor == "" { return nil } - return models.GetDriver(hypervisor) + return computemodels.GetDriver(hypervisor) } func (b *BaseHostDesc) fillStorages(host *computemodels.SHost) error { @@ -766,7 +770,6 @@ func (b *BaseHostDesc) fillStorages(host *computemodels.SHost) error { cs := &api.CandidateStorage{ SStorage: storage, ActualFreeCapacity: storage.Capacity - storage.ActualCapacityUsed, - Schedtags: storage.GetSchedtags(), } if b.GetHypervisorDriver() == nil { cs.FreeCapacity = storage.GetFreeCapacity() diff --git a/pkg/scheduler/cache/candidate/builder.go b/pkg/scheduler/cache/candidate/builder.go index 445f8df446..3565e0b8c4 100644 --- a/pkg/scheduler/cache/candidate/builder.go +++ b/pkg/scheduler/cache/candidate/builder.go @@ -15,25 +15,52 @@ package candidate import ( + "fmt" + "strings" + gosync "sync" + "sync/atomic" "time" + "yunion.io/x/log" + "yunion.io/x/pkg/errors" + "yunion.io/x/pkg/util/sets" + "yunion.io/x/pkg/util/workqueue" "yunion.io/x/pkg/utils" "yunion.io/x/sqlchemy" computeapi "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db" computemodels "yunion.io/x/onecloud/pkg/compute/models" + "yunion.io/x/onecloud/pkg/scheduler/data_manager/cloudaccount" + "yunion.io/x/onecloud/pkg/scheduler/data_manager/cloudprovider" + o "yunion.io/x/onecloud/pkg/scheduler/options" ) type baseBuilder struct { resourceType string + builder IResourceBuilder + + hosts []computemodels.SHost + hostDict map[string]*computemodels.SHost isolatedDevicesDict map[string][]interface{} + + hostCloudproviers map[string]*computemodels.SCloudprovider + hostCloudaccounts map[string]*computemodels.SCloudaccount } -func newBaseBuilder(resourceType string) *baseBuilder { +type InitFunc func(ids []computemodels.SHost, errChan chan error) + +type IResourceBuilder interface { + FetchHosts(ids []string) ([]computemodels.SHost, error) + InitFuncs() []InitFunc + BuildOne(host *computemodels.SHost, getter *networkGetter, desc *BaseHostDesc) (interface{}, error) +} + +func newBaseBuilder(resourceType string, builder IResourceBuilder) *baseBuilder { return &baseBuilder{ resourceType: resourceType, + builder: builder, } } @@ -41,6 +68,119 @@ func (b *baseBuilder) Type() string { return b.resourceType } +func (b *baseBuilder) Do(ids []string) ([]interface{}, error) { + err := b.init(ids) + if err != nil { + return nil, err + } + netGetter := newNetworkGetter() + descs, err := b.build(netGetter) + if err != nil { + return nil, err + } + return descs, nil +} + +func (b *baseBuilder) init(ids []string) error { + if err := b.setHosts(ids); err != nil { + return errors.Wrap(err, "set host objects") + } + wg := &WaitGroupWrapper{} + errMessageChannel := make(chan error, 12) + defer close(errMessageChannel) + setFuncs := []func(){ + // func() { b.setHosts(ids, errMessageChannel) }, + func() { + b.setIsolatedDevs(ids, errMessageChannel) + }, + func() { + b.setCloudproviderAccounts(b.hosts, errMessageChannel) + }, + func() { + for _, f := range b.builder.InitFuncs() { + f(b.hosts, errMessageChannel) + } + }, + } + + for _, f := range setFuncs { + wg.Wrap(f) + } + + if ok := waitTimeOut(wg, time.Duration(20*time.Second)); !ok { + log.Errorln("HostBuilder waitgroup timeout.") + } + + if len(errMessageChannel) != 0 { + errMessages := make([]string, 0) + lengthChan := len(errMessageChannel) + for ; lengthChan > 0; lengthChan-- { + msg := fmt.Sprintf("%s", <-errMessageChannel) + log.Errorf("Get error from chan: %s", msg) + errMessages = append(errMessages, msg) + } + return fmt.Errorf("%s\n", strings.Join(errMessages, ";")) + } + + return nil +} + +func (b *baseBuilder) build(netGetter *networkGetter) ([]interface{}, error) { + schedDescs := make([]interface{}, len(b.hosts)) + errs := []error{} + var descResultLock gosync.Mutex + var descedLen int32 + + buildOne := func(i int) { + if i >= len(b.hosts) { + log.Errorf("invalid host index[%d] in b.hosts: %v", i, b.hosts) + return + } + host := b.hosts[i] + desc, err := b.buildOne(&host, netGetter) + if err != nil { + descResultLock.Lock() + errs = append(errs, err) + descResultLock.Unlock() + return + } + descResultLock.Lock() + schedDescs[atomic.AddInt32(&descedLen, 1)-1] = desc + descResultLock.Unlock() + } + + workqueue.Parallelize(o.Options.HostBuildParallelizeSize, len(b.hosts), buildOne) + schedDescs = schedDescs[:descedLen] + if len(errs) > 0 { + //return nil, errors.NewAggregate(errs) + err := errors.NewAggregate(errs) + log.Errorf("Build schedule desc of %s error: %s", b.resourceType, err) + } + + return schedDescs, nil +} + +func (b *baseBuilder) buildOne(host *computemodels.SHost, netGetter *networkGetter) (interface{}, error) { + baseDesc, err := newBaseHostDesc(b, host, netGetter) + if err != nil { + return nil, err + } + + return b.builder.BuildOne(host, netGetter, baseDesc) +} + +func (b *baseBuilder) setHosts(ids []string) error { + hostObjs, err := b.builder.FetchHosts(ids) + if err != nil { + return errors.Wrap(err, "FetchHosts") + } + + hostDict := ToDict(hostObjs) + b.hosts = hostObjs + b.hostDict = hostDict + return nil +} + func (b *baseBuilder) getIsolatedDevices(hostID string) (devs []computemodels.SIsolatedDevice) { devObjs, ok := b.isolatedDevicesDict[hostID] devs = make([]computemodels.SIsolatedDevice, 0) @@ -70,6 +210,58 @@ func (b *baseBuilder) setIsolatedDevs(ids []string, errMessageChannel chan error b.isolatedDevicesDict = dict } +func (b *baseBuilder) setCloudproviderAccounts(hosts []computemodels.SHost, errCh chan error) { + providerSets := sets.NewString() + for _, host := range hosts { + mId := host.ManagerId + if mId != "" { + providerSets.Insert(mId) + } + } + providerObjs := make([]computemodels.SCloudprovider, 0) + for _, pId := range providerSets.List() { + pObj, ok := cloudprovider.Manager.GetResource(pId) + if !ok { + errCh <- errors.Errorf("Not found cloudprovider by id: %q", pId) + return + } + providerObjs = append(providerObjs, pObj) + } + providerDict := ToDict(providerObjs) + + accountSets := sets.NewString() + for _, provider := range providerObjs { + accountSets.Insert(provider.CloudaccountId) + } + accountObjs := make([]computemodels.SCloudaccount, 0) + for _, aId := range accountSets.List() { + aObj, ok := cloudaccount.Manager.GetResource(aId) + if !ok { + errCh <- errors.Errorf("Not found cloudaccount by id: %q", aId) + return + } + accountObjs = append(accountObjs, aObj) + } + accountDict := ToDict(accountObjs) + + b.hostCloudproviers = make(map[string]*computemodels.SCloudprovider, 0) + b.hostCloudaccounts = make(map[string]*computemodels.SCloudaccount, 0) + for _, host := range hosts { + pId := host.ManagerId + provider, ok := providerDict[pId] + if !ok { + continue + } + b.hostCloudproviers[host.GetId()] = provider + aId := provider.CloudaccountId + account, ok := accountDict[aId] + if !ok { + continue + } + b.hostCloudaccounts[host.GetId()] = account + } +} + func FetchModelIds(q *sqlchemy.SQuery) ([]string, error) { rs, err := q.Rows() if err != nil { diff --git a/pkg/scheduler/cache/candidate/common.go b/pkg/scheduler/cache/candidate/common.go index f4e8bbf939..4c9badd31a 100644 --- a/pkg/scheduler/cache/candidate/common.go +++ b/pkg/scheduler/cache/candidate/common.go @@ -15,12 +15,11 @@ package candidate import ( - //"yunion.io/x/log" - "yunion.io/x/pkg/util/sets" computeapi "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" "yunion.io/x/onecloud/pkg/compute/models" ) @@ -46,6 +45,14 @@ var ( ) ) +func GetHostIds(hosts []models.SHost) []string { + ids := make([]string, len(hosts)) + for i, host := range hosts { + ids[i] = host.GetId() + } + return ids +} + func FetchGuestByHostIDs(ids []string) ([]models.SGuest, error) { gs := make([]models.SGuest, 0) q := models.GuestManager.Query().In("host_id", ids) @@ -64,3 +71,13 @@ func IsGuestCreating(g models.SGuest) bool { func IsGuestPendingDelete(g models.SGuest) bool { return g.PendingDeleted } + +func ToDict[O lockman.ILockedObject](objs []O) map[string]*O { + ret := make(map[string]*O, 0) + for _, obj := range objs { + tmpObj := obj + objPtr := &tmpObj + ret[obj.GetId()] = objPtr + } + return ret +} diff --git a/pkg/scheduler/cache/candidate/hosts.go b/pkg/scheduler/cache/candidate/hosts.go index 41f6a96bf8..b455c85452 100644 --- a/pkg/scheduler/cache/candidate/hosts.go +++ b/pkg/scheduler/cache/candidate/hosts.go @@ -16,16 +16,12 @@ package candidate import ( "encoding/json" - "fmt" - "strings" gosync "sync" - "sync/atomic" "time" "yunion.io/x/log" "yunion.io/x/pkg/errors" "yunion.io/x/pkg/util/sets" - "yunion.io/x/pkg/util/workqueue" "yunion.io/x/pkg/utils" "yunion.io/x/sqlchemy" @@ -211,51 +207,6 @@ func NewGuestReservedResourceUsedByBuilder(b *HostBuilder, host *computemodels.S return } -type HostBuilder struct { - *baseBuilder - - residentTenantDict map[string]map[string]interface{} - - hosts []computemodels.SHost - hostDict map[string]interface{} - - guests []computemodels.SGuest - guestDict map[string]interface{} - guestIDs []string - - hostStorages []computemodels.SHoststorage - //hostStoragesDict map[string][]*computemodels.SStorage - storages []interface{} - storageStatesSizeDict map[string]map[string]interface{} - - hostGuests map[string][]interface{} - hostBackupGuests map[string][]interface{} - - //groupGuests []interface{} - //groups []interface{} - //groupDict map[string]interface{} - //hostGroupCountDict HostGroupCountDict - - //hostMetadatas []interface{} - //hostMetadatasDict map[string][]interface{} - //guestMetadatas []interface{} - //guestMetadatasDict map[string][]interface{} - - //diskStats []models.StorageCapacity - // isolatedDevicesDict map[string][]interface{} - - cpuIOLoads map[string]map[string]float64 - - schedtags []computemodels.SSchedtag - zoneSkus map[string][]computemodels.SServerSku -} - -func newHostBuilder() *HostBuilder { - return &HostBuilder{ - baseBuilder: newBaseBuilder(HostDescBuilder), - } -} - func (h *HostDesc) String() string { s, _ := json.Marshal(h) return string(s) @@ -454,77 +405,55 @@ func waitTimeOut(wg *WaitGroupWrapper, timeout time.Duration) bool { } } -func (b *HostBuilder) init(ids []string) error { - wg := &WaitGroupWrapper{} - errMessageChannel := make(chan error, 12) - defer close(errMessageChannel) - setFuncs := []func(){ - func() { b.setHosts(ids, errMessageChannel) }, - func() { b.setSchedtags(ids, errMessageChannel) }, - func() { - b.setGuests(ids, errMessageChannel) - b.setIsolatedDevs(ids, errMessageChannel) - }, - } +type HostBuilder struct { + *baseBuilder - for _, f := range setFuncs { - wg.Wrap(f) - } + residentTenantDict map[string]map[string]interface{} - if ok := waitTimeOut(wg, time.Duration(20*time.Second)); !ok { - log.Errorln("HostBuilder waitgroup timeout.") - } + guests []computemodels.SGuest + guestDict map[string]interface{} + guestIDs []string - if len(errMessageChannel) != 0 { - errMessages := make([]string, 0) - lengthChan := len(errMessageChannel) - for ; lengthChan > 0; lengthChan-- { - msg := fmt.Sprintf("%s", <-errMessageChannel) - log.Errorf("Get error from chan: %s", msg) - errMessages = append(errMessages, msg) - } - return fmt.Errorf("%s\n", strings.Join(errMessages, ";")) - } + hostStorages []computemodels.SHoststorage + //hostStoragesDict map[string][]*computemodels.SStorage + storages []interface{} + storageStatesSizeDict map[string]map[string]interface{} - return nil + hostGuests map[string][]interface{} + hostBackupGuests map[string][]interface{} + + //groupGuests []interface{} + //groups []interface{} + //groupDict map[string]interface{} + //hostGroupCountDict HostGroupCountDict + + //hostMetadatas []interface{} + //hostMetadatasDict map[string][]interface{} + //guestMetadatas []interface{} + //guestMetadatasDict map[string][]interface{} + + //diskStats []models.StorageCapacity + // isolatedDevicesDict map[string][]interface{} + + cpuIOLoads map[string]map[string]float64 } -func (b *HostBuilder) setHosts(ids []string, errMessageChannel chan error) { +func newHostBuilder() *HostBuilder { + builder := new(HostBuilder) + builder.baseBuilder = newBaseBuilder(HostDescBuilder, builder) + return builder +} + +func (b *HostBuilder) FetchHosts(ids []string) ([]computemodels.SHost, error) { hosts := computemodels.HostManager.Query() q := hosts.In("id", ids).NotEquals("host_type", computeapi.HOST_TYPE_BAREMETAL) hostObjs := make([]computemodels.SHost, 0) err := computedb.FetchModelObjects(computemodels.HostManager, q, &hostObjs) - if err != nil { - errMessageChannel <- err - return - } - - hostDict, err := utils.ToDict(hostObjs, func(obj interface{}) (string, error) { - host, ok := obj.(computemodels.SHost) - if !ok { - return "", utils.ConvertError(obj, "computemodels.Host") - } - return host.Id, nil - }) - if err != nil { - errMessageChannel <- err - return - } - b.hosts = hostObjs - b.hostDict = hostDict - return + return hostObjs, err } -func (b *HostBuilder) setSchedtags(ids []string, errMessageChannel chan error) { - tags := make([]computemodels.SSchedtag, 0) - if err := computemodels.SchedtagManager.Query().All(&tags); err != nil { - errMessageChannel <- err - return - } - b.schedtags = tags -} - -func (b *HostBuilder) setGuests(ids []string, errMessageChannel chan error) { +func (b *HostBuilder) setGuests(hosts []computemodels.SHost, errMessageChannel chan error) { + ids := GetHostIds(hosts) guests, err := FetchGuestByHostIDs(ids) if err != nil { errMessageChannel <- err @@ -743,9 +672,7 @@ func (b *HostBuilder) setGuests(ids []string, errMessageChannel chan error) { }*/ func (b *HostBuilder) Clone() BuildActor { - return &HostBuilder{ - baseBuilder: newBaseBuilder(HostDescBuilder), - } + return newHostBuilder() } func (b *HostBuilder) AllIDs() ([]string, error) { @@ -754,59 +681,14 @@ func (b *HostBuilder) AllIDs() ([]string, error) { return FetchModelIds(q) } -func (b *HostBuilder) Do(ids []string) ([]interface{}, error) { - err := b.init(ids) - if err != nil { - return nil, err +func (b *HostBuilder) InitFuncs() []InitFunc { + return []InitFunc{ + // b.setSchedtags, + b.setGuests, } - netGetter := newNetworkGetter() - descs, err := b.build(netGetter) - if err != nil { - return nil, err - } - return descs, nil } -func (b *HostBuilder) build(netGetter *networkGetter) ([]interface{}, error) { - schedDescs := make([]interface{}, len(b.hosts)) - errs := []error{} - var descResultLock gosync.Mutex - var descedLen int32 - - buildOne := func(i int) { - if i >= len(b.hosts) { - log.Errorf("invalid host index[%d] in b.hosts: %v", i, b.hosts) - return - } - host := b.hosts[i] - desc, err := b.buildOne(&host, netGetter) - if err != nil { - descResultLock.Lock() - errs = append(errs, err) - descResultLock.Unlock() - return - } - descResultLock.Lock() - schedDescs[atomic.AddInt32(&descedLen, 1)-1] = desc - descResultLock.Unlock() - } - - workqueue.Parallelize(o.Options.HostBuildParallelizeSize, len(b.hosts), buildOne) - schedDescs = schedDescs[:descedLen] - if len(errs) > 0 { - //return nil, errors.NewAggregate(errs) - err := errors.NewAggregate(errs) - log.V(4).Warningf("Build schedule descs error: %s", err) - } - - return schedDescs, nil -} - -func (b *HostBuilder) buildOne(host *computemodels.SHost, netGetter *networkGetter) (interface{}, error) { - baseDesc, err := newBaseHostDesc(b.baseBuilder, host, netGetter) - if err != nil { - return nil, err - } +func (b *HostBuilder) BuildOne(host *computemodels.SHost, getter *networkGetter, baseDesc *BaseHostDesc) (interface{}, error) { desc := &HostDesc{ BaseHostDesc: baseDesc, } diff --git a/pkg/scheduler/core/types.go b/pkg/scheduler/core/types.go index 3642d90007..053077339c 100644 --- a/pkg/scheduler/core/types.go +++ b/pkg/scheduler/core/types.go @@ -75,7 +75,6 @@ type CandidatePropertyGetter interface { SharedDomains() []string Region() *computemodels.SCloudregion HostType() string - HostSchedtags() []computemodels.SSchedtag Sku(string) *sku.ServerSku Storages() []*api.CandidateStorage Networks() []*api.CandidateNetwork diff --git a/pkg/scheduler/data_manager/cloudaccount/cloudaccount.go b/pkg/scheduler/data_manager/cloudaccount/cloudaccount.go new file mode 100644 index 0000000000..f0a6e4f2ae --- /dev/null +++ b/pkg/scheduler/data_manager/cloudaccount/cloudaccount.go @@ -0,0 +1,31 @@ +package cloudaccount + +import ( + "time" + + "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.SCloudaccount] + +func init() { + Manager = NewResourceManager() +} + +func NewResourceManager() common.IResourceManager[models.SCloudaccount] { + cm := common.NewCommonResourceManager( + "cloudaccount", + 15*time.Minute, + NewResourceStore(), + ) + return cm +} + +func NewResourceStore() common.IResourceStore[models.SCloudaccount] { + return common.NewResourceStore[models.SCloudaccount]( + models.CloudaccountManager, + compute.Cloudaccounts, + ) +} diff --git a/pkg/scheduler/data_manager/cloudaccount/doc.go b/pkg/scheduler/data_manager/cloudaccount/doc.go new file mode 100644 index 0000000000..ffa8edf420 --- /dev/null +++ b/pkg/scheduler/data_manager/cloudaccount/doc.go @@ -0,0 +1 @@ +package cloudaccount // import "yunion.io/x/onecloud/pkg/scheduler/data_manager/cloudaccount" diff --git a/pkg/scheduler/data_manager/cloudprovider/cloudprovider.go b/pkg/scheduler/data_manager/cloudprovider/cloudprovider.go new file mode 100644 index 0000000000..dc2f225c44 --- /dev/null +++ b/pkg/scheduler/data_manager/cloudprovider/cloudprovider.go @@ -0,0 +1,31 @@ +package cloudprovider + +import ( + "time" + + "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] + +func init() { + Manager = NewResourceManager() +} + +func NewResourceManager() common.IResourceManager[models.SCloudprovider] { + cm := common.NewCommonResourceManager( + "cloudprovider", + 15*time.Minute, + NewResourceStore(), + ) + return cm +} + +func NewResourceStore() common.IResourceStore[models.SCloudprovider] { + return common.NewResourceStore[models.SCloudprovider]( + models.CloudproviderManager, + compute.Cloudproviders, + ) +} diff --git a/pkg/scheduler/data_manager/cloudprovider/doc.go b/pkg/scheduler/data_manager/cloudprovider/doc.go new file mode 100644 index 0000000000..37ba75ceb8 --- /dev/null +++ b/pkg/scheduler/data_manager/cloudprovider/doc.go @@ -0,0 +1 @@ +package cloudprovider // import "yunion.io/x/onecloud/pkg/scheduler/data_manager/cloudprovider" diff --git a/pkg/scheduler/data_manager/cloudregion/cloudregion.go b/pkg/scheduler/data_manager/cloudregion/cloudregion.go new file mode 100644 index 0000000000..9c82098c75 --- /dev/null +++ b/pkg/scheduler/data_manager/cloudregion/cloudregion.go @@ -0,0 +1,31 @@ +package cloudregion + +import ( + "time" + + "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.SCloudregion] + +func init() { + Manager = NewResourceManager() +} + +func NewResourceManager() common.IResourceManager[models.SCloudregion] { + cm := common.NewCommonResourceManager( + "cloudregion", + 15*time.Minute, + NewResourceStore(), + ) + return cm +} + +func NewResourceStore() common.IResourceStore[models.SCloudregion] { + return common.NewResourceStore[models.SCloudregion]( + models.CloudregionManager, + compute.Cloudregions, + ) +} diff --git a/pkg/scheduler/data_manager/cloudregion/doc.go b/pkg/scheduler/data_manager/cloudregion/doc.go new file mode 100644 index 0000000000..bebb635431 --- /dev/null +++ b/pkg/scheduler/data_manager/cloudregion/doc.go @@ -0,0 +1 @@ +package cloudregion // import "yunion.io/x/onecloud/pkg/scheduler/data_manager/cloudregion" diff --git a/pkg/scheduler/data_manager/common/common.go b/pkg/scheduler/data_manager/common/common.go new file mode 100644 index 0000000000..ac8f59d8dd --- /dev/null +++ b/pkg/scheduler/data_manager/common/common.go @@ -0,0 +1,310 @@ +package common + +import ( + "context" + "reflect" + "strings" + "sync" + "time" + + "yunion.io/x/jsonutils" + "yunion.io/x/log" + "yunion.io/x/pkg/errors" + "yunion.io/x/pkg/util/wait" + + "yunion.io/x/onecloud/pkg/cloudcommon/consts" + "yunion.io/x/onecloud/pkg/cloudcommon/db" + "yunion.io/x/onecloud/pkg/cloudcommon/db/lockman" + "yunion.io/x/onecloud/pkg/mcclient/auth" + "yunion.io/x/onecloud/pkg/mcclient/informer" +) + +type IResourceManager[O lockman.ILockedObject] interface { + GetKeyword() string + + GetRefreshInterval() time.Duration + GetStore() IResourceStore[O] + GetResource(id string) (O, bool) + + SyncOnce() error + + Start(ctx context.Context) +} + +type IResourceStore[O lockman.ILockedObject] interface { + GetInformerResourceManager() informer.IResourceManager + + Init() error + Get(id string) (O, bool) + GetAll() []O + GetByPrefix(prefixId string) []O + Add(obj *jsonutils.JSONDict) + Update(oldObj, newObj *jsonutils.JSONDict) + Delete(obj *jsonutils.JSONDict) +} + +type CommonResourceManager[O lockman.ILockedObject] struct { + keyword string + refreshInterval time.Duration + store IResourceStore[O] +} + +func NewCommonResourceManager[O lockman.ILockedObject]( + keyword string, + refreshInterval time.Duration, + store IResourceStore[O], +) *CommonResourceManager[O] { + return &CommonResourceManager[O]{ + keyword: keyword, + refreshInterval: refreshInterval, + store: store, + } +} + +func (m *CommonResourceManager[O]) Start(ctx context.Context) { + go func() { + Start[O](ctx, m) + }() +} + +func (m *CommonResourceManager[O]) GetKeyword() string { + return m.keyword +} + +func (m CommonResourceManager[O]) GetStore() IResourceStore[O] { + return m.store +} + +func (m CommonResourceManager[O]) GetResource(id string) (O, bool) { + return m.store.Get(id) +} + +func (m CommonResourceManager[O]) GetAll() []O { + return m.store.GetAll() +} + +func (m *CommonResourceManager[O]) GetRefreshInterval() time.Duration { + return m.refreshInterval +} + +func (m *CommonResourceManager[O]) SyncOnce() error { + return m.GetStore().Init() +} + +type ResourceStore[O lockman.ILockedObject] struct { + dataMap *sync.Map + modelMan db.IModelManager + res informer.IResourceManager + getId func(O) string + getWatchId func(*jsonutils.JSONDict) string +} + +func NewResourceStore[O lockman.ILockedObject]( + modelMan db.IModelManager, + res informer.IResourceManager, +) IResourceStore[O] { + return newResourceStore[O](modelMan, res, nil, nil) +} + +func NewJointResourceStore[O lockman.ILockedObject]( + modelMan db.IModelManager, + res informer.IResourceManager, + getId func(O) string, + getWatchId func(*jsonutils.JSONDict) string, +) IResourceStore[O] { + return newResourceStore(modelMan, res, getId, getWatchId) +} + +func newResourceStore[O lockman.ILockedObject]( + modelMan db.IModelManager, + res informer.IResourceManager, + getId func(O) string, + getWatchId func(*jsonutils.JSONDict) string, +) IResourceStore[O] { + if getId == nil { + getId = func(o O) string { + return o.GetId() + } + } + if getWatchId == nil { + getWatchId = func(o *jsonutils.JSONDict) string { + id, _ := o.GetString("id") + return id + } + } + return &ResourceStore[O]{ + dataMap: new(sync.Map), + modelMan: modelMan, + res: res, + getId: getId, + getWatchId: getWatchId, + } +} + +func (s *ResourceStore[O]) GetInformerResourceManager() informer.IResourceManager { + return s.res +} + +func (s *ResourceStore[O]) Init() error { + objs := make([]O, 0) + q := s.modelMan.Query() + if err := db.FetchModelObjects(s.modelMan, q, &objs); err != nil { + return err + } + for _, obj := range objs { + s.dataMap.Store(s.getId(obj), obj) + } + return nil +} + +func (s *ResourceStore[O]) Get(id string) (O, bool) { + obj, ok := s.dataMap.Load(id) + if !ok { + var ret O + return ret, false + } + return obj.(O), true +} + +func (s *ResourceStore[O]) GetByPrefix(prefixId string) []O { + ret := make([]O, 0) + s.dataMap.Range(func(key, value any) bool { + if strings.HasPrefix(key.(string), prefixId) { + ret = append(ret, value.(O)) + } + return true + }) + return ret +} + +func (s *ResourceStore[O]) GetAll() []O { + ret := make([]O, 0) + s.dataMap.Range(func(key, value any) bool { + ret = append(ret, value.(O)) + return true + }) + return ret +} + +func (s *ResourceStore[O]) Add(obj *jsonutils.JSONDict) { + id := s.getWatchId(obj) + if id != "" { + dbObj, err := s.modelMan.FetchById(id) + if err == nil { + v := reflect.ValueOf(dbObj) + tmpObj := v.Elem().Interface() + s.dataMap.Store(id, tmpObj) + log.Infof("Add %s %s", s.modelMan.Keyword(), obj.String()) + } else { + log.Errorf("Fetch %s by id %s error when created: %v", s.modelMan.Keyword(), id, err) + } + } +} + +func (s *ResourceStore[O]) removeIgnoreKeys(obj *jsonutils.JSONDict) *jsonutils.JSONDict { + // ignore keys updated by cloudaccount + for _, key := range []string{ + "probe_at", + "update_version", + "updated_at", + } { + obj.Remove(key) + } + return obj +} + +func (s *ResourceStore[O]) Update(oldObj, newObj *jsonutils.JSONDict) { + id := s.getWatchId(newObj) + oldObj = s.removeIgnoreKeys(oldObj) + newObj = s.removeIgnoreKeys(newObj) + isEq := oldObj.String() == newObj.String() + if id != "" && !isEq { + dbObj, err := s.modelMan.FetchById(id) + if err == nil { + v := reflect.ValueOf(dbObj) + tmpObj := v.Elem().Interface() + s.dataMap.Store(id, tmpObj) + log.Infof("Update %s %s", s.modelMan.Keyword(), newObj.String()) + } else { + log.Errorf("Fetch %s by id %s error when updated: %v", s.modelMan.Keyword(), id, err) + } + } +} + +func (s *ResourceStore[O]) Delete(obj *jsonutils.JSONDict) { + id := s.getWatchId(obj) + if id != "" { + s.dataMap.Delete(id) + log.Infof("Delete %s %s", s.modelMan.Keyword(), obj.String()) + } +} + +func Start[O lockman.ILockedObject](ctx context.Context, resMan IResourceManager[O]) { + startWatch(ctx, resMan) + startSync(resMan) +} + +func startWatch[O lockman.ILockedObject](ctx context.Context, resMan IResourceManager[O]) { + s := auth.GetAdminSession(ctx, consts.GetRegion()) + informer.NewWatchManagerBySessionBg(s, func(man *informer.SWatchManager) error { + res := resMan.GetStore().GetInformerResourceManager() + if err := man.For(res).AddEventHandler(ctx, newEventHandler(res, resMan)); err != nil { + return errors.Wrapf(err, "watch resource %s", res.KeyString()) + } + return nil + }) +} + +func startSync[O lockman.ILockedObject](resMan IResourceManager[O]) { + wait.Forever(func() { + log.Infof("%s data start sync", resMan.GetKeyword()) + startTime := time.Now() + if err := syncOnce(resMan); err != nil { + log.Errorf("%s sync data error: %v", resMan.GetKeyword(), err) + return + } + log.Infof("%s finish sync, elapsed %s", resMan.GetKeyword(), time.Since(startTime)) + }, resMan.GetRefreshInterval()) +} + +func syncOnce[O lockman.ILockedObject](resMan IResourceManager[O]) error { + if err := resMan.SyncOnce(); err != nil { + return errors.Wrapf(err, "sync once of %s", resMan.GetKeyword()) + } + return nil +} + +type eventHandler[O lockman.ILockedObject] struct { + resMan informer.IResourceManager + dataMan IResourceManager[O] +} + +func newEventHandler[O lockman.ILockedObject](resMan informer.IResourceManager, dataMan IResourceManager[O]) informer.EventHandler { + return &eventHandler[O]{ + resMan: resMan, + dataMan: dataMan, + } +} + +func (e eventHandler[O]) keyword() string { + return e.resMan.GetKeyword() +} + +func (e eventHandler[O]) store() IResourceStore[O] { + return e.dataMan.GetStore() +} + +func (e eventHandler[O]) OnAdd(obj *jsonutils.JSONDict) { + log.Debugf("%s [CREATED]: \n%s", e.keyword(), obj.String()) + e.store().Add(obj) +} + +func (e eventHandler[O]) OnUpdate(oldObj, newObj *jsonutils.JSONDict) { + log.Debugf("%s [UPDATED]: \n[NEW]: %s\n[OLD]: %s", e.keyword(), newObj.String(), oldObj.String()) + e.store().Update(oldObj, newObj) +} + +func (e eventHandler[O]) OnDelete(obj *jsonutils.JSONDict) { + log.Debugf("%s [DELETED]: \n%s", e.keyword(), obj.String()) + e.store().Delete(obj) +} diff --git a/pkg/scheduler/data_manager/common/doc.go b/pkg/scheduler/data_manager/common/doc.go new file mode 100644 index 0000000000..12c7395e4b --- /dev/null +++ b/pkg/scheduler/data_manager/common/doc.go @@ -0,0 +1 @@ +package common // import "yunion.io/x/onecloud/pkg/scheduler/data_manager/common" diff --git a/pkg/scheduler/data_manager/hostwire/doc.go b/pkg/scheduler/data_manager/hostwire/doc.go new file mode 100644 index 0000000000..5b80d0d5ad --- /dev/null +++ b/pkg/scheduler/data_manager/hostwire/doc.go @@ -0,0 +1 @@ +package hostwire // import "yunion.io/x/onecloud/pkg/scheduler/data_manager/hostwire" diff --git a/pkg/scheduler/data_manager/hostwire/hostwire.go b/pkg/scheduler/data_manager/hostwire/hostwire.go new file mode 100644 index 0000000000..221d5da343 --- /dev/null +++ b/pkg/scheduler/data_manager/hostwire/hostwire.go @@ -0,0 +1,54 @@ +package hostwire + +import ( + "fmt" + "time" + + "yunion.io/x/jsonutils" + + "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.SHostwire] + +func GetManager() common.IResourceManager[models.SHostwire] { + if manager != nil { + return manager + } + manager = NewResourceManager() + return manager +} + +func NewResourceManager() common.IResourceManager[models.SHostwire] { + cm := common.NewCommonResourceManager( + "hostwire", + 10*time.Minute, + NewResourceStore(), + ) + return cm +} + +func GetId(hostId, wireId string) string { + return fmt.Sprintf("%s/%s", hostId, wireId) +} + +func NewResourceStore() common.IResourceStore[models.SHostwire] { + return common.NewJointResourceStore( + models.HostwireManager, + compute.Hostwires, + func(o models.SHostwire) string { + return GetId(o.HostId, o.WireId) + }, + func(o *jsonutils.JSONDict) string { + hostId, _ := o.GetString("host_id") + wireId, _ := o.GetString("wire_id") + return GetId(hostId, wireId) + }, + ) +} + +func GetByHost(hostId string) []models.SHostwire { + return GetManager().GetStore().GetByPrefix(hostId) +} diff --git a/pkg/scheduler/data_manager/netinterface/doc.go b/pkg/scheduler/data_manager/netinterface/doc.go new file mode 100644 index 0000000000..e726156219 --- /dev/null +++ b/pkg/scheduler/data_manager/netinterface/doc.go @@ -0,0 +1 @@ +package netinterface // import "yunion.io/x/onecloud/pkg/scheduler/data_manager/netinterface" diff --git a/pkg/scheduler/data_manager/netinterface/netinterface.go b/pkg/scheduler/data_manager/netinterface/netinterface.go new file mode 100644 index 0000000000..4289196880 --- /dev/null +++ b/pkg/scheduler/data_manager/netinterface/netinterface.go @@ -0,0 +1,55 @@ +package netinterface + +import ( + "fmt" + "time" + + "yunion.io/x/jsonutils" + + "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.SNetInterface] + +func GetManager() common.IResourceManager[models.SNetInterface] { + if manager != nil { + return manager + } + manager = NewResourceManager() + return manager +} + +func NewResourceManager() common.IResourceManager[models.SNetInterface] { + cm := common.NewCommonResourceManager( + "netinterface", + 10*time.Minute, + NewResourceStore(), + ) + return cm +} + +func GetId(hostId, wireId, mac string) string { + return fmt.Sprintf("%s/%s/%s", hostId, wireId, mac) +} + +func NewResourceStore() common.IResourceStore[models.SNetInterface] { + return common.NewJointResourceStore( + models.NetInterfaceManager, + compute.Networkinterfaces, + func(o models.SNetInterface) string { + return GetId(o.BaremetalId, o.WireId, o.Mac) + }, + func(o *jsonutils.JSONDict) string { + hostId, _ := o.GetString("host_id") + wireId, _ := o.GetString("wire_id") + mac, _ := o.GetString("mac") + return GetId(hostId, wireId, mac) + }, + ) +} + +func GetByHost(hostId string) []models.SNetInterface { + return GetManager().GetStore().GetByPrefix(hostId) +} diff --git a/pkg/scheduler/data_manager/network/doc.go b/pkg/scheduler/data_manager/network/doc.go new file mode 100644 index 0000000000..437ae7530c --- /dev/null +++ b/pkg/scheduler/data_manager/network/doc.go @@ -0,0 +1 @@ +package network // import "yunion.io/x/onecloud/pkg/scheduler/data_manager/network" diff --git a/pkg/scheduler/data_manager/network/network.go b/pkg/scheduler/data_manager/network/network.go new file mode 100644 index 0000000000..532ef2e66b --- /dev/null +++ b/pkg/scheduler/data_manager/network/network.go @@ -0,0 +1,31 @@ +package network + +import ( + "time" + + "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.SNetwork] + +func init() { + Manager = NewResourceManager() +} + +func NewResourceManager() common.IResourceManager[models.SNetwork] { + cm := common.NewCommonResourceManager( + "network", + 10*time.Minute, + NewResourceStore(), + ) + return cm +} + +func NewResourceStore() common.IResourceStore[models.SNetwork] { + return common.NewResourceStore[models.SNetwork]( + models.NetworkManager, + compute.Networks, + ) +} diff --git a/pkg/scheduler/data_manager/wire/doc.go b/pkg/scheduler/data_manager/wire/doc.go new file mode 100644 index 0000000000..8861057af5 --- /dev/null +++ b/pkg/scheduler/data_manager/wire/doc.go @@ -0,0 +1 @@ +package wire // import "yunion.io/x/onecloud/pkg/scheduler/data_manager/wire" diff --git a/pkg/scheduler/data_manager/wire/wire.go b/pkg/scheduler/data_manager/wire/wire.go new file mode 100644 index 0000000000..90043544c7 --- /dev/null +++ b/pkg/scheduler/data_manager/wire/wire.go @@ -0,0 +1,31 @@ +package wire + +import ( + "time" + + "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.SWire] + +func init() { + Manager = NewResourceManager() +} + +func NewResourceManager() common.IResourceManager[models.SWire] { + cm := common.NewCommonResourceManager( + "wire", + 10*time.Minute, + NewResourceStore(), + ) + return cm +} + +func NewResourceStore() common.IResourceStore[models.SWire] { + return common.NewResourceStore[models.SWire]( + models.WireManager, + compute.Wires, + ) +} diff --git a/pkg/scheduler/data_manager/zone/doc.go b/pkg/scheduler/data_manager/zone/doc.go new file mode 100644 index 0000000000..134545a390 --- /dev/null +++ b/pkg/scheduler/data_manager/zone/doc.go @@ -0,0 +1 @@ +package zone // import "yunion.io/x/onecloud/pkg/scheduler/data_manager/zone" diff --git a/pkg/scheduler/data_manager/zone/zone.go b/pkg/scheduler/data_manager/zone/zone.go new file mode 100644 index 0000000000..f15db869c7 --- /dev/null +++ b/pkg/scheduler/data_manager/zone/zone.go @@ -0,0 +1,31 @@ +package zone + +import ( + "time" + + "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.SZone] + +func init() { + Manager = NewResourceManager() +} + +func NewResourceManager() common.IResourceManager[models.SZone] { + cm := common.NewCommonResourceManager( + "zone", + 15*time.Minute, + NewResourceStore(), + ) + return cm +} + +func NewResourceStore() common.IResourceStore[models.SZone] { + return common.NewResourceStore[models.SZone]( + models.ZoneManager, + compute.Zones, + ) +} diff --git a/pkg/scheduler/service/service.go b/pkg/scheduler/service/service.go index 5d80d854c5..d85dfca2a4 100644 --- a/pkg/scheduler/service/service.go +++ b/pkg/scheduler/service/service.go @@ -38,8 +38,16 @@ import ( _ "yunion.io/x/onecloud/pkg/compute/hostdrivers" computemodels "yunion.io/x/onecloud/pkg/compute/models" _ "yunion.io/x/onecloud/pkg/scheduler/algorithmprovider" + "yunion.io/x/onecloud/pkg/scheduler/data_manager/cloudaccount" + "yunion.io/x/onecloud/pkg/scheduler/data_manager/cloudprovider" + "yunion.io/x/onecloud/pkg/scheduler/data_manager/cloudregion" + "yunion.io/x/onecloud/pkg/scheduler/data_manager/hostwire" + "yunion.io/x/onecloud/pkg/scheduler/data_manager/netinterface" + "yunion.io/x/onecloud/pkg/scheduler/data_manager/network" "yunion.io/x/onecloud/pkg/scheduler/data_manager/schedtag" skuman "yunion.io/x/onecloud/pkg/scheduler/data_manager/sku" + "yunion.io/x/onecloud/pkg/scheduler/data_manager/wire" + "yunion.io/x/onecloud/pkg/scheduler/data_manager/zone" schedhandler "yunion.io/x/onecloud/pkg/scheduler/handler" schedman "yunion.io/x/onecloud/pkg/scheduler/manager" o "yunion.io/x/onecloud/pkg/scheduler/options" @@ -95,8 +103,25 @@ func StartService() error { startSched := func() { stopEverything := make(chan struct{}) + ctx := context.Background() go skuman.Start(utils.ToDuration(o.Options.SkuRefreshInterval)) - go schedtag.Start(context.Background(), utils.ToDuration("30s")) + go schedtag.Start(ctx, utils.ToDuration("30s")) + + for _, f := range []func(ctx context.Context){ + cloudregion.Manager.Start, + zone.Manager.Start, + cloudprovider.Manager.Start, + cloudaccount.Manager.Start, + wire.Manager.Start, + network.Manager.Start, + hostwire.GetManager().Start, + netinterface.GetManager().Start, + } { + f(ctx) + } + + time.Sleep(5 * time.Second) + schedman.InitAndStart(stopEverything) } startSched()