mirror of
https://github.com/yunionio/cloudpods.git
synced 2026-09-24 16:03:43 +08:00
scheduler: k8s scheduler extender
This commit is contained in:
@@ -0,0 +1,35 @@
|
||||
package k8s
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"k8s.io/api/core/v1"
|
||||
|
||||
"yunion.io/x/onecloud/pkg/scheduler/algorithm/predicates/guest"
|
||||
"yunion.io/x/onecloud/pkg/scheduler/cache/candidate"
|
||||
)
|
||||
|
||||
type HostStatusPredicate struct{}
|
||||
|
||||
func (p *HostStatusPredicate) Clone() IPredicate {
|
||||
return &HostStatusPredicate{}
|
||||
}
|
||||
|
||||
func (p *HostStatusPredicate) Name() string {
|
||||
return "host_status"
|
||||
}
|
||||
|
||||
func (p *HostStatusPredicate) PreExecute(pod *v1.Pod, node *v1.Node, host *candidate.HostDesc) bool {
|
||||
return true
|
||||
}
|
||||
|
||||
func (p *HostStatusPredicate) Execute(pod *v1.Pod, node *v1.Node, host *candidate.HostDesc) (bool, error) {
|
||||
if host.Status != guest.ExpectedStatus {
|
||||
return false, fmt.Errorf("Host status is %s", host.Status)
|
||||
}
|
||||
|
||||
if !host.Enabled {
|
||||
return false, fmt.Errorf("Host is disabled")
|
||||
}
|
||||
return true, nil
|
||||
}
|
||||
@@ -0,0 +1,76 @@
|
||||
package k8s
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"k8s.io/api/core/v1"
|
||||
|
||||
"yunion.io/x/onecloud/pkg/scheduler/cache/candidate"
|
||||
)
|
||||
|
||||
var PredicatesManager *SPredicatesManager
|
||||
|
||||
func init() {
|
||||
PredicatesManager = newPredicatesManager()
|
||||
|
||||
PredicatesManager.Register(
|
||||
&HostStatusPredicate{},
|
||||
&NetworkPredicate{},
|
||||
)
|
||||
}
|
||||
|
||||
type IPredicate interface {
|
||||
Name() string
|
||||
Clone() IPredicate
|
||||
PreExecute(pod *v1.Pod, node *v1.Node, host *candidate.HostDesc) bool
|
||||
Execute(pod *v1.Pod, node *v1.Node, host *candidate.HostDesc) (bool, error)
|
||||
}
|
||||
|
||||
type SPredicatesManager struct {
|
||||
predicates []IPredicate
|
||||
}
|
||||
|
||||
func newPredicatesManager() *SPredicatesManager {
|
||||
man := &SPredicatesManager{
|
||||
predicates: make([]IPredicate, 0),
|
||||
}
|
||||
return man
|
||||
}
|
||||
|
||||
func (man *SPredicatesManager) Register(pres ...IPredicate) *SPredicatesManager {
|
||||
for _, pre := range pres {
|
||||
if !man.Has(pre) {
|
||||
man.predicates = append(man.predicates, pre)
|
||||
}
|
||||
}
|
||||
return man
|
||||
}
|
||||
|
||||
func (man *SPredicatesManager) Has(newPre IPredicate) bool {
|
||||
if len(man.predicates) == 0 {
|
||||
return false
|
||||
}
|
||||
for _, pre := range man.predicates {
|
||||
if pre.Name() == newPre.Name() {
|
||||
return true
|
||||
}
|
||||
}
|
||||
return false
|
||||
}
|
||||
|
||||
func (man *SPredicatesManager) DoFilter(pod *v1.Pod, node *v1.Node, host *candidate.HostDesc) (bool, error) {
|
||||
for _, pre := range man.predicates {
|
||||
tmpPre := pre.Clone()
|
||||
if !tmpPre.PreExecute(pod, node, host) {
|
||||
continue
|
||||
}
|
||||
fit, err := tmpPre.Execute(pod, node, host)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
if !fit {
|
||||
return false, fmt.Errorf("Filtered by %s", tmpPre.Name())
|
||||
}
|
||||
}
|
||||
return true, nil
|
||||
}
|
||||
@@ -0,0 +1,105 @@
|
||||
package k8s
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
|
||||
"k8s.io/api/core/v1"
|
||||
|
||||
"yunion.io/x/pkg/util/errors"
|
||||
|
||||
"yunion.io/x/onecloud/pkg/scheduler/cache/candidate"
|
||||
"yunion.io/x/onecloud/pkg/scheduler/db/models"
|
||||
)
|
||||
|
||||
const (
|
||||
// k8s annotations for create pod
|
||||
YUNION_CNI_NETWORK_ANNOTATION = "cni.yunion.io/network"
|
||||
YUNION_CNI_IPADDR_ANNOTATION = "cni.yunion.io/ip"
|
||||
)
|
||||
|
||||
type NetworkPredicate struct {
|
||||
network string
|
||||
ipAddr string
|
||||
}
|
||||
|
||||
func (p *NetworkPredicate) Clone() IPredicate {
|
||||
return &NetworkPredicate{}
|
||||
}
|
||||
|
||||
func (p *NetworkPredicate) Name() string {
|
||||
return "network"
|
||||
}
|
||||
|
||||
func (p *NetworkPredicate) PreExecute(pod *v1.Pod, node *v1.Node, host *candidate.HostDesc) bool {
|
||||
net, netCont := pod.Annotations[YUNION_CNI_NETWORK_ANNOTATION]
|
||||
ipAddr, ipCont := pod.Annotations[YUNION_CNI_IPADDR_ANNOTATION]
|
||||
p.network = net
|
||||
p.ipAddr = ipAddr
|
||||
return netCont || ipCont
|
||||
}
|
||||
|
||||
func (p *NetworkPredicate) Execute(pod *v1.Pod, node *v1.Node, host *candidate.HostDesc) (bool, error) {
|
||||
hostNets := host.Networks
|
||||
if p.network != "" {
|
||||
err := p.checkByNetworks(hostNets)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
}
|
||||
if p.ipAddr != "" {
|
||||
err := p.checkNetworksIP(p.ipAddr, hostNets)
|
||||
if err != nil {
|
||||
return false, err
|
||||
}
|
||||
}
|
||||
return true, nil
|
||||
}
|
||||
|
||||
func (p NetworkPredicate) checkByNetworks(nets []*models.NetworkSchedResult) error {
|
||||
if len(nets) == 0 {
|
||||
return fmt.Errorf("Network is empty")
|
||||
}
|
||||
errs := make([]error, 0)
|
||||
for _, net := range nets {
|
||||
err := p.checkByNetwork(net)
|
||||
if err == nil {
|
||||
return nil
|
||||
}
|
||||
errs = append(errs, err)
|
||||
}
|
||||
return errors.NewAggregate(errs)
|
||||
}
|
||||
|
||||
func (p NetworkPredicate) checkByNetwork(net *models.NetworkSchedResult) error {
|
||||
if net.Ports <= 0 {
|
||||
return fmt.Errorf("Network %s no free IPs", net.Name)
|
||||
}
|
||||
if !(p.network == net.Name || p.network == net.ID) {
|
||||
return fmt.Errorf("Network %s:%s or id not match %s", net.Name, net.ID, p.network)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (p NetworkPredicate) checkNetworksIP(ip string, nets []*models.NetworkSchedResult) error {
|
||||
if len(nets) == 0 {
|
||||
return fmt.Errorf("Network is empty")
|
||||
}
|
||||
errs := make([]error, 0)
|
||||
for _, net := range nets {
|
||||
err := p.checkNetworkIP(ip, net)
|
||||
if err == nil {
|
||||
return nil
|
||||
}
|
||||
errs = append(errs, err)
|
||||
}
|
||||
return errors.NewAggregate(errs)
|
||||
}
|
||||
|
||||
func (p NetworkPredicate) checkNetworkIP(ip string, net *models.NetworkSchedResult) error {
|
||||
if ok, err := net.ContainsIp(ip); err != nil {
|
||||
return err
|
||||
} else if !ok {
|
||||
return fmt.Errorf("Network %s not contains ip %s", net.Name, ip)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
Reference in New Issue
Block a user