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 (
v1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/labels"
"k8s.io/klog/v2"
k8sframework "k8s.io/kubernetes/pkg/scheduler/framework"
"volcano.sh/volcano/pkg/scheduler/api"
)
type PodFilter func(*v1.Pod) bool
type PodsLister interface {
List(labels.Selector) ([]*v1.Pod, error)
FilteredList(podFilter PodFilter, selector labels.Selector) ([]*v1.Pod, error)
}
type PodLister struct {
Session *Session
CachedPods map[api.TaskID]*v1.Pod
Tasks map[api.TaskID]*api.TaskInfo
TaskWithAffinity map[api.TaskID]*api.TaskInfo
}
type PodAffinityLister struct {
pl *PodLister
}
func HaveAffinity(pod *v1.Pod) bool {
affinity := pod.Spec.Affinity
return affinity != nil &&
(affinity.NodeAffinity != nil ||
affinity.PodAffinity != nil ||
affinity.PodAntiAffinity != nil)
}
func NewPodLister(ssn *Session) *PodLister {
pl := &PodLister{
Session: ssn,
CachedPods: make(map[api.TaskID]*v1.Pod),
Tasks: make(map[api.TaskID]*api.TaskInfo),
TaskWithAffinity: make(map[api.TaskID]*api.TaskInfo),
}
for _, job := range pl.Session.Jobs {
for status, tasks := range job.TaskStatusIndex {
if !api.AllocatedStatus(status) {
continue
}
for _, task := range tasks {
pl.Tasks[task.UID] = task
pod := pl.copyTaskPod(task)
pl.CachedPods[task.UID] = pod
if HaveAffinity(task.Pod) {
pl.TaskWithAffinity[task.UID] = task
}
}
}
}
return pl
}
func NewPodListerFromNode(ssn *Session) *PodLister {
pl := &PodLister{
Session: ssn,
CachedPods: make(map[api.TaskID]*v1.Pod),
Tasks: make(map[api.TaskID]*api.TaskInfo),
TaskWithAffinity: make(map[api.TaskID]*api.TaskInfo),
}
for _, node := range pl.Session.Nodes {
for _, task := range node.Tasks {
if !api.AllocatedStatus(task.Status) && task.Status != api.Releasing {
continue
}
pl.Tasks[task.UID] = task
pod := pl.copyTaskPod(task)
pl.CachedPods[task.UID] = pod
if HaveAffinity(task.Pod) {
pl.TaskWithAffinity[task.UID] = task
}
}
}
return pl
}
func (pl *PodLister) copyTaskPod(task *api.TaskInfo) *v1.Pod {
pod := task.Pod.DeepCopy()
pod.Spec.NodeName = task.NodeName
return pod
}
func (pl *PodLister) GetPod(task *api.TaskInfo) *v1.Pod {
if task.NodeName == task.Pod.Spec.NodeName {
return task.Pod
}
pod, found := pl.CachedPods[task.UID]
if !found {
pod = pl.copyTaskPod(task)
klog.Warningf("DeepCopy for pod %s/%s at PodLister.GetPod is unexpected", pod.Namespace, pod.Name)
}
return pod
}
func (pl *PodLister) UpdateTask(task *api.TaskInfo, nodeName string) *v1.Pod {
pod, found := pl.CachedPods[task.UID]
if !found {
pod = pl.copyTaskPod(task)
pl.CachedPods[task.UID] = pod
}
pod.Spec.NodeName = nodeName
if !api.AllocatedStatus(task.Status) {
delete(pl.Tasks, task.UID)
if HaveAffinity(task.Pod) {
delete(pl.TaskWithAffinity, task.UID)
}
} else {
pl.Tasks[task.UID] = task
if HaveAffinity(task.Pod) {
pl.TaskWithAffinity[task.UID] = task
}
}
return pod
}
func (pl *PodLister) List(selector labels.Selector) ([]*v1.Pod, error) {
var pods []*v1.Pod
for _, task := range pl.Tasks {
pod := pl.GetPod(task)
if selector.Matches(labels.Set(pod.Labels)) {
pods = append(pods, pod)
}
}
return pods, nil
}
func (pl *PodLister) filteredListWithTaskSet(taskSet map[api.TaskID]*api.TaskInfo, podFilter PodFilter, selector labels.Selector) ([]*v1.Pod, error) {
var pods []*v1.Pod
for _, task := range taskSet {
pod := pl.GetPod(task)
if podFilter(pod) && selector.Matches(labels.Set(pod.Labels)) {
pods = append(pods, pod)
}
}
return pods, nil
}
func (pl *PodLister) FilteredList(podFilter PodFilter, selector labels.Selector) ([]*v1.Pod, error) {
return pl.filteredListWithTaskSet(pl.Tasks, podFilter, selector)
}
func (pl *PodLister) AffinityFilteredList(podFilter PodFilter, selector labels.Selector) ([]*v1.Pod, error) {
return pl.filteredListWithTaskSet(pl.TaskWithAffinity, podFilter, selector)
}
func (pl *PodLister) AffinityLister() *PodAffinityLister {
pal := &PodAffinityLister{
pl: pl,
}
return pal
}
func (pal *PodAffinityLister) List(selector labels.Selector) ([]*v1.Pod, error) {
return pal.pl.List(selector)
}
func (pal *PodAffinityLister) FilteredList(podFilter PodFilter, selector labels.Selector) ([]*v1.Pod, error) {
return pal.pl.AffinityFilteredList(podFilter, selector)
}
func GenerateNodeMapAndSlice(nodes map[string]*api.NodeInfo) map[string]*k8sframework.NodeInfo {
nodeMap := make(map[string]*k8sframework.NodeInfo)
for _, node := range nodes {
nodeInfo := k8sframework.NewNodeInfo(node.Pods()...)
nodeInfo.SetNode(node.Node)
nodeMap[node.Name] = nodeInfo
nodeMap[node.Name].ImageStates = node.CloneImageSummary()
}
return nodeMap
}
type CachedNodeInfo struct {
Session *Session
}
func (c *CachedNodeInfo) GetNodeInfo(name string) (*v1.Node, error) {
node, found := c.Session.Nodes[name]
if !found {
return nil, errors.NewNotFound(v1.Resource("node"), name)
}
return node.Node, nil
}
type NodeLister struct {
Session *Session
}
func (nl *NodeLister) List() ([]*v1.Node, error) {
var nodes []*v1.Node
for _, node := range nl.Session.Nodes {
nodes = append(nodes, node.Node)
}
return nodes, nil
}