support network schedtag

This commit is contained in:
Zexi Li
2019-04-18 09:48:18 +08:00
committed by Zexi
parent ccb3bcbd78
commit e01c4e590f
27 changed files with 824 additions and 330 deletions
@@ -16,83 +16,18 @@ package predicates
import (
"fmt"
"sort"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/util/errors"
"yunion.io/x/pkg/utils"
computeapi "yunion.io/x/onecloud/pkg/apis/compute"
schedapi "yunion.io/x/onecloud/pkg/apis/scheduler"
"yunion.io/x/onecloud/pkg/compute/models"
"yunion.io/x/onecloud/pkg/scheduler/algorithm/plugin"
"yunion.io/x/onecloud/pkg/scheduler/api"
"yunion.io/x/onecloud/pkg/scheduler/core"
)
type DiskStoragesMap map[int][]*PredicatedStorage
func (m DiskStoragesMap) getAllTags(isPrefer bool) []computeapi.SchedtagConfig {
ret := make([]computeapi.SchedtagConfig, 0)
for _, ss := range m {
for _, s := range ss {
var tags []computeapi.SchedtagConfig
if isPrefer {
tags = s.PreferTags
} else {
tags = s.AvoidTags
}
ret = append(ret, tags...)
}
}
return ret
}
func (m DiskStoragesMap) GetPreferTags() []computeapi.SchedtagConfig {
return m.getAllTags(true)
}
func (m DiskStoragesMap) GetAvoidTags() []computeapi.SchedtagConfig {
return m.getAllTags(false)
}
type CandidateDiskStoragesMap map[string]DiskStoragesMap
type PredicatedStorage struct {
*api.CandidateStorage
PreferTags []computeapi.SchedtagConfig
AvoidTags []computeapi.SchedtagConfig
}
func newPredicatedStorage(s *api.CandidateStorage, preferTags, avoidTags []computeapi.SchedtagConfig) *PredicatedStorage {
return &PredicatedStorage{
CandidateStorage: s,
PreferTags: preferTags,
AvoidTags: avoidTags,
}
}
func (s *PredicatedStorage) isNoTag() bool {
return len(s.PreferTags) == 0 && len(s.AvoidTags) == 0
}
func (s *PredicatedStorage) hasPreferTags() bool {
return len(s.PreferTags) != 0
}
func (s *PredicatedStorage) hasAvoidTags() bool {
return len(s.AvoidTags) != 0
}
type DiskSchedtagPredicate struct {
BasePredicate
plugin.BasePlugin
SchedtagPredicate *SchedtagPredicate
CandidateDiskStoragesMap CandidateDiskStoragesMap
Hypervisor string
*BaseSchedtagPredicate
}
func (p *DiskSchedtagPredicate) Name() string {
@@ -101,210 +36,109 @@ func (p *DiskSchedtagPredicate) Name() string {
func (p *DiskSchedtagPredicate) Clone() core.FitPredicate {
return &DiskSchedtagPredicate{
CandidateDiskStoragesMap: make(map[string]DiskStoragesMap),
BaseSchedtagPredicate: NewBaseSchedtagPredicate(),
}
}
func (p *DiskSchedtagPredicate) getSchedtagDisks(disks []*computeapi.DiskConfig) ([]*computeapi.DiskConfig, []*computeapi.DiskConfig) {
noTagDisk := make([]*computeapi.DiskConfig, 0)
tagDisk := make([]*computeapi.DiskConfig, 0)
for _, d := range disks {
if len(d.Schedtags) != 0 {
tagDisk = append(tagDisk, d)
} else {
noTagDisk = append(noTagDisk, d)
}
}
return noTagDisk, tagDisk
}
func (p *DiskSchedtagPredicate) PreExecute(u *core.Unit, cs []core.Candidater) (bool, error) {
disks := u.SchedData().Disks
if len(disks) == 0 {
return false, nil
return p.BaseSchedtagPredicate.PreExecute(p, u, cs)
}
type diskW struct {
*computeapi.DiskConfig
}
func (d diskW) Keyword() string {
return "disk"
}
func (d diskW) GetSchedtags() []*computeapi.SchedtagConfig {
return d.DiskConfig.Schedtags
}
func (p *DiskSchedtagPredicate) GetInputs(u *core.Unit) []ISchedtagCustomer {
ret := make([]ISchedtagCustomer, 0)
for _, disk := range u.SchedData().Disks {
ret = append(ret, &diskW{disk})
}
p.Hypervisor = computeapi.HOSTTYPE_HYPERVISOR[u.SchedData().Hypervisor]
// always select each storages to disks
u.AppendSelectPlugin(p)
return true, nil
return ret
}
type schedtagStorageW struct {
candidater *api.CandidateStorage
disk *computeapi.DiskConfig
}
func (w schedtagStorageW) IndexKey() string {
return fmt.Sprintf("%s:%s", w.candidater.GetName(), w.candidater.StorageType)
}
func (w schedtagStorageW) GetDynamicSchedDesc() *jsonutils.JSONDict {
ret := jsonutils.NewDict()
storageSchedDesc := w.candidater.GetDynamicConditionInput()
diskSchedDesc := w.disk.JSON(w.disk)
ret.Add(storageSchedDesc, models.StorageManager.Keyword())
ret.Add(diskSchedDesc, models.DiskManager.Keyword())
func (p *DiskSchedtagPredicate) GetResources(c core.Candidater) []ISchedtagCandidateResource {
ret := make([]ISchedtagCandidateResource, 0)
for _, storage := range c.Getter().Storages() {
ret = append(ret, storage)
}
return ret
}
func (w schedtagStorageW) GetSchedtags() []models.SSchedtag {
return w.candidater.Schedtags
}
func (w schedtagStorageW) ResourceType() string {
return models.StorageManager.KeywordPlural()
}
func (p *DiskSchedtagPredicate) check(d *computeapi.DiskConfig, s *api.CandidateStorage, u *core.Unit, c core.Candidater) (*PredicatedStorage, error) {
allTags, err := GetAllSchedtags(models.StorageManager.KeywordPlural())
if err != nil {
return nil, err
}
tagPredicate := NewSchedtagPredicate(d.Schedtags, allTags)
shouldExec := u.ShouldExecuteSchedtagFilter(c.Getter().Id())
ps := newPredicatedStorage(s, nil, nil)
if shouldExec {
if err := tagPredicate.Check(
schedtagStorageW{
candidater: s,
disk: d,
},
); err != nil {
return nil, err
func (p *DiskSchedtagPredicate) IsResourceFitInput(u *core.Unit, res ISchedtagCandidateResource, input ISchedtagCustomer) error {
storage := res.(*api.CandidateStorage)
d := input.(*diskW)
if d.Storage != "" {
if storage.Id != d.Storage && storage.Name != d.Storage {
return fmt.Errorf("Storage name %s != (%s:%s)", d.Storage, storage.Name, storage.Id)
}
ps.PreferTags = tagPredicate.GetPreferTags()
ps.AvoidTags = tagPredicate.GetAvoidTags()
}
return ps, nil
}
func (p *DiskSchedtagPredicate) checkStorages(d *computeapi.DiskConfig, storages []*api.CandidateStorage, u *core.Unit, c core.Candidater) ([]*PredicatedStorage, error) {
errs := make([]error, 0)
ret := make([]*PredicatedStorage, 0)
for _, s := range storages {
ps, err := p.check(d, s, u, c)
if err != nil {
// append err, storage not suit disk
errs = append(errs, err)
continue
if !(len(d.Backend) == 0 || d.Backend == computeapi.STORAGE_LOCAL) {
if storage.StorageType != d.Backend {
return fmt.Errorf("Storage %s backend %s != %s", storage.Name, storage.StorageType, d.Backend)
}
ret = append(ret, ps)
}
if len(ret) == 0 {
return nil, errors.NewAggregate(errs)
}
return ret, nil
}
func (p *DiskSchedtagPredicate) GetDiskStoragesMap(candidateId string) DiskStoragesMap {
ret, ok := p.CandidateDiskStoragesMap[candidateId]
if !ok {
ret = make(map[int][]*PredicatedStorage)
p.CandidateDiskStoragesMap[candidateId] = ret
storageTypes := p.GetHypervisorDriver().GetStorageTypes()
if len(storageTypes) != 0 && !utils.IsInStringArray(storage.StorageType, storageTypes) {
return fmt.Errorf("Storage %s storage type %s not in %v", storage.Name, storage.StorageType, storageTypes)
}
return ret
return nil
}
func (p *DiskSchedtagPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) {
h := NewPredicateHelper(p, u, c)
storages := c.Getter().Storages()
ds := p.GetDiskStoragesMap(c.IndexKey())
disks := u.SchedData().Disks
for idx, d := range disks {
fitStorages := make([]*api.CandidateStorage, 0)
for _, s := range storages {
if p.isStorageFitDisk(s, d) {
fitStorages = append(fitStorages, s)
}
}
if len(fitStorages) == 0 {
h.Exclude(fmt.Sprintf("Not found available storages for disk backend %q", d.Backend))
break
}
matchedStorages, err := p.checkStorages(d, fitStorages, u, c)
if err != nil {
h.Exclude(err.Error())
}
ds[idx] = matchedStorages
}
return h.GetResult()
return p.BaseSchedtagPredicate.Execute(p, u, c)
}
func (p *DiskSchedtagPredicate) OnPriorityEnd(u *core.Unit, c core.Candidater) {
storageTags := []models.SSchedtag{}
for _, s := range c.Getter().Storages() {
storageTags = append(storageTags, s.Schedtags...)
}
ds := p.GetDiskStoragesMap(c.IndexKey())
avoidTags := ds.GetAvoidTags()
preferTags := ds.GetPreferTags()
avoidCountMap := GetSchedtagCount(avoidTags, storageTags, api.AggregateStrategyAvoid)
preferCountMap := GetSchedtagCount(preferTags, storageTags, api.AggregateStrategyPrefer)
setScore := SetCandidateScoreBySchedtag
setScore(u, c, preferCountMap, true)
setScore(u, c, avoidCountMap, false)
p.BaseSchedtagPredicate.OnPriorityEnd(p, u, c)
}
func (p *DiskSchedtagPredicate) OnSelectEnd(u *core.Unit, c core.Candidater, count int64) {
res := u.GetAllocatedResource(c.IndexKey())
diskStorages := p.GetDiskStoragesMap(c.IndexKey())
res.Disks = make([]*schedapi.CandidateDisk, len(diskStorages))
disks := u.SchedData().Disks
for idx, ds := range diskStorages {
res.Disks[idx] = p.allocatedDiskResource(c, disks[idx], ds)
}
p.BaseSchedtagPredicate.OnSelectEnd(p, u, c, count)
}
func (p *DiskSchedtagPredicate) allocatedDiskResource(c core.Candidater, disk *computeapi.DiskConfig, storages []*PredicatedStorage) *schedapi.CandidateDisk {
sortStorages := p.selectStorages(disk, storages)
func (p *DiskSchedtagPredicate) DoSelect(
c core.Candidater,
input ISchedtagCustomer,
res []ISchedtagCandidateResource,
) []ISchedtagCandidateResource {
return p.GetUsedStorages(res, input.(*diskW).Backend)
}
func (p *DiskSchedtagPredicate) GetCandidateResourceSortScore(selectRes ISchedtagCandidateResource) int {
return selectRes.(*api.CandidateStorage).GetFreeCapacity()
}
func (p *DiskSchedtagPredicate) AddSelectResult(index int, selectRes []ISchedtagCandidateResource, output *core.AllocatedResource) {
storageIds := []string{}
for _, s := range sortStorages {
storageIds = append(storageIds, s.Id)
for _, res := range selectRes {
storageIds = append(storageIds, res.GetId())
}
log.Debugf("Suggestion %s storages %v for disk: %d", c.Getter().Name(), storageIds, disk.Index)
return &schedapi.CandidateDisk{
Index: disk.Index,
ret := &schedapi.CandidateDisk{
Index: index,
StorageIds: storageIds,
}
log.Debugf("Suggestion storages %v for disk%d", storageIds, index)
output.Disks = append(output.Disks, ret)
}
func (p *DiskSchedtagPredicate) selectStorages(d *computeapi.DiskConfig, storages []*PredicatedStorage) []*api.CandidateStorage {
preferStorages := []*api.CandidateStorage{}
noTagStorages := []*api.CandidateStorage{}
avoidStorages := []*api.CandidateStorage{}
for _, storage := range storages {
candi := storage.CandidateStorage
if storage.isNoTag() {
noTagStorages = append(noTagStorages, candi)
} else if storage.hasPreferTags() {
preferStorages = append(preferStorages, candi)
} else if storage.hasAvoidTags() {
avoidStorages = append(avoidStorages, candi)
}
func (p *DiskSchedtagPredicate) GetUsedStorages(res []ISchedtagCandidateResource, backend string) []ISchedtagCandidateResource {
storages := make([]ISchedtagCandidateResource, 0)
for _, s := range res {
storages = append(storages, s.(*api.CandidateStorage))
}
sortStorages := []*api.CandidateStorage{}
sortStorages = append(sortStorages, preferStorages...)
sortStorages = append(sortStorages, noTagStorages...)
sortStorages = append(sortStorages, avoidStorages...)
return p.SortLeastUsedStorage(sortStorages, d.Backend)
}
func (p *DiskSchedtagPredicate) SortLeastUsedStorage(storages []*api.CandidateStorage, backend string) []*api.CandidateStorage {
backendStorages := make([]*api.CandidateStorage, 0)
backendStorages := make([]ISchedtagCandidateResource, 0)
if backend != "" {
for _, s := range storages {
if s.StorageType == backend {
if s.(*api.CandidateStorage).StorageType == backend {
backendStorages = append(backendStorages, s)
}
}
@@ -314,51 +148,5 @@ func (p *DiskSchedtagPredicate) SortLeastUsedStorage(storages []*api.CandidateSt
if len(backendStorages) == 0 {
backendStorages = storages
}
return p.sortByLeastUsedStorages(backendStorages)
}
type sortStorages []*api.CandidateStorage
func (s sortStorages) Len() int {
return len(s)
}
func (s sortStorages) Less(i, j int) bool {
s1, s2 := s[i], s[j]
cap1 := s1.GetFreeCapacity()
cap2 := s2.GetFreeCapacity()
return cap1 > cap2
}
func (s sortStorages) Swap(i, j int) {
s[i], s[j] = s[j], s[i]
}
func (p *DiskSchedtagPredicate) sortByLeastUsedStorages(storages []*api.CandidateStorage) []*api.CandidateStorage {
ss := sortStorages(storages)
sort.Sort(ss)
return []*api.CandidateStorage(ss)
}
func (p *DiskSchedtagPredicate) GetHypervisorDriver() models.IGuestDriver {
return models.GetDriver(p.Hypervisor)
}
func (p *DiskSchedtagPredicate) isStorageFitDisk(storage *api.CandidateStorage, d *computeapi.DiskConfig) bool {
if d.Storage != "" {
if storage.Id == d.Storage || storage.Name == d.Storage {
return true
}
return false
}
if storage.StorageType == d.Backend {
return true
}
for _, stype := range p.GetHypervisorDriver().GetStorageTypes() {
if storage.StorageType == stype {
return true
}
}
return false
return backendStorages
}
@@ -27,6 +27,7 @@ import (
computeapi "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/scheduler/algorithm/plugin"
"yunion.io/x/onecloud/pkg/scheduler/algorithm/predicates"
"yunion.io/x/onecloud/pkg/scheduler/api"
"yunion.io/x/onecloud/pkg/scheduler/core"
)
@@ -96,7 +97,7 @@ func (p *NetworkPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []cor
errMsgs = append(errMsgs, errMsg)
}
if !isMatchServerType(&n) {
if !isMatchServerType(n.SNetwork) {
appendError(predicates.ErrServerTypeIsNotMatch)
}
@@ -119,7 +120,7 @@ func (p *NetworkPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []cor
if len(errMsgs) == 0 {
// add resource
reservedNetworks := 0
counter := counterOfNetwork(u, &n, reservedNetworks)
counter := counterOfNetwork(u, n.SNetwork, reservedNetworks)
p.SelectedNetworks.Store(n.GetId(), counter.GetCount())
counters.Add(counter)
found = true
@@ -146,7 +147,7 @@ func (p *NetworkPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []cor
}
isNetworkAvaliable := func(n *computeapi.NetworkConfig, counters *core.MinCounters,
networks []models.SNetwork) string {
networks []*api.CandidateNetwork) string {
if n.Network == "" {
counters0 := core.NewCounters()
ret_msg := isRandomNetworkAvailable(n.Private, n.Exit, n.Wire, counters0)
@@ -173,7 +174,7 @@ func (p *NetworkPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []cor
} else {
// add resource
reservedNetworks := 0
counter := counterOfNetwork(u, &net, reservedNetworks)
counter := counterOfNetwork(u, net.SNetwork, reservedNetworks)
if counter.GetCount() < int64(d.Count) {
errMsgs = append(errMsgs, fmt.Sprintf("%s: ports not enough, free: %d, required: %d", net.Name, counter.GetCount(), d.Count))
continue
@@ -25,6 +25,7 @@ import (
"yunion.io/x/onecloud/pkg/compute/models"
"yunion.io/x/onecloud/pkg/scheduler/api"
"yunion.io/x/onecloud/pkg/scheduler/cache/candidate"
)
@@ -72,13 +73,13 @@ func (p *NetworkPredicate) Execute(cli *kubernetes.Clientset, pod *v1.Pod, node
return true, nil
}
func (p NetworkPredicate) checkByNetworks(nets []models.SNetwork) error {
func (p NetworkPredicate) checkByNetworks(nets []*api.CandidateNetwork) error {
if len(nets) == 0 {
return fmt.Errorf("Network is empty")
}
errs := make([]error, 0)
for _, net := range nets {
err := p.checkByNetwork(&net)
err := p.checkByNetwork(net.SNetwork)
if err == nil {
return nil
}
@@ -97,13 +98,13 @@ func (p NetworkPredicate) checkByNetwork(net *models.SNetwork) error {
return nil
}
func (p NetworkPredicate) checkNetworksIP(ip string, nets []models.SNetwork) error {
func (p NetworkPredicate) checkNetworksIP(ip string, nets []*api.CandidateNetwork) error {
if len(nets) == 0 {
return fmt.Errorf("Network is empty")
}
errs := make([]error, 0)
for _, net := range nets {
err := p.checkNetworkIP(ip, &net)
err := p.checkNetworkIP(ip, net.SNetwork)
if err == nil {
return nil
}
@@ -0,0 +1,152 @@
package predicates
import (
"fmt"
"yunion.io/x/log"
"yunion.io/x/pkg/util/netutils"
"yunion.io/x/pkg/utils"
computeapi "yunion.io/x/onecloud/pkg/apis/compute"
schedapi "yunion.io/x/onecloud/pkg/apis/scheduler"
"yunion.io/x/onecloud/pkg/scheduler/api"
"yunion.io/x/onecloud/pkg/scheduler/core"
)
type NetworkSchedtagPredicate struct {
*BaseSchedtagPredicate
}
func (p *NetworkSchedtagPredicate) Name() string {
return "network_schedtag"
}
func (p *NetworkSchedtagPredicate) Clone() core.FitPredicate {
return &NetworkSchedtagPredicate{
BaseSchedtagPredicate: NewBaseSchedtagPredicate(),
}
}
func (p *NetworkSchedtagPredicate) PreExecute(u *core.Unit, cs []core.Candidater) (bool, error) {
return p.BaseSchedtagPredicate.PreExecute(p, u, cs)
}
type netW struct {
*computeapi.NetworkConfig
}
func (n netW) Keyword() string {
return "net"
}
func (n netW) GetSchedtags() []*computeapi.SchedtagConfig {
return n.NetworkConfig.Schedtags
}
func (p *NetworkSchedtagPredicate) GetInputs(u *core.Unit) []ISchedtagCustomer {
ret := make([]ISchedtagCustomer, 0)
for _, net := range u.SchedData().Networks {
ret = append(ret, &netW{net})
}
return ret
}
func (p *NetworkSchedtagPredicate) GetResources(c core.Candidater) []ISchedtagCandidateResource {
ret := make([]ISchedtagCandidateResource, 0)
for _, network := range c.Getter().Networks() {
ret = append(ret, network)
}
return ret
}
func (p *NetworkSchedtagPredicate) IsResourceFitInput(u *core.Unit, res ISchedtagCandidateResource, input ISchedtagCustomer) error {
network := res.(*api.CandidateNetwork)
net := input.(*netW)
if net.Network != "" {
if network.Id != net.Network && network.Name != net.Network {
return fmt.Errorf("Network name %s != (%s:%s)", net.Network, network.Name, network.Id)
}
}
if net.Wire != "" {
if network.WireId != net.Wire {
return fmt.Errorf("Wire %s != %s", net.Wire, network.WireId)
}
}
netTypes := p.GetNetworkTypes(net.NetType)
if net.Network == "" && !utils.IsInStringArray(network.ServerType, netTypes) {
return fmt.Errorf("Network %s type %s not in %v", network.Name, network.ServerType, netTypes)
}
schedData := u.SchedData()
if net.Private {
if network.IsPublic {
return fmt.Errorf("Network %s is public", network.Name)
}
if network.ProjectId != schedData.Project {
return fmt.Errorf("Network project %s not owner by %s", network.ProjectId, schedData.Project)
}
} else {
if !network.IsPublic {
return fmt.Errorf("Network %s is private", network.Name)
}
}
if len(net.Address) > 0 {
ipAddr, err := netutils.NewIPV4Addr(net.Address)
if err != nil {
return fmt.Errorf("Invalid ip address %s: %v", net.Address, err)
}
if !network.GetIPRange().Contains(ipAddr) {
return fmt.Errorf("Address %s not in range", net.Address)
}
}
free := network.GetFreeAddressCount()
req := u.SchedData().Count
if free < req {
return fmt.Errorf("Network %s no free IPs, free %d, require %d", network.Name, free, req)
}
return nil
}
func (p *NetworkSchedtagPredicate) GetNetworkTypes(specifyType string) []string {
netTypes := p.GetHypervisorDriver().GetRandomNetworkTypes()
if len(specifyType) > 0 {
netTypes = []string{specifyType}
}
return netTypes
}
func (p *NetworkSchedtagPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) {
return p.BaseSchedtagPredicate.Execute(p, u, c)
}
func (p *NetworkSchedtagPredicate) OnPriorityEnd(u *core.Unit, c core.Candidater) {
p.BaseSchedtagPredicate.OnPriorityEnd(p, u, c)
}
func (p *NetworkSchedtagPredicate) OnSelectEnd(u *core.Unit, c core.Candidater, count int64) {
p.BaseSchedtagPredicate.OnSelectEnd(p, u, c, count)
}
func (p *NetworkSchedtagPredicate) GetCandidateResourceSortScore(selectRes ISchedtagCandidateResource) int {
return selectRes.(*api.CandidateNetwork).GetFreeAddressCount()
}
func (p *NetworkSchedtagPredicate) DoSelect(
c core.Candidater,
input ISchedtagCustomer,
res []ISchedtagCandidateResource,
) []ISchedtagCandidateResource {
return res
}
func (p *NetworkSchedtagPredicate) AddSelectResult(index int, selectRes []ISchedtagCandidateResource, output *core.AllocatedResource) {
networkIds := []string{}
for _, res := range selectRes {
networkIds = append(networkIds, res.GetId())
}
ret := &schedapi.CandidateNet{
Index: index,
NetworkIds: networkIds,
}
log.Debugf("Suggestion networks %v for net%d", networkIds, index)
output.Nets = append(output.Nets, ret)
}
@@ -16,11 +16,18 @@ package predicates
import (
"fmt"
"sort"
"strings"
"yunion.io/x/jsonutils"
"yunion.io/x/log"
"yunion.io/x/pkg/util/errors"
computeapi "yunion.io/x/onecloud/pkg/apis/compute"
"yunion.io/x/onecloud/pkg/compute/models"
"yunion.io/x/onecloud/pkg/scheduler/algorithm"
"yunion.io/x/onecloud/pkg/scheduler/algorithm/plugin"
"yunion.io/x/onecloud/pkg/scheduler/api"
"yunion.io/x/onecloud/pkg/scheduler/cache/candidate"
"yunion.io/x/onecloud/pkg/scheduler/core"
"yunion.io/x/onecloud/pkg/scheduler/data_manager"
@@ -179,3 +186,339 @@ func (h *PredicateHelper) UseReserved() bool {
}
return usable
}
type PredicatedSchedtagResource struct {
ISchedtagCandidateResource
PreferTags []computeapi.SchedtagConfig
AvoidTags []computeapi.SchedtagConfig
}
type SchedtagInputResourcesMap map[int][]*PredicatedSchedtagResource
func (m SchedtagInputResourcesMap) getAllTags(isPrefer bool) []computeapi.SchedtagConfig {
ret := make([]computeapi.SchedtagConfig, 0)
for _, ss := range m {
for _, s := range ss {
var tags []computeapi.SchedtagConfig
if isPrefer {
tags = s.PreferTags
} else {
tags = s.AvoidTags
}
ret = append(ret, tags...)
}
}
return ret
}
func (m SchedtagInputResourcesMap) GetPreferTags() []computeapi.SchedtagConfig {
return m.getAllTags(true)
}
func (m SchedtagInputResourcesMap) GetAvoidTags() []computeapi.SchedtagConfig {
return m.getAllTags(false)
}
type CandidateInputResourcesMap map[string]SchedtagInputResourcesMap
type ISchedtagCandidateResource interface {
GetName() string
GetId() string
Keyword() string
GetSchedtags() []models.SSchedtag
GetSchedtagJointManager() models.ISchedtagJointManager
GetDynamicConditionInput() *jsonutils.JSONDict
}
type ISchedtagPredicateInstance interface {
core.FitPredicate
OnPriorityEnd(u *core.Unit, c core.Candidater)
OnSelectEnd(u *core.Unit, c core.Candidater, count int64)
GetInputs(u *core.Unit) []ISchedtagCustomer
GetResources(c core.Candidater) []ISchedtagCandidateResource
IsResourceFitInput(unit *core.Unit, res ISchedtagCandidateResource, input ISchedtagCustomer) error
DoSelect(c core.Candidater, input ISchedtagCustomer, res []ISchedtagCandidateResource) []ISchedtagCandidateResource
AddSelectResult(index int, selectRes []ISchedtagCandidateResource, output *core.AllocatedResource)
GetCandidateResourceSortScore(candidate ISchedtagCandidateResource) int
}
type BaseSchedtagPredicate struct {
BasePredicate
plugin.BasePlugin
CandidateInputResourcesMap CandidateInputResourcesMap
Hypervisor string
}
func NewBaseSchedtagPredicate() *BaseSchedtagPredicate {
return &BaseSchedtagPredicate{
CandidateInputResourcesMap: make(map[string]SchedtagInputResourcesMap),
}
}
func (p *PredicatedSchedtagResource) isNoTag() bool {
return len(p.PreferTags) == 0 && len(p.AvoidTags) == 0
}
func (p *PredicatedSchedtagResource) hasPreferTags() bool {
return len(p.PreferTags) != 0
}
func (p *PredicatedSchedtagResource) hasAvoidTags() bool {
return len(p.AvoidTags) != 0
}
type ISchedtagCustomer interface {
JSON(interface{}) *jsonutils.JSONDict
Keyword() string
GetSchedtags() []*computeapi.SchedtagConfig
}
type SchedtagResourceW struct {
candidater ISchedtagCandidateResource
input ISchedtagCustomer
}
func (w SchedtagResourceW) IndexKey() string {
return fmt.Sprintf("%s:%s", w.candidater.GetName(), w.candidater.GetId())
}
func (w SchedtagResourceW) ResourceType() string {
return getSchedtagResourceType(w.candidater)
}
func getSchedtagResourceType(candidater ISchedtagCandidateResource) string {
return candidater.GetSchedtagJointManager().GetMasterManager().KeywordPlural()
}
func (w SchedtagResourceW) GetSchedtags() []models.SSchedtag {
return w.candidater.GetSchedtags()
}
func (w SchedtagResourceW) GetDynamicSchedDesc() *jsonutils.JSONDict {
ret := jsonutils.NewDict()
resSchedDesc := w.candidater.GetDynamicConditionInput()
inputSchedDesc := w.input.JSON(w.input)
ret.Add(resSchedDesc, w.candidater.Keyword())
ret.Add(inputSchedDesc, w.input.Keyword())
return ret
}
func (p *BaseSchedtagPredicate) GetHypervisorDriver() models.IGuestDriver {
return models.GetDriver(p.Hypervisor)
}
func (p *BaseSchedtagPredicate) check(input ISchedtagCustomer, candidate ISchedtagCandidateResource, u *core.Unit, c core.Candidater) (*PredicatedSchedtagResource, error) {
allTags, err := GetAllSchedtags(getSchedtagResourceType(candidate))
if err != nil {
return nil, err
}
tagPredicate := NewSchedtagPredicate(input.GetSchedtags(), allTags)
shouldExec := u.ShouldExecuteSchedtagFilter(c.Getter().Id())
res := &PredicatedSchedtagResource{
ISchedtagCandidateResource: candidate,
}
if shouldExec {
if err := tagPredicate.Check(
SchedtagResourceW{
candidater: candidate,
input: input,
},
); err != nil {
return nil, err
}
res.PreferTags = tagPredicate.GetPreferTags()
res.AvoidTags = tagPredicate.GetAvoidTags()
}
return res, nil
}
func (p *BaseSchedtagPredicate) checkResources(input ISchedtagCustomer, ress []ISchedtagCandidateResource, u *core.Unit, c core.Candidater) ([]*PredicatedSchedtagResource, error) {
errs := make([]error, 0)
ret := make([]*PredicatedSchedtagResource, 0)
for _, res := range ress {
ps, err := p.check(input, res, u, c)
if err != nil {
// append err, resource not suit input customer
errs = append(errs, err)
continue
}
ret = append(ret, ps)
}
if len(ret) == 0 {
return nil, errors.NewAggregate(errs)
}
return ret, nil
}
func (p *BaseSchedtagPredicate) GetInputResourcesMap(candidateId string) SchedtagInputResourcesMap {
ret, ok := p.CandidateInputResourcesMap[candidateId]
if !ok {
ret = make(map[int][]*PredicatedSchedtagResource)
p.CandidateInputResourcesMap[candidateId] = ret
}
return ret
}
func (p *BaseSchedtagPredicate) PreExecute(sp ISchedtagPredicateInstance, u *core.Unit, cs []core.Candidater) (bool, error) {
input := sp.GetInputs(u)
if len(input) == 0 {
return false, nil
}
p.Hypervisor = computeapi.HOSTTYPE_HYPERVISOR[u.SchedData().Hypervisor]
// always do select step
u.AppendSelectPlugin(sp)
return true, nil
}
func (p *BaseSchedtagPredicate) Execute(
sp ISchedtagPredicateInstance,
u *core.Unit,
c core.Candidater,
) (bool, []core.PredicateFailureReason, error) {
inputs := sp.GetInputs(u)
resources := sp.GetResources(c)
h := NewPredicateHelper(sp, u, c)
inputRes := p.GetInputResourcesMap(c.IndexKey())
filterErrs := make([]error, 0)
for idx, input := range inputs {
fitResources := make([]ISchedtagCandidateResource, 0)
errs := make([]error, 0)
for _, res := range resources {
if err := sp.IsResourceFitInput(u, res, input); err == nil {
fitResources = append(fitResources, res)
} else {
errs = append(errs, err)
}
}
if len(fitResources) == 0 {
h.Exclude(fmt.Sprintf("Not found available resources for %s %s: %s", input.Keyword(), input.JSON(input), errors.NewAggregate(errs)))
break
}
if len(errs) > 0 {
filterErrs = append(filterErrs, errors.NewAggregate(errs))
}
matchedResources, err := p.checkResources(input, fitResources, u, c)
if err != nil {
log.Errorf("=========err: %v, filterErrs: %v", err, filterErrs)
aggErr := errors.NewAggregate(filterErrs)
errMsg := fmt.Sprintf("schedtag: %v", err.Error())
if aggErr != nil {
errMsg = fmt.Sprintf("%s; filter: %v", errMsg, aggErr.Error())
}
h.Exclude(errMsg)
}
inputRes[idx] = matchedResources
}
return h.GetResult()
}
func (p *BaseSchedtagPredicate) OnPriorityEnd(sp ISchedtagPredicateInstance, u *core.Unit, c core.Candidater) {
resTags := []models.SSchedtag{}
for _, res := range sp.GetResources(c) {
resTags = append(resTags, res.GetSchedtags()...)
}
inputRes := p.GetInputResourcesMap(c.IndexKey())
avoidTags := inputRes.GetAvoidTags()
preferTags := inputRes.GetPreferTags()
avoidCountMap := GetSchedtagCount(avoidTags, resTags, api.AggregateStrategyAvoid)
preferCountMap := GetSchedtagCount(preferTags, resTags, api.AggregateStrategyPrefer)
setScore := SetCandidateScoreBySchedtag
setScore(u, c, preferCountMap, true)
setScore(u, c, avoidCountMap, false)
}
func (p *BaseSchedtagPredicate) OnSelectEnd(sp ISchedtagPredicateInstance, u *core.Unit, c core.Candidater, count int64) {
inputRes := p.GetInputResourcesMap(c.IndexKey())
output := u.GetAllocatedResource(c.IndexKey())
inputs := sp.GetInputs(u)
for idx, res := range inputRes {
selRes := p.selectResource(sp, c, inputs[idx], res)
sortRes := newSortCandidateResource(sp, selRes)
sort.Sort(sortRes)
//log.Debugf("sort result: %s", sortRes.DebugString())
sp.AddSelectResult(idx, sortRes.res, output)
}
}
type sortCandidateResource struct {
predicate ISchedtagPredicateInstance
res []ISchedtagCandidateResource
}
func newSortCandidateResource(predicate ISchedtagPredicateInstance, res []ISchedtagCandidateResource) *sortCandidateResource {
return &sortCandidateResource{
predicate: predicate,
res: res,
}
}
func (s *sortCandidateResource) Len() int {
return len(s.res)
}
func (s *sortCandidateResource) DebugString() string {
var debugStr string
for _, i := range s.res {
debugStr = fmt.Sprintf("%s %d", debugStr, s.predicate.GetCandidateResourceSortScore(i))
}
return debugStr
}
// desc order
func (s *sortCandidateResource) Less(i, j int) bool {
res1, res2 := s.res[i], s.res[j]
v1 := s.predicate.GetCandidateResourceSortScore(res1)
v2 := s.predicate.GetCandidateResourceSortScore(res2)
return v1 > v2
}
func (s *sortCandidateResource) Swap(i, j int) {
s.res[i], s.res[j] = s.res[j], s.res[i]
}
func (p *BaseSchedtagPredicate) selectResource(
sp ISchedtagPredicateInstance,
c core.Candidater,
input ISchedtagCustomer,
ress []*PredicatedSchedtagResource,
) []ISchedtagCandidateResource {
preferRes := make([]ISchedtagCandidateResource, 0)
noTagRes := make([]ISchedtagCandidateResource, 0)
avoidRes := make([]ISchedtagCandidateResource, 0)
for _, res := range ress {
if res.isNoTag() {
noTagRes = append(noTagRes, res.ISchedtagCandidateResource)
} else if res.hasPreferTags() {
preferRes = append(preferRes, res.ISchedtagCandidateResource)
} else if res.hasAvoidTags() {
avoidRes = append(avoidRes, res.ISchedtagCandidateResource)
}
}
for _, ress := range [][]ISchedtagCandidateResource{
preferRes,
noTagRes,
avoidRes,
} {
if len(ress) == 0 {
continue
}
if ret := sp.DoSelect(c, input, ress); ret != nil {
return ret
}
}
return nil
}
@@ -265,15 +265,16 @@ func (c *SchedtagChecker) Check(p ISchedtagPredicate, candidate ISchedtagCandida
log.V(10).Debugf("[SchedtagChecker] check candidate: %s requireTags: %#v, execludeTags: %#v, candidateTags: %#v", candidate.IndexKey(), requireTags, execludeTags, candidateTags)
candiInfo := fmt.Sprintf("%s:%s", candidate.ResourceType(), candidate.IndexKey())
if len(execludeTags) > 0 {
if ok, tag := c.HasIntersection(execludeTags, candidateTags); ok {
return fmt.Errorf("Execlude by schedtag: '%s:%s'", tag.Name, tag.Id)
return fmt.Errorf("schedtag %q exclude %s", tag.Name, candiInfo)
}
}
if len(requireTags) > 0 {
if ok, tag := c.Contains(candidateTags, requireTags); !ok {
return fmt.Errorf("Need schedtag: '%s'", tag.Id)
return fmt.Errorf("%s need schedtag: %q", candiInfo, tag.Id)
}
}