Add scheduler service.

This commit is contained in:
Zexi Li
2018-07-30 11:18:02 +08:00
parent 8ba5ffe1b4
commit 270298fae2
150 changed files with 19220 additions and 0 deletions
@@ -0,0 +1,218 @@
package predicates
import (
"fmt"
"github.com/yunionio/log"
"github.com/yunionio/onecloud/pkg/scheduler/api"
"github.com/yunionio/onecloud/pkg/scheduler/core"
"github.com/yunionio/onecloud/pkg/scheduler/db/models"
)
// NOTE: Aggregate Description
// require: Must be scheduled to the specified host
// prefer: Priority to the specified host
// avoid: Try to avoid scheduling to the specified host
// exclude: Do not allow scheduling on the specified host
// AggregatePredicate is designed to quickly filter unavailable
// hosts and improve scheduling efficiency by tabbing whether
// the host is available.
type AggregatePredicate struct {
BasePredicate
AggregateHosts hostsAggregatesMap
RequireAggregates []api.Aggregate
ExcludeAggregates []api.Aggregate
AvoidAggregates []api.Aggregate
PreferAggregates []api.Aggregate
AggregateMap map[string]api.Aggregate
}
type hostAggregates []*models.Aggregate
type hostsAggregatesMap map[string]hostAggregates
func (p *AggregatePredicate) Name() string {
return "host_aggregate"
}
func (p *AggregatePredicate) Clone() core.FitPredicate {
return &AggregatePredicate{
AggregateMap: make(map[string]api.Aggregate, 0),
}
}
func hostsAggregatesInfo(cs []core.Candidater) (hostsAggregatesMap, []*models.Aggregate) {
ret := make(map[string]hostAggregates, 0)
allAggs := make([]*models.Aggregate, 0)
for _, c := range cs {
hostAggs := c.GetHostAggregates()
ret[c.IndexKey()] = hostAggs
}
if len(cs) > 0 {
allAggs = cs[0].GetAggregates()
}
return ret, allAggs
}
func (p *AggregatePredicate) PreExecute(u *core.Unit, cs []core.Candidater) (bool, error) {
data := u.SchedData()
if len(data.Candidates) > 0 {
return false, nil
}
hsMap, allAggs := hostsAggregatesInfo(cs)
p.AggregateHosts = hsMap
appendedAggIds := make(map[string]int, len(data.Aggregates))
for _, aggregate := range data.Aggregates {
switch aggregate.Strategy {
case api.AggregateStrategyRequire:
p.RequireAggregates = append(p.RequireAggregates, aggregate)
case api.AggregateStrategyExclude:
p.ExcludeAggregates = append(p.ExcludeAggregates, aggregate)
case api.AggregateStrategyPrefer:
p.PreferAggregates = append(p.PreferAggregates, aggregate)
case api.AggregateStrategyAvoid:
p.AvoidAggregates = append(p.AvoidAggregates, aggregate)
}
p.AggregateMap[aggregate.Idx] = aggregate
appendedAggIds[aggregate.Idx] = 1
}
for _, aggregate := range allAggs {
_, nameOk := appendedAggIds[aggregate.Name]
_, idOk := appendedAggIds[aggregate.ID]
if !(nameOk || idOk) {
agg := api.Aggregate{Idx: aggregate.ID, Strategy: aggregate.DefaultStrategy}
switch agg.Strategy {
case api.AggregateStrategyRequire:
p.RequireAggregates = append(p.RequireAggregates, agg)
case api.AggregateStrategyExclude:
p.ExcludeAggregates = append(p.ExcludeAggregates, agg)
case api.AggregateStrategyPrefer:
p.PreferAggregates = append(p.PreferAggregates, agg)
case api.AggregateStrategyAvoid:
p.AvoidAggregates = append(p.AvoidAggregates, agg)
}
}
}
u.AppendSelectPlugin(p)
return true, nil
}
func getHostAggregateCount(inAggs []api.Aggregate, hAggs []*models.Aggregate, strategy string) (countMap map[string]int) {
countMap = make(map[string]int)
in := func(hAgg *models.Aggregate, inAggs []api.Aggregate) bool {
for _, agg := range inAggs {
if agg.Idx == hAgg.ID || agg.Idx == hAgg.Name {
return true
}
}
return false
}
for _, hAgg := range hAggs {
if in(hAgg, inAggs) {
countMap[fmt.Sprintf("%s:%s:%s", hAgg.ID, hAgg.Name, strategy)]++
}
}
return
}
func (p *AggregatePredicate) OnSelect(u *core.Unit, c core.Candidater) bool {
hostAggs, ok := p.AggregateHosts[c.IndexKey()]
if !ok {
return true
}
avoidCountMap := getHostAggregateCount(p.AvoidAggregates, hostAggs, api.AggregateStrategyAvoid)
preferCountMap := getHostAggregateCount(p.PreferAggregates, hostAggs, api.AggregateStrategyPrefer)
setScore := func(aggCountMap map[string]int, postiveScore bool) {
stepScore := core.PriorityStep
if !postiveScore {
stepScore = -stepScore
}
for n, count := range aggCountMap {
u.IncreaseScore(c.IndexKey(), n, count*stepScore)
}
}
setScore(avoidCountMap, false)
setScore(preferCountMap, true)
return true
}
func (p *AggregatePredicate) OnSelectEnd(u *core.Unit, c core.Candidater, count int64) {}
func (p *AggregatePredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) {
h := NewPredicateHelper(p, u, c)
if errMsg := p.exec(h); len(errMsg) > 0 {
h.Exclude(errMsg)
}
return h.GetResult()
}
func (p *AggregatePredicate) exec(h *PredicateHelper) string {
ahs := p.AggregateHosts
candidateID := h.Candidate.IndexKey()
log.V(10).Debugf(">>>> ExcludeAggregates: %#v, RequireAggregates: %#v, AvoidAggregates: %#v, PreferAggregates: %#v, candidateID: %v", p.ExcludeAggregates, p.RequireAggregates, p.AvoidAggregates, p.PreferAggregates, candidateID)
if len(p.ExcludeAggregates) > 0 {
inExclude := func(a *models.Aggregate) bool {
for _, agg := range p.ExcludeAggregates {
if agg.Idx == a.ID || agg.Idx == a.Name {
return true
}
}
return false
}
if ah, ok := ahs[candidateID]; ok {
for _, a := range ah {
if inExclude(a) {
return fmt.Sprintf("exclude by aggregate: '%s:%s'", a.Name, a.ID)
}
}
}
}
if len(p.RequireAggregates) > 0 {
var as []*models.Aggregate = nil
if ah, ok := ahs[candidateID]; ok {
as = ah
}
inRequire := func(agg api.Aggregate) bool {
for _, a := range as {
if a.ID == agg.Idx || a.Name == agg.Idx {
return true
}
}
return false
}
for _, agg := range p.RequireAggregates {
if !inRequire(agg) {
return fmt.Sprintf("need aggregate: '%s'", agg.Idx)
}
}
}
return ""
}
@@ -0,0 +1,18 @@
package baremetal
import (
o "github.com/yunionio/onecloud/cmd/scheduler/options"
"github.com/yunionio/onecloud/pkg/scheduler/algorithm/predicates"
"github.com/yunionio/onecloud/pkg/scheduler/core"
)
type BasePredicate struct {
predicates.BasePredicate
}
func (p *BasePredicate) PreExecute(u *core.Unit, cs []core.Candidater) (bool, error) {
if o.GetOptions().DisableBaremetalPredicates {
return false, nil
}
return true, nil
}
@@ -0,0 +1,39 @@
package baremetal
import (
"github.com/yunionio/onecloud/pkg/scheduler/algorithm/predicates"
"github.com/yunionio/onecloud/pkg/scheduler/core"
)
type CPUPredicate struct {
BasePredicate
}
func (p *CPUPredicate) Name() string {
return "baremetal_cpu"
}
func (p *CPUPredicate) Clone() core.FitPredicate {
return &CPUPredicate{}
}
func (p *CPUPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) {
h := predicates.NewPredicateHelper(p, u, c)
d := u.SchedData()
freeCPUCount := h.GetInt64("FreeCPUCount", 0)
reqCPUCount := d.VCPUCount
if freeCPUCount < d.VCPUCount {
totalCPUCount := h.GetInt64("CPUCount", 0)
h.AppendInsufficientResourceError(reqCPUCount, totalCPUCount, freeCPUCount)
h.SetCapacity(0)
} else {
if reqCPUCount/freeCPUCount != 1 {
h.Exclude2("cpu", freeCPUCount, reqCPUCount)
} else {
h.SetCapacity(1)
}
}
return h.GetResult()
}
@@ -0,0 +1,39 @@
package baremetal
import (
"github.com/yunionio/onecloud/pkg/scheduler/algorithm/predicates"
"github.com/yunionio/onecloud/pkg/scheduler/core"
)
type MemoryPredicate struct {
BasePredicate
}
func (p *MemoryPredicate) Name() string {
return "baremetal_memory"
}
func (p *MemoryPredicate) Clone() core.FitPredicate {
return &MemoryPredicate{}
}
func (p *MemoryPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) {
h := predicates.NewPredicateHelper(p, u, c)
d := u.SchedData()
freeMemSize := h.GetInt64("FreeMemSize", 0)
reqMemSize := d.VMEMSize
if freeMemSize < reqMemSize {
totalMemSize := h.GetInt64("MemSize", 0)
h.AppendInsufficientResourceError(reqMemSize, totalMemSize, freeMemSize)
h.SetCapacity(0)
} else {
if reqMemSize/freeMemSize != 1 {
h.Exclude2("memory", freeMemSize, reqMemSize)
} else {
h.SetCapacity(1)
}
}
return h.GetResult()
}
@@ -0,0 +1,150 @@
package baremetal
import (
"fmt"
"strings"
"sync"
"github.com/yunionio/onecloud/pkg/scheduler/algorithm/predicates"
"github.com/yunionio/onecloud/pkg/scheduler/api"
"github.com/yunionio/onecloud/pkg/scheduler/core"
"github.com/yunionio/pkg/utils"
)
type NetworkPredicate struct {
BasePredicate
SelectedNetworks sync.Map
}
func (p *NetworkPredicate) Name() string {
return "baremetal_network"
}
func (p *NetworkPredicate) Clone() core.FitPredicate {
return &NetworkPredicate{}
}
func (p *NetworkPredicate) PreExecute(u *core.Unit, cs []core.Candidater) (bool, error) {
notIgnore, _ := p.BasePredicate.PreExecute(u, cs)
if !notIgnore {
return false, nil
}
u.AppendSelectPlugin(p)
d := u.SchedData()
if len(d.HostID) > 0 && len(d.Networks) == 0 {
return false, nil
}
return true, nil
}
func (p *NetworkPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) {
h := predicates.NewPredicateHelper(p, u, c)
schedData := u.SchedData()
candidate, err := h.BaremetalCandidate()
if err != nil {
return false, nil, err
}
counters := core.NewCounters()
isMigrate := func() bool {
return len(schedData.HostID) > 0
}
isRandomNetworkAvailable := func(private bool, exit bool, wire string) string {
var errMsgs []string
for _, network := range candidate.Networks {
appendError := func(errMsg string) {
errMsgs = append(errMsgs, fmt.Sprintf("%s: %s", network.ID, errMsg))
}
if !((network.Ports > 0 || isMigrate()) && network.IsExit == exit) {
appendError(predicates.ErrNoPorts)
}
if wire != "" && !utils.HasPrefix(wire, network.Wire) && !utils.HasPrefix(wire, network.WireID) { // re
appendError(predicates.ErrWireIsNotMatch)
}
if (!private && network.IsPublic) || (private && !network.IsPublic && network.TenantID == schedData.OwnerTenantID) {
// TODO: support reservedNetworks
reservedNetworks := 0
restPort := int64(network.Ports - reservedNetworks)
if restPort == 0 {
appendError("not enough network port")
continue
}
counter := u.CounterManager.GetOrCreate("net:"+network.ID, func() core.Counter {
return core.NewNormalCounter(restPort)
})
u.SharedResourceManager.Add(network.ID, counter)
counters.Add(counter)
p.SelectedNetworks.Store(network.ID, counter.GetCount())
return ""
} else {
appendError(predicates.ErrNotOwner)
}
}
return strings.Join(errMsgs, ";")
}
filterByRandomNetwork := func() {
if err_msg := isRandomNetworkAvailable(false, false, ""); err_msg != "" {
h.AppendPredicateFailMsg(err_msg)
}
h.SetCapacityCounter(counters)
}
isNetworkAvaliable := func(network *api.Network) string {
if network.Idx == "" {
return isRandomNetworkAvailable(network.Private, network.Exit, network.Wire)
}
for _, net := range candidate.Networks {
if (network.Idx == net.ID || network.Idx == net.Name) && (net.IsPublic || net.TenantID == schedData.OwnerTenantID) && (net.Ports > 0 || isMigrate()) {
h.SetCapacity(1)
return ""
}
}
return predicates.ErrUnknown
}
filterBySpecifiedNetworks := func() {
var errMsgs []string
for _, network := range schedData.Networks {
if err_msg := isNetworkAvaliable(network); err_msg != "" {
errMsgs = append(errMsgs, err_msg)
}
}
if len(errMsgs) > 0 {
h.AppendPredicateFailMsg(strings.Join(errMsgs, ", "))
}
}
loadUnknownNetworks := func() bool {
return true // TODO: ???
}
loadUnknownNetworks()
// Randomly assign networks if no network is specified.
if len(schedData.Networks) == 0 {
filterByRandomNetwork()
} else {
filterBySpecifiedNetworks()
}
return h.GetResult()
}
func (p *NetworkPredicate) OnSelect(u *core.Unit, c core.Candidater) bool {
u.SetFiltedData(c.IndexKey(), "networks", p.SelectedNetworks)
return true
}
func (p *NetworkPredicate) OnSelectEnd(u *core.Unit, c core.Candidater, count int64) {
}
@@ -0,0 +1,51 @@
package baremetal
import (
"github.com/yunionio/onecloud/pkg/scheduler/algorithm/predicates"
"github.com/yunionio/onecloud/pkg/scheduler/core"
"github.com/yunionio/pkg/util/sets"
)
var (
ExpectedStatus = sets.NewString("running", "start_convert")
)
type StatusPredicate struct {
BasePredicate
}
func (p *StatusPredicate) Name() string {
return "baremetal_status"
}
func (p *StatusPredicate) Clone() core.FitPredicate {
return &StatusPredicate{}
}
func (p *StatusPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) {
h := predicates.NewPredicateHelper(p, u, c)
bm, err := h.BaremetalCandidate()
if err != nil {
return false, nil, err
}
if !ExpectedStatus.Has(bm.Status) {
h.Exclude2("status", bm.Status, ExpectedStatus)
return h.GetResult()
}
if !bm.Enabled {
h.Exclude2("enable_status", "disable", "enable")
return h.GetResult()
}
if bm.ServerID == "" {
h.SetCapacity(1)
} else {
h.AppendPredicateFailMsg(predicates.ErrBaremetalHasAlreadyBeenOccupied)
h.SetCapacity(0)
}
return h.GetResult()
}
@@ -0,0 +1,46 @@
package baremetal
import (
"fmt"
//"github.com/yunionio/log"
"github.com/yunionio/onecloud/pkg/scheduler/algorithm/predicates"
"github.com/yunionio/onecloud/pkg/scheduler/core"
"github.com/yunionio/onecloud/pkg/scheduler/util/baremetal"
)
type StoragePredicate struct {
BasePredicate
}
func (p *StoragePredicate) Name() string {
return "baremetal_storage"
}
func (p *StoragePredicate) Clone() core.FitPredicate {
return &StoragePredicate{}
}
func (p *StoragePredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) {
h := predicates.NewPredicateHelper(p, u, c)
schedData := u.SchedData()
candidate, err := h.BaremetalCandidate()
if err != nil {
return false, nil, err
}
layouts, err := baremetal.CalculateLayout(
schedData.BaremetalDiskConfigs,
candidate.Storages,
)
if err == nil && baremetal.CheckDisksAllocable(layouts, schedData.Disks) {
h.SetCapacity(int64(1))
} else {
h.SetCapacity(int64(0))
h.AppendPredicateFailMsg(fmt.Sprintf("%s err: %v", predicates.ErrNoEnoughStorage, err))
}
return h.GetResult()
}
@@ -0,0 +1,79 @@
package predicates
import (
"fmt"
)
// Here are all the errors that may appear in the preselection predicates.
const (
ErrServerTypeIsNotMatch = `server type is not match`
ErrExitIsNotMatch = `exit is not match`
ErrWireIsNotMatch = `wire is not match`
ErrNoPorts = `no ports`
ErrNotOwner = `not owner`
ErrNoEnoughStorage = `no enough storage`
ErrNoAvailableNetwork = `no available network on this host`
ErrNoEnoughAvailableGPUs = `no enough available GPUs`
ErrNotSupportNest = `nested function not supported`
ErrRequireMvs = `require mvs`
ErrRequireNoMvs = `require not mvs`
ErrHostIsSpecifiedForMigration = `host_id specified for migration`
ErrMoreThanOneSizeUnspecificSplit = `more than 1 size unspecific split`
ErrNoMoreSpaceForUnspecificSplit = `no more space for an unspecific split`
ErrSubtotalOfSplitExceedsDiskSize = `subtotal of split exceeds disk size`
ErrBaremetalHasAlreadyBeenOccupied = `baremetal has already been occupied`
ErrUnknown = `unknown error`
)
// InsufficientResourceError is an error type that indicates what kind of resource limit is
// hit and caused the unfitting failure.
type InsufficientResourceError struct {
// resourceName is the name of the resource that is insufficient
ResourceName string
requested int64
total int64
free int64
}
func NewInsufficientResourceError(resourceName string, requested, total, free int64) *InsufficientResourceError {
return &InsufficientResourceError{
ResourceName: resourceName,
requested: requested,
total: total,
free: free,
}
}
func (ire *InsufficientResourceError) Error() string {
return fmt.Sprintf("no enough resource: %s, requested: %d, total: %d, free: %d",
ire.ResourceName, ire.requested, ire.total, ire.free)
}
func (ire *InsufficientResourceError) GetReason() string {
return ire.Error()
}
type UnexceptedResourceError struct {
message string
}
func NewUnexceptedResourceError(message string) *UnexceptedResourceError {
return &UnexceptedResourceError{
message: message,
}
}
func Error(message string) *UnexceptedResourceError {
return NewUnexceptedResourceError(message)
}
func (ure *UnexceptedResourceError) Error() string {
return ure.message
}
func (ure *UnexceptedResourceError) GetReason() string {
return ure.Error()
}
@@ -0,0 +1,54 @@
package guest
import (
"github.com/yunionio/onecloud/pkg/scheduler/algorithm/predicates"
"github.com/yunionio/onecloud/pkg/scheduler/core"
)
// CPUPredicate check the current resources of the CPU is available,
// it returns the maximum available capacity.
type CPUPredicate struct {
predicates.BasePredicate
}
func (f *CPUPredicate) Name() string {
return "host_cpu"
}
func (f *CPUPredicate) Clone() core.FitPredicate {
return &CPUPredicate{}
}
func (f *CPUPredicate) PreExecute(u *core.Unit, cs []core.Candidater) (bool, error) {
if u.IsPublicCloudProvider() {
return false, nil
}
data := u.SchedData()
if data.VCPUCount <= 0 {
return false, nil
}
return true, nil
}
func (f *CPUPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) {
h := predicates.NewPredicateHelper(f, u, c)
d := u.SchedData()
hc, err := h.HostCandidate()
if err != nil {
return false, nil, err
}
useRsvd := h.UseReserved()
freeCPUCount := hc.GetFreeCPUCount(useRsvd)
reqCPUCount := d.VCPUCount
if freeCPUCount < reqCPUCount {
totalCPUCount := hc.GetTotalCPUCount(useRsvd)
h.AppendInsufficientResourceError(reqCPUCount, totalCPUCount, freeCPUCount)
}
h.SetCapacity(freeCPUCount / reqCPUCount)
return h.GetResult()
}
@@ -0,0 +1,103 @@
package guest
import (
"fmt"
"github.com/yunionio/onecloud/pkg/scheduler/algorithm/predicates"
"github.com/yunionio/onecloud/pkg/scheduler/core"
)
// GroupPredicate filter the packet based on the label information,
// the same group of guests should avoid schedule on same host.
type GroupPredicate struct {
predicates.BasePredicate
ExcludeGroups []string
RequireGroups []string
AvoidGroups []string
PreferGroups []string
}
func (p *GroupPredicate) Name() string {
return "host_group"
}
func (p *GroupPredicate) Clone() core.FitPredicate {
return &GroupPredicate{}
}
func (p *GroupPredicate) PreExecute(u *core.Unit, cs []core.Candidater) (bool, error) {
d := u.SchedData()
if len(d.GroupRelations) == 0 {
return false, nil
}
for _, r := range d.GroupRelations {
if r.Strategy == "exclude" {
p.ExcludeGroups = append(p.ExcludeGroups, r.GroupID)
} else if r.Strategy == "require" {
p.RequireGroups = append(p.RequireGroups, r.GroupID)
} else if r.Strategy == "avoid" {
p.AvoidGroups = append(p.AvoidGroups, r.GroupID)
} else if r.Strategy == "prefer" {
p.PreferGroups = append(p.PreferGroups, r.GroupID)
}
}
u.AppendSelectPlugin(p)
return true, nil
}
func (p *GroupPredicate) OnSelect(u *core.Unit, c core.Candidater) bool {
if len(p.ExcludeGroups) > 0 {
return false
}
if len(p.RequireGroups) > 0 {
// TODO: what?
}
if len(p.AvoidGroups) > 0 {
u.IncreaseScore(c.IndexKey(),
p.Name()+":avoid", -core.PriorityStep*len(p.AvoidGroups),
)
}
if len(p.PreferGroups) > 0 {
u.IncreaseScore(c.IndexKey(),
p.Name()+":prefer", core.PriorityStep*len(p.PreferGroups),
)
}
return true
}
func (p *GroupPredicate) OnSelectEnd(u *core.Unit, c core.Candidater, count int64) {
}
func (p *GroupPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) {
h := predicates.NewPredicateHelper(p, u, c)
g, err := h.GetGroupCounts()
if err != nil {
return false, nil, err
}
if len(p.ExcludeGroups) > 0 {
for _, groupId := range p.ExcludeGroups {
if g.ExistsGroup(groupId) {
h.Exclude(fmt.Sprintf("exclude by %v:exclude", groupId))
break
}
}
} else if len(p.RequireGroups) > 0 {
for _, groupId := range p.RequireGroups {
if !g.ExistsGroup(groupId) {
h.Exclude(fmt.Sprintf("exclude by %v:require", groupId))
break
}
}
}
return h.GetResult()
}
@@ -0,0 +1,52 @@
package guest
import (
"github.com/yunionio/log"
"github.com/yunionio/onecloud/pkg/scheduler/algorithm/predicates"
"github.com/yunionio/onecloud/pkg/scheduler/api"
"github.com/yunionio/onecloud/pkg/scheduler/core"
)
const (
CONTAINER_ALLOWED_TAG = "container"
)
// HypervisorPredicate is to select candidates match guest hyperviosr
// runtime
type HypervisorPredicate struct {
predicates.BasePredicate
}
func (f *HypervisorPredicate) Name() string {
return "host_hypervisor_runtime"
}
func (f *HypervisorPredicate) Clone() core.FitPredicate {
return &HypervisorPredicate{}
}
func hostHasContainerTag(c core.Candidater) bool {
aggs := c.GetHostAggregates()
for _, agg := range aggs {
if agg.Name == CONTAINER_ALLOWED_TAG {
return true
}
}
return false
}
func (f *HypervisorPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) {
h := predicates.NewPredicateHelper(f, u, c)
hostType := c.Get("HostType")
guestNeedType := u.SchedData().Hypervisor
if guestNeedType != hostType {
if guestNeedType == api.SchedTypeContainer && hostHasContainerTag(c) {
log.Debugf("Host %q has %q tag, allow it run container", c.IndexKey(), CONTAINER_ALLOWED_TAG)
return h.GetResult()
}
h.Exclude2(f.Name(), hostType, guestNeedType)
}
return h.GetResult()
}
@@ -0,0 +1,112 @@
package guest
import (
"fmt"
"github.com/yunionio/onecloud/pkg/scheduler/algorithm/predicates"
"github.com/yunionio/onecloud/pkg/scheduler/core"
)
// IsolatedDevicePredicate check mode, and number of scheduled
// device configurations and current resources.
type IsolatedDevicePredicate struct {
predicates.BasePredicate
}
func (f *IsolatedDevicePredicate) Name() string {
return "host_isolated_device"
}
func (f *IsolatedDevicePredicate) Clone() core.FitPredicate {
return &IsolatedDevicePredicate{}
}
func (f *IsolatedDevicePredicate) PreExecute(u *core.Unit, cs []core.Candidater) (bool, error) {
data := u.SchedData()
if len(data.IsolatedDevices) == 0 {
return false, nil
}
return true, nil
}
func (f *IsolatedDevicePredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) {
h := predicates.NewPredicateHelper(f, u, c)
reqIsoDevs := u.SchedData().IsolatedDevices
hc, err := h.HostCandidate()
if err != nil {
return false, nil, err
}
minCapacity := int64(0xFFFFFFFF)
// check by specify device id
for _, dev := range reqIsoDevs {
if len(dev.ID) == 0 {
continue
}
if fDev := hc.GetIsolatedDevice(dev.ID); fDev != nil {
if len(fDev.GuestID) != 0 {
h.Exclude(fmt.Sprintf("IsolatedDevice %q already used by guest %q", dev.ID, fDev.GuestID))
return h.GetResult()
}
} else {
h.Exclude(fmt.Sprintf("Not found IsolatedDevice %q", dev.ID))
return h.GetResult()
}
minCapacity = 1
}
reqCount := len(reqIsoDevs)
freeCount := len(hc.UnusedIsolatedDevices())
totalCount := len(hc.IsolatedDevices)
// check host isolated device count
if freeCount < reqCount {
h.AppendInsufficientResourceError(int64(reqCount), int64(totalCount), int64(freeCount))
h.Exclude(fmt.Sprintf(
"IsolatedDevice count not enough, request: %d, hostTotal: %d, hostFree: %d",
reqCount, totalCount, freeCount))
return h.GetResult()
}
// check host device by type
devTypeRequest := make(map[string]int, 0)
for _, dev := range reqIsoDevs {
if len(dev.Type) != 0 {
devTypeRequest[dev.Type] += 1
}
}
for devType, reqCount := range devTypeRequest {
freeCount := len(hc.UnusedIsolatedDevicesByType(devType))
if freeCount < reqCount {
h.Exclude(fmt.Sprintf("IsolatedDevice type %q not enough, request: %d, hostFree: %d", devType, reqCount, freeCount))
return h.GetResult()
}
cap := freeCount / reqCount
if int64(cap) < minCapacity {
minCapacity = int64(cap)
}
}
// check host device by model
devVendorModelRequest := make(map[string]int, 0)
for _, dev := range reqIsoDevs {
if len(dev.Model) != 0 {
devVendorModelRequest[fmt.Sprintf("%s:%s", dev.Vendor, dev.Model)] += 1
}
}
for vendorModel, reqCount := range devVendorModelRequest {
freeCount := len(hc.UnusedIsolatedDevicesByVendorModel(vendorModel))
if freeCount < reqCount {
h.Exclude(fmt.Sprintf("IsolatedDevice vendor:model %q not enough, request: %d, hostFree: %d", vendorModel, reqCount, freeCount))
return h.GetResult()
}
cap := freeCount / reqCount
if int64(cap) < minCapacity {
minCapacity = int64(cap)
}
}
h.SetCapacity(minCapacity)
return h.GetResult()
}
@@ -0,0 +1,55 @@
package guest
import (
"github.com/yunionio/onecloud/pkg/scheduler/algorithm/predicates"
"github.com/yunionio/onecloud/pkg/scheduler/core"
)
// MemoryPredicate filter current resources free memory capacity is meet,
// if it is satisfied to return the size of the memory that
// can carry the scheduling request.
type MemoryPredicate struct {
predicates.BasePredicate
}
func (p *MemoryPredicate) Name() string {
return "host_memory"
}
func (p *MemoryPredicate) Clone() core.FitPredicate {
return &MemoryPredicate{}
}
func (p *MemoryPredicate) PreExecute(u *core.Unit, cs []core.Candidater) (bool, error) {
if u.IsPublicCloudProvider() {
return false, nil
}
data := u.SchedData()
if data.VMEMSize <= 0 {
return false, nil
}
return true, nil
}
func (p *MemoryPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) {
h := predicates.NewPredicateHelper(p, u, c)
d := u.SchedData()
hc, err := h.HostCandidate()
if err != nil {
return false, nil, err
}
useRsvd := h.UseReserved()
freeMemSize := hc.GetFreeMemSize(useRsvd)
reqMemSize := d.VMEMSize
if freeMemSize < reqMemSize {
totalMemSize := hc.GetTotalMemSize(useRsvd)
h.AppendInsufficientResourceError(reqMemSize, totalMemSize, freeMemSize)
}
h.SetCapacity(freeMemSize / reqMemSize)
return h.GetResult()
}
@@ -0,0 +1,33 @@
package guest
import (
"github.com/yunionio/onecloud/pkg/scheduler/algorithm/predicates"
"github.com/yunionio/onecloud/pkg/scheduler/core"
)
// MigratePredicate filters whether the current candidate can be migrated.
type MigratePredicate struct {
predicates.BasePredicate
}
func (p *MigratePredicate) Name() string {
return "host_migrate"
}
func (p *MigratePredicate) Clone() core.FitPredicate {
return &MigratePredicate{}
}
func (p *MigratePredicate) PreExecute(u *core.Unit, cs []core.Candidater) (bool, error) {
return len(u.SchedData().HostID) > 0, nil
}
func (p *MigratePredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) {
h := predicates.NewPredicateHelper(p, u, c)
if u.SchedData().HostID == c.IndexKey() {
h.Exclude(predicates.ErrHostIsSpecifiedForMigration)
}
return h.GetResult()
}
@@ -0,0 +1,40 @@
package guest
import (
"github.com/yunionio/onecloud/pkg/scheduler/algorithm/predicates"
"github.com/yunionio/onecloud/pkg/scheduler/core"
)
// NestPredicate will filter whether the current host is turned on KVM,
//if the scheduling specified settings are inconsistent, then the host
// will be filtered out.
type NestPredicate struct {
predicates.BasePredicate
}
func (p *NestPredicate) Name() string {
return "host_nest"
}
func (p *NestPredicate) Clone() core.FitPredicate {
return &NestPredicate{}
}
func (p *NestPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) {
h := predicates.NewPredicateHelper(p, u, c)
hc, err := h.HostCandidate()
if err != nil {
return false, nil, err
}
d := u.SchedData()
if d.Meta["kvm"] == "enabled" {
if hc.Metadata["nest"] != "enabled" {
h.Exclude(predicates.ErrNotSupportNest)
}
}
return h.GetResult()
}
@@ -0,0 +1,207 @@
package guest
import (
"fmt"
"strings"
"sync"
"github.com/yunionio/onecloud/pkg/scheduler/algorithm/predicates"
"github.com/yunionio/onecloud/pkg/scheduler/api"
"github.com/yunionio/onecloud/pkg/scheduler/core"
networks "github.com/yunionio/onecloud/pkg/scheduler/db/models"
"github.com/yunionio/pkg/utils"
)
// NetworkPredicate will filter the current network information with
// the specified scheduling information to match, if not specified will
// randomly match the available network resources.
type NetworkPredicate struct {
predicates.BasePredicate
SelectedNetworks sync.Map
}
func (p *NetworkPredicate) Name() string {
return "host_network"
}
func (p *NetworkPredicate) Clone() core.FitPredicate {
return &NetworkPredicate{}
}
func (p *NetworkPredicate) PreExecute(u *core.Unit, cs []core.Candidater) (bool, error) {
data := u.SchedData()
if len(data.HostID) > 0 && len(data.Networks) == 0 {
return false, nil
}
return true, nil
}
func (p *NetworkPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) {
h := predicates.NewPredicateHelper(p, u, c)
hc, err := h.HostCandidate()
if err != nil {
return false, nil, err
}
d := u.SchedData()
isMigrate := func() bool {
return len(d.HostID) > 0
}
// ServerType's value is 'guest' or ''(support all type) will return true.
isMatchServerType := func(network *networks.NetworkSchedResult) bool {
return network.ServerType == "guest" || network.ServerType == ""
}
counterOfNetwork := func(u *core.Unit, n *networks.NetworkSchedResult, r int) core.Counter {
counter := u.CounterManager.GetOrCreate("net:"+n.ID, func() core.Counter {
return core.NewNormalCounter(int64(n.Ports - r))
})
u.SharedResourceManager.Add(n.ID, counter)
return counter
}
isRandomNetworkAvailable := func(private bool, exit bool, wire string,
counters core.MultiCounter) string {
var fullErrMsgs []string
found := false
for _, n := range hc.Networks {
errMsgs := []string{}
appendError := func(errMsg string) {
errMsgs = append(errMsgs, errMsg)
}
if !isMatchServerType(n) {
appendError(predicates.ErrServerTypeIsNotMatch)
}
if n.IsExit != exit {
appendError(predicates.ErrExitIsNotMatch)
}
if !(n.Ports > 0 || isMigrate()) {
appendError(predicates.ErrNoPorts)
}
if wire != "" && !utils.HasPrefix(wire, n.Wire) && !utils.HasPrefix(wire, n.WireID) { // re
appendError(predicates.ErrWireIsNotMatch)
}
if !((!private && n.IsPublic) || (private && !n.IsPublic && n.TenantID == d.OwnerTenantID)) {
appendError(predicates.ErrNotOwner)
}
if len(errMsgs) == 0 {
// add resource
reservedNetworks := 0
counter := counterOfNetwork(u, n, reservedNetworks)
p.SelectedNetworks.Store(n.ID, counter.GetCount())
counters.Add(counter)
found = true
if counters.GetCount() >= d.Count {
break
}
} else {
fullErrMsgs = append(fullErrMsgs,
fmt.Sprintf("%s: %s", n.ID, strings.Join(errMsgs, ",")),
)
}
}
if !found {
return strings.Join(fullErrMsgs, "; ")
}
return ""
}
filterByRandomNetwork := func() {
counters := core.NewCounters()
if err_msg := isRandomNetworkAvailable(false, false, "", counters); err_msg != "" {
h.AppendPredicateFailMsg(err_msg)
}
h.SetCapacityCounter(counters)
}
isNetworkAvaliable := func(n *api.Network, counters *core.MinCounters,
networks []*networks.NetworkSchedResult) string {
if n.Idx == "" {
counters0 := core.NewCounters()
ret_msg := isRandomNetworkAvailable(n.Private, n.Exit, n.Wire, counters0)
counters.Add(counters0)
return ret_msg
}
if len(hc.Networks) == 0 {
return predicates.ErrNoAvailableNetwork
}
errMsgs := make([]string, 0)
for _, net := range hc.Networks {
if !isMatchServerType(net) {
errMsgs = append(errMsgs, fmt.Sprintf("%v(%v): server type not matched", net.Name, net.ID))
continue
}
if !(n.Idx == net.ID || n.Idx == net.Name) {
//errMsgs = append(errMsgs, fmt.Sprintf("%v(%v): id/name not matched", net.Name, net.ID))
} else if !(net.IsPublic || net.TenantID == d.OwnerTenantID) {
errMsgs = append(errMsgs, fmt.Sprintf("%v(%v): not owner (%v != %v)", net.Name, net.ID, net.TenantID, d.OwnerTenantID))
} else if !(net.Ports > 0 || isMigrate()) {
errMsgs = append(errMsgs, fmt.Sprintf("%v(%v): ports use up", net.Name, net.ID))
} else {
// add resource
reservedNetworks := 0
counter := counterOfNetwork(u, net, reservedNetworks)
p.SelectedNetworks.Store(net.ID, counter.GetCount())
counters.Add(counter)
return ""
}
}
if len(errMsgs) == 0 {
return predicates.ErrUnknown
}
return strings.Join(errMsgs, "; ")
}
filterBySpecifiedNetworks := func() {
counters := core.NewMinCounters()
var errMsgs []string
for _, n := range d.Networks {
if err_msg := isNetworkAvaliable(n, counters, hc.Networks); err_msg != "" {
errMsgs = append(errMsgs, err_msg)
}
}
if len(errMsgs) > 0 {
h.AppendPredicateFailMsg(strings.Join(errMsgs, ", "))
} else {
h.SetCapacityCounter(counters)
}
}
if len(d.Networks) == 0 {
filterByRandomNetwork()
} else {
filterBySpecifiedNetworks()
}
return h.GetResult()
}
func (p *NetworkPredicate) OnSelect(u *core.Unit, c core.Candidater) bool {
u.SetFiltedData(c.IndexKey(), "networks", p.SelectedNetworks)
return true
}
func (p *NetworkPredicate) OnSelectEnd(u *core.Unit, c core.Candidater, count int64) {
}
@@ -0,0 +1,48 @@
package guest
import (
"github.com/yunionio/onecloud/pkg/scheduler/algorithm/predicates"
"github.com/yunionio/onecloud/pkg/scheduler/core"
)
const (
ExpectedStatus = "running"
ExpectedHostStatus = "online"
ExpectedEnableStatus = "enable"
)
// StatusPredicate is to filter the current state of host is available,
// not available host's capacity will be set to 0 and filtered out.
type StatusPredicate struct {
predicates.BasePredicate
}
func (p *StatusPredicate) Name() string {
return "host_status"
}
func (p *StatusPredicate) Clone() core.FitPredicate {
return &StatusPredicate{}
}
func (p *StatusPredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) {
h := predicates.NewPredicateHelper(p, u, c)
curStatus := h.Get("Status").(string)
curHostStatus := h.Get("HostStatus").(string)
curEnableStatus := h.Get("EnableStatus").(string)
if curStatus != ExpectedStatus {
h.Exclude2("status", curStatus, ExpectedStatus)
}
if curHostStatus != ExpectedHostStatus {
h.Exclude2("host_status", curHostStatus, ExpectedHostStatus)
}
if curEnableStatus != ExpectedEnableStatus {
h.Exclude2("enable_status", curEnableStatus, ExpectedEnableStatus)
}
return h.GetResult()
}
@@ -0,0 +1,134 @@
package guest
import (
"fmt"
"strings"
"github.com/yunionio/onecloud/pkg/scheduler/algorithm/predicates"
"github.com/yunionio/onecloud/pkg/scheduler/core"
"github.com/yunionio/pkg/utils"
)
// StoragePredicate used to filter whether the storage capacity of the
// current candidate matches the type of the disk. If not matched, the
// storage capacity will be set to 0.
type StoragePredicate struct {
predicates.BasePredicate
}
func (p *StoragePredicate) Name() string {
return "host_storage"
}
func (p *StoragePredicate) Clone() core.FitPredicate {
return &StoragePredicate{}
}
func (p *StoragePredicate) PreExecute(u *core.Unit, cs []core.Candidater) (bool, error) {
if u.IsPublicCloudProvider() {
return false, nil
}
return true, nil
}
func (p *StoragePredicate) Execute(u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) {
h := predicates.NewPredicateHelper(p, u, c)
hc, err := h.HostCandidate()
if err != nil {
return false, nil, err
}
d := u.SchedData()
isMigrate := func() bool {
return len(d.HostID) > 0
}
isLocalhostBackend := func(backend string) bool {
return utils.IsLocalStorage(backend)
}
isStorageAccessible := func(storage string) bool {
for _, s := range hc.Storages {
if storage == s.ID || storage == s.Name {
return true
}
}
return false
}
getStorageCapacity := func(backend string, reqMaxSize int64, reqTotalSize int64, useRsvd bool) (int64, int64) {
totalFree := hc.GetFreeStorageSizeOfType(backend, useRsvd)
capacity := totalFree / utils.Max(reqTotalSize, 1)
return capacity, totalFree
}
getReqSizeStr := func(backend string) string {
ss := make([]string, 0, len(d.Disks))
for _, disk := range d.Disks {
if disk.Backend == backend {
ss = append(ss, fmt.Sprintf("%v", disk.Size))
}
}
return strings.Join(ss, "+")
}
getStorageFreeStr := func(backend string, useRsvd bool) string {
ss := []string{}
for _, s := range hc.Storages {
if s.StorageType == backend {
total := int64(float64(s.Capacity) * s.Cmtbound)
free := total - s.UsedCapacity - s.WasteCapacity
ss = append(ss, fmt.Sprintf("(%v-%v-%v=%v)", total, s.UsedCapacity, s.WasteCapacity, free))
}
}
return strings.Join(ss, " + ")
}
sizeRequest := make(map[string]map[string]int64, 0)
storeRequest := make(map[string]int64, 0)
for _, disk := range d.Disks {
if isMigrate() && !isLocalhostBackend(disk.Backend) {
storeRequest[*disk.Storage] = 1
} else {
if _, ok := sizeRequest[disk.Backend]; !ok {
sizeRequest[disk.Backend] = map[string]int64{"max": -1, "total": 0}
}
max := sizeRequest[disk.Backend]["max"]
if max < disk.Size {
sizeRequest[disk.Backend]["max"] = disk.Size
}
sizeRequest[disk.Backend]["total"] += disk.Size
}
}
for store := range storeRequest {
if !isStorageAccessible(store) {
h.Exclude(fmt.Sprintf("storage %v not accessible", store))
return h.GetResult()
}
}
useRsvd := h.UseReserved()
minCapacity := int64(0xFFFFFFFF)
for be, req := range sizeRequest {
capacity, totalFree := getStorageCapacity(be, req["max"], req["total"], useRsvd)
if capacity == 0 {
s := fmt.Sprintf("no enough %q storage, req=%v(%v), free=%v(%v)",
be, req["total"], getReqSizeStr(be), totalFree, getStorageFreeStr(be, useRsvd))
h.AppendPredicateFailMsg(s)
}
minCapacity = utils.Min(minCapacity, capacity)
}
h.SetCapacity(minCapacity)
if minCapacity <= 0 {
h.AppendPredicateFailMsg(predicates.ErrNoEnoughStorage)
}
return h.GetResult()
}
@@ -0,0 +1,166 @@
package predicates
import (
"fmt"
"strings"
"github.com/yunionio/log"
"github.com/yunionio/onecloud/pkg/scheduler/algorithm"
"github.com/yunionio/onecloud/pkg/scheduler/cache/candidate"
"github.com/yunionio/onecloud/pkg/scheduler/core"
"github.com/yunionio/onecloud/pkg/scheduler/data_manager"
)
// BasePredicate is a default struct for all the predicates that will
// include it and implement it's Name() and PreExecute() methods.
type BasePredicate struct{}
func (b *BasePredicate) Name() string {
return "base_predicate_should_not_be_called"
}
func (b *BasePredicate) PreExecute(unit *core.Unit, candis []core.Candidater) (bool, error) {
return true, nil
}
type PredicateHelper struct {
predicate core.FitPredicate
predicateFails []core.PredicateFailureReason
capacity int64
Unit *core.Unit
Candidate core.Candidater
}
func (h *PredicateHelper) getResult() (bool, []core.PredicateFailureReason, error) {
if len(h.predicateFails) > 0 {
return false, h.predicateFails, nil
}
if h.capacity == 0 {
return false, []core.PredicateFailureReason{}, nil
}
return true, nil, nil
}
func (h *PredicateHelper) GetResult() (bool, []core.PredicateFailureReason, error) {
ok, reasons, err := h.getResult()
if !ok {
log.Warningf("[Filter Result] candidate: %q, filter: %q, is_ok: %v, reason: %q, error: %v\n",
h.Candidate.IndexKey(), h.predicate.Name(), ok, getReasonsString(reasons), err)
}
return ok, reasons, err
}
func getReasonsString(reasons []core.PredicateFailureReason) string {
if len(reasons) == 0 {
return ""
}
ss := make([]string, 0, len(reasons))
for _, reason := range reasons {
ss = append(ss, reason.GetReason())
}
return strings.Join(ss, ", ")
}
func NewPredicateHelper(pre core.FitPredicate, unit *core.Unit, candi core.Candidater) *PredicateHelper {
h := &PredicateHelper{
predicate: pre,
capacity: core.EmptyCapacity,
predicateFails: []core.PredicateFailureReason{},
Unit: unit,
Candidate: candi,
}
return h
}
func (h *PredicateHelper) GetFailedResult(err error) (bool, []core.PredicateFailureReason, error) {
return false, nil, err
}
func (h *PredicateHelper) AppendPredicateFail(reason core.PredicateFailureReason) {
h.predicateFails = append(h.predicateFails, reason)
}
func (h *PredicateHelper) AppendPredicateFailMsg(reason string) {
h.AppendPredicateFail(NewUnexceptedResourceError(reason))
}
func (h *PredicateHelper) AppendInsufficientResourceError(req, total, free int64) {
h.AppendPredicateFail(
NewInsufficientResourceError(h.Candidate.Get("Name").(string), req, total, free))
}
// SetCapacity returns the current resource capacity calculated by a filter.
// And 'capacity' default is -1.
func (h *PredicateHelper) SetCapacity(capacity int64) {
if capacity < 0 {
capacity = 0
}
h.SetCapacityCounter(core.NewNormalCounter(capacity))
}
func (h *PredicateHelper) SetCapacityCounter(counter core.Counter) {
capacity := counter.GetCount()
if capacity < core.EmptyCapacity {
capacity = core.EmptyCapacity
}
h.capacity = capacity
h.Unit.SetCapacity(h.Candidate.IndexKey(), h.predicate.Name(), counter)
}
func (h *PredicateHelper) Exclude(reason string) {
h.SetCapacity(0)
h.AppendPredicateFailMsg(reason)
}
func (h *PredicateHelper) Exclude2(predicateName string, current, expected interface{}) {
h.Exclude(fmt.Sprintf("%s is '%v', expected '%v'", predicateName, current, expected))
}
func (h *PredicateHelper) Get(key string) interface{} {
return h.Candidate.Get(key)
}
func (h *PredicateHelper) GetInt64(key string, def int64) int64 {
value := h.Get(key)
if value == nil {
return def
}
return value.(int64)
}
func (h *PredicateHelper) GetGroupCounts() (*data_manager.GroupResAlgorithmResult, error) {
value := h.Get("Groups")
if value == nil {
return nil, nil
}
if r, ok := value.(*data_manager.GroupResAlgorithmResult); ok {
return r, nil
}
return nil, fmt.Errorf("type error: not *data_manager.GroupResAlgorithmResult (GetGroupCounts)")
}
func (h *PredicateHelper) HostCandidate() (*candidate.HostDesc, error) {
return algorithm.ToHostCandidate(h.Candidate)
}
func (h *PredicateHelper) BaremetalCandidate() (*candidate.BaremetalDesc, error) {
return algorithm.ToBaremetalCandidate(h.Candidate)
}
// UseReserved check whether the unit can use guest reserved resource
func (h *PredicateHelper) UseReserved() bool {
usable := false
data := h.Unit.SchedData()
isoDevs := data.IsolatedDevices
if len(isoDevs) > 0 {
usable = true
}
return usable
}