Copyright 2018 The Kubernetes Authors.
Licensed under the Apache License, Version 2.0 (the "License");
you may not use this file except in compliance with the License.
You may obtain a copy of the License at
http://www.apache.org/licenses/LICENSE-2.0
Unless required by applicable law or agreed to in writing, software
distributed under the License is distributed on an "AS IS" BASIS,
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
See the License for the specific language governing permissions and
limitations under the License.
*/
package framework
import (
k8sframework "k8s.io/kubernetes/pkg/scheduler/framework"
"volcano.sh/apis/pkg/apis/scheduling"
"volcano.sh/volcano/pkg/controllers/job/helpers"
"volcano.sh/volcano/pkg/scheduler/api"
"volcano.sh/volcano/pkg/scheduler/util"
)
func (ssn *Session) AddJobOrderFn(name string, cf api.CompareFn) {
ssn.jobOrderFns[name] = cf
}
func (ssn *Session) AddQueueOrderFn(name string, qf api.CompareFn) {
ssn.queueOrderFns[name] = qf
}
func (ssn *Session) AddClusterOrderFn(name string, qf api.CompareFn) {
ssn.clusterOrderFns[name] = qf
}
func (ssn *Session) AddTaskOrderFn(name string, cf api.CompareFn) {
ssn.taskOrderFns[name] = cf
}
func (ssn *Session) AddPreemptableFn(name string, cf api.EvictableFn) {
ssn.preemptableFns[name] = cf
}
func (ssn *Session) AddReclaimableFn(name string, rf api.EvictableFn) {
ssn.reclaimableFns[name] = rf
}
func (ssn *Session) AddJobReadyFn(name string, vf api.ValidateFn) {
ssn.jobReadyFns[name] = vf
}
func (ssn *Session) AddJobPipelinedFn(name string, vf api.VoteFn) {
ssn.jobPipelinedFns[name] = vf
}
func (ssn *Session) AddPredicateFn(name string, pf api.PredicateFn) {
ssn.predicateFns[name] = pf
}
func (ssn *Session) AddPrePredicateFn(name string, pf api.PrePredicateFn) {
ssn.prePredicateFns[name] = pf
}
func (ssn *Session) AddBestNodeFn(name string, pf api.BestNodeFn) {
ssn.bestNodeFns[name] = pf
}
func (ssn *Session) AddNodeOrderFn(name string, pf api.NodeOrderFn) {
ssn.nodeOrderFns[name] = pf
}
func (ssn *Session) AddBatchNodeOrderFn(name string, pf api.BatchNodeOrderFn) {
ssn.batchNodeOrderFns[name] = pf
}
func (ssn *Session) AddNodeMapFn(name string, pf api.NodeMapFn) {
ssn.nodeMapFns[name] = pf
}
func (ssn *Session) AddNodeReduceFn(name string, pf api.NodeReduceFn) {
ssn.nodeReduceFns[name] = pf
}
func (ssn *Session) AddOverusedFn(name string, fn api.ValidateFn) {
ssn.overusedFns[name] = fn
}
func (ssn *Session) AddPreemptiveFn(name string, fn api.ValidateWithCandidateFn) {
ssn.preemptiveFns[name] = fn
}
func (ssn *Session) AddAllocatableFn(name string, fn api.AllocatableFn) {
ssn.allocatableFns[name] = fn
}
func (ssn *Session) AddJobValidFn(name string, fn api.ValidateExFn) {
ssn.jobValidFns[name] = fn
}
func (ssn *Session) AddJobEnqueueableFn(name string, fn api.VoteFn) {
ssn.jobEnqueueableFns[name] = fn
}
func (ssn *Session) AddJobEnqueuedFn(name string, fn api.JobEnqueuedFn) {
ssn.jobEnqueuedFns[name] = fn
}
func (ssn *Session) AddTargetJobFn(name string, fn api.TargetJobFn) {
ssn.targetJobFns[name] = fn
}
func (ssn *Session) AddReservedNodesFn(name string, fn api.ReservedNodesFn) {
ssn.reservedNodesFns[name] = fn
}
func (ssn *Session) AddVictimTasksFns(name string, fns []api.VictimTasksFn) {
ssn.victimTasksFns[name] = fns
}
func (ssn *Session) AddJobStarvingFns(name string, fn api.ValidateFn) {
ssn.jobStarvingFns[name] = fn
}
func (ssn *Session) Reclaimable(reclaimer *api.TaskInfo, reclaimees []*api.TaskInfo) []*api.TaskInfo {
var victims []*api.TaskInfo
var init bool
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
if !isEnabled(plugin.EnabledReclaimable) {
continue
}
rf, found := ssn.reclaimableFns[plugin.Name]
if !found {
continue
}
candidates, abstain := rf(reclaimer, reclaimees)
if abstain == 0 {
continue
}
if len(candidates) == 0 {
victims = nil
break
}
if !init {
victims = candidates
init = true
} else {
var intersection []*api.TaskInfo
for _, v := range victims {
for _, c := range candidates {
if v.UID == c.UID {
intersection = append(intersection, v)
}
}
}
victims = intersection
}
}
if victims != nil {
return victims
}
}
return victims
}
func (ssn *Session) Preemptable(preemptor *api.TaskInfo, preemptees []*api.TaskInfo) []*api.TaskInfo {
var victims []*api.TaskInfo
var init bool
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
if !isEnabled(plugin.EnabledPreemptable) {
continue
}
pf, found := ssn.preemptableFns[plugin.Name]
if !found {
continue
}
candidates, abstain := pf(preemptor, preemptees)
if abstain == 0 {
continue
}
if len(candidates) == 0 {
victims = nil
break
}
if !init {
victims = candidates
init = true
} else {
var intersection []*api.TaskInfo
for _, v := range victims {
for _, c := range candidates {
if v.UID == c.UID {
intersection = append(intersection, v)
}
}
}
victims = intersection
}
}
if victims != nil {
return victims
}
}
return victims
}
func (ssn *Session) Overused(queue *api.QueueInfo) bool {
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
if !isEnabled(plugin.EnabledOverused) {
continue
}
of, found := ssn.overusedFns[plugin.Name]
if !found {
continue
}
if of(queue) {
return true
}
}
}
return false
}
func (ssn *Session) Preemptive(queue *api.QueueInfo, candidate *api.TaskInfo) bool {
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
of, found := ssn.preemptiveFns[plugin.Name]
if !isEnabled(plugin.EnablePreemptive) {
continue
}
if !found {
continue
}
if !of(queue, candidate) {
return false
}
}
}
return true
}
func (ssn *Session) Allocatable(queue *api.QueueInfo, candidate *api.TaskInfo) bool {
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
if !isEnabled(plugin.EnabledAllocatable) {
continue
}
af, found := ssn.allocatableFns[plugin.Name]
if !found {
continue
}
if !af(queue, candidate) {
return false
}
}
}
return true
}
func (ssn *Session) JobReady(obj interface{}) bool {
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
if !isEnabled(plugin.EnabledJobReady) {
continue
}
jrf, found := ssn.jobReadyFns[plugin.Name]
if !found {
continue
}
if !jrf(obj) {
return false
}
}
}
return true
}
func (ssn *Session) JobPipelined(obj interface{}) bool {
var hasFound bool
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
if !isEnabled(plugin.EnabledJobPipelined) {
continue
}
jrf, found := ssn.jobPipelinedFns[plugin.Name]
if !found {
continue
}
res := jrf(obj)
if res < 0 {
return false
}
if res > 0 {
hasFound = true
}
}
if hasFound {
return true
}
}
return true
}
func (ssn *Session) JobStarving(obj interface{}) bool {
var hasFound bool
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
if !isEnabled(plugin.EnabledJobStarving) {
continue
}
jrf, found := ssn.jobStarvingFns[plugin.Name]
if !found {
continue
}
hasFound = true
if !jrf(obj) {
return false
}
}
if hasFound {
return true
}
}
return false
}
func (ssn *Session) JobValid(obj interface{}) *api.ValidateResult {
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
jrf, found := ssn.jobValidFns[plugin.Name]
if !found {
continue
}
if vr := jrf(obj); vr != nil && !vr.Pass {
return vr
}
}
}
return nil
}
func (ssn *Session) JobEnqueueable(obj interface{}) bool {
var hasFound bool
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
if !isEnabled(plugin.EnabledJobEnqueued) {
continue
}
fn, found := ssn.jobEnqueueableFns[plugin.Name]
if !found {
continue
}
res := fn(obj)
if res < 0 {
return false
}
if res > 0 {
hasFound = true
}
}
if hasFound {
return true
}
}
return true
}
func (ssn *Session) JobEnqueued(obj interface{}) {
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
if !isEnabled(plugin.EnabledJobEnqueued) {
continue
}
fn, found := ssn.jobEnqueuedFns[plugin.Name]
if !found {
continue
}
fn(obj)
}
}
}
func (ssn *Session) TargetJob(jobs []*api.JobInfo) *api.JobInfo {
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
if !isEnabled(plugin.EnabledTargetJob) {
continue
}
fn, found := ssn.targetJobFns[plugin.Name]
if !found {
continue
}
return fn(jobs)
}
}
return nil
}
func (ssn *Session) VictimTasks(tasks []*api.TaskInfo) map[*api.TaskInfo]bool {
victimSet := make(map[*api.TaskInfo]bool)
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
if !isEnabled(plugin.EnabledVictim) {
continue
}
fns, found := ssn.victimTasksFns[plugin.Name]
if !found {
continue
}
for _, fn := range fns {
victimTasks := fn(tasks)
for _, victim := range victimTasks {
victimSet[victim] = true
}
}
}
if len(victimSet) > 0 {
return victimSet
}
}
return victimSet
}
func (ssn *Session) ReservedNodes() {
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
if !isEnabled(plugin.EnabledReservedNodes) {
continue
}
fn, found := ssn.reservedNodesFns[plugin.Name]
if !found {
continue
}
fn()
}
}
}
func (ssn *Session) JobOrderFn(l, r interface{}) bool {
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
if !isEnabled(plugin.EnabledJobOrder) {
continue
}
jof, found := ssn.jobOrderFns[plugin.Name]
if !found {
continue
}
if j := jof(l, r); j != 0 {
return j < 0
}
}
}
lv := l.(*api.JobInfo)
rv := r.(*api.JobInfo)
if lv.CreationTimestamp.Equal(&rv.CreationTimestamp) {
return lv.UID < rv.UID
}
return lv.CreationTimestamp.Before(&rv.CreationTimestamp)
}
func (ssn *Session) ClusterOrderFn(l, r interface{}) bool {
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
if !isEnabled(plugin.EnabledClusterOrder) {
continue
}
cof, found := ssn.clusterOrderFns[plugin.Name]
if !found {
continue
}
if j := cof(l, r); j != 0 {
return j < 0
}
}
}
lv := l.(*scheduling.Cluster)
rv := r.(*scheduling.Cluster)
return lv.Name < rv.Name
}
func (ssn *Session) QueueOrderFn(l, r interface{}) bool {
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
if !isEnabled(plugin.EnabledQueueOrder) {
continue
}
qof, found := ssn.queueOrderFns[plugin.Name]
if !found {
continue
}
if j := qof(l, r); j != 0 {
return j < 0
}
}
}
lv := l.(*api.QueueInfo)
rv := r.(*api.QueueInfo)
if lv.Queue.CreationTimestamp.Equal(&rv.Queue.CreationTimestamp) {
return lv.UID < rv.UID
}
return lv.Queue.CreationTimestamp.Before(&rv.Queue.CreationTimestamp)
}
func (ssn *Session) TaskCompareFns(l, r interface{}) int {
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
if !isEnabled(plugin.EnabledTaskOrder) {
continue
}
tof, found := ssn.taskOrderFns[plugin.Name]
if !found {
continue
}
if j := tof(l, r); j != 0 {
return j
}
}
}
return 0
}
func (ssn *Session) TaskOrderFn(l, r interface{}) bool {
if res := ssn.TaskCompareFns(l, r); res != 0 {
return res < 0
}
lv := l.(*api.TaskInfo)
rv := r.(*api.TaskInfo)
return helpers.CompareTask(lv, rv)
}
func (ssn *Session) PredicateFn(task *api.TaskInfo, node *api.NodeInfo) error {
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
if !isEnabled(plugin.EnabledPredicate) {
continue
}
pfn, found := ssn.predicateFns[plugin.Name]
if !found {
continue
}
err := pfn(task, node)
if err != nil {
return err
}
}
}
return nil
}
func (ssn *Session) PrePredicateFn(task *api.TaskInfo) error {
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
if !isEnabled(plugin.EnabledPredicate) {
continue
}
pfn, found := ssn.prePredicateFns[plugin.Name]
if !found {
continue
}
err := pfn(task)
if err != nil {
return err
}
}
}
return nil
}
func (ssn *Session) BestNodeFn(task *api.TaskInfo, nodeScores map[float64][]*api.NodeInfo) *api.NodeInfo {
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
if !isEnabled(plugin.EnabledBestNode) {
continue
}
pfn, found := ssn.bestNodeFns[plugin.Name]
if !found {
continue
}
if bestNode := pfn(task, nodeScores); bestNode != nil {
return bestNode
}
}
}
return nil
}
func (ssn *Session) NodeOrderFn(task *api.TaskInfo, node *api.NodeInfo) (float64, error) {
priorityScore := 0.0
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
if !isEnabled(plugin.EnabledNodeOrder) {
continue
}
pfn, found := ssn.nodeOrderFns[plugin.Name]
if !found {
continue
}
score, err := pfn(task, node)
if err != nil {
return 0, err
}
priorityScore += score
}
}
return priorityScore, nil
}
func (ssn *Session) BatchNodeOrderFn(task *api.TaskInfo, nodes []*api.NodeInfo) (map[string]float64, error) {
priorityScore := make(map[string]float64, len(nodes))
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
if !isEnabled(plugin.EnabledNodeOrder) {
continue
}
pfn, found := ssn.batchNodeOrderFns[plugin.Name]
if !found {
continue
}
score, err := pfn(task, nodes)
if err != nil {
return nil, err
}
for nodeName, score := range score {
priorityScore[nodeName] += score
}
}
}
return priorityScore, nil
}
func isEnabled(enabled *bool) bool {
return enabled != nil && *enabled
}
func (ssn *Session) NodeOrderMapFn(task *api.TaskInfo, node *api.NodeInfo) (map[string]float64, float64, error) {
nodeScoreMap := map[string]float64{}
var priorityScore float64
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
if !isEnabled(plugin.EnabledNodeOrder) {
continue
}
if pfn, found := ssn.nodeOrderFns[plugin.Name]; found {
score, err := pfn(task, node)
if err != nil {
return nodeScoreMap, priorityScore, err
}
priorityScore += score
}
if pfn, found := ssn.nodeMapFns[plugin.Name]; found {
score, err := pfn(task, node)
if err != nil {
return nodeScoreMap, priorityScore, err
}
nodeScoreMap[plugin.Name] = score
}
}
}
return nodeScoreMap, priorityScore, nil
}
func (ssn *Session) NodeOrderReduceFn(task *api.TaskInfo, pluginNodeScoreMap map[string]k8sframework.NodeScoreList) (map[string]float64, error) {
nodeScoreMap := map[string]float64{}
for _, tier := range ssn.Tiers {
for _, plugin := range tier.Plugins {
if !isEnabled(plugin.EnabledNodeOrder) {
continue
}
pfn, found := ssn.nodeReduceFns[plugin.Name]
if !found {
continue
}
if err := pfn(task, pluginNodeScoreMap[plugin.Name]); err != nil {
return nodeScoreMap, err
}
for _, hp := range pluginNodeScoreMap[plugin.Name] {
nodeScoreMap[hp.Name] += float64(hp.Score)
}
}
}
return nodeScoreMap, nil
}
func (ssn *Session) BuildVictimsPriorityQueue(victims []*api.TaskInfo) *util.PriorityQueue {
victimsQueue := util.NewPriorityQueue(func(l, r interface{}) bool {
lv := l.(*api.TaskInfo)
rv := r.(*api.TaskInfo)
if lv.Job == rv.Job {
return !ssn.TaskOrderFn(l, r)
}
lvJob, lvJobFound := ssn.Jobs[lv.Job]
rvJob, rvJobFound := ssn.Jobs[rv.Job]
if lvJobFound && rvJobFound && lvJob.Queue != rvJob.Queue {
return !ssn.QueueOrderFn(ssn.Queues[lvJob.Queue], ssn.Queues[rvJob.Queue])
}
return !ssn.JobOrderFn(lvJob, rvJob)
})
for _, victim := range victims {
victimsQueue.Push(victim)
}
return victimsQueue
}