Fix up usage and tests, split into multiple files.

Doing this in multiple commits in an attempt to preserve the file movement history.
This commit is contained in:
Daniel Smith
2014-06-28 15:35:51 -07:00
parent 21e63cf75a
commit 0760e9bc2c
12 changed files with 444 additions and 192 deletions

View File

@@ -14,129 +14,14 @@ See the License for the specific language governing permissions and
limitations under the License.
*/
package registry
package scheduler
import (
"fmt"
"math/rand"
"github.com/GoogleCloudPlatform/kubernetes/pkg/api"
"github.com/GoogleCloudPlatform/kubernetes/pkg/labels"
)
// Anything that can list minions for a scheduler.
type MinionLister interface {
List() (machines []string, err error)
}
// Make a MinionLister from a []string
type StringMinionLister []string
func (s StringMinionLister) List() ([]string, error) {
return []string(s), nil
}
// Scheduler is an interface implemented by things that know how to schedule pods onto machines.
// Scheduler is an interface implemented by things that know how to schedule pods
// onto machines.
type Scheduler interface {
Schedule(api.Pod, MinionLister) (string, error)
}
// RandomScheduler choses machines uniformly at random.
type RandomScheduler struct {
random rand.Rand
}
func MakeRandomScheduler(random rand.Rand) Scheduler {
return &RandomScheduler{
random: random,
}
}
func (s *RandomScheduler) Schedule(pod api.Pod, minionLister MinionLister) (string, error) {
machines, err := minionLister.List()
if err != nil {
return "", err
}
return machines[s.random.Int()%len(machines)], nil
}
// RoundRobinScheduler chooses machines in order.
type RoundRobinScheduler struct {
currentIndex int
}
func MakeRoundRobinScheduler() Scheduler {
return &RoundRobinScheduler{
currentIndex: -1,
}
}
func (s *RoundRobinScheduler) Schedule(pod api.Pod, minionLister MinionLister) (string, error) {
machines, err := minionLister.List()
if err != nil {
return "", err
}
s.currentIndex = (s.currentIndex + 1) % len(machines)
result := machines[s.currentIndex]
return result, nil
}
type FirstFitScheduler struct {
registry PodRegistry
random *rand.Rand
}
func MakeFirstFitScheduler(registry PodRegistry, random *rand.Rand) Scheduler {
return &FirstFitScheduler{
registry: registry,
random: random,
}
}
func (s *FirstFitScheduler) containsPort(pod api.Pod, port api.Port) bool {
for _, container := range pod.DesiredState.Manifest.Containers {
for _, podPort := range container.Ports {
if podPort.HostPort == port.HostPort {
return true
}
}
}
return false
}
func (s *FirstFitScheduler) Schedule(pod api.Pod, minionLister MinionLister) (string, error) {
machines, err := minionLister.List()
if err != nil {
return "", err
}
machineToPods := map[string][]api.Pod{}
pods, err := s.registry.ListPods(labels.Everything())
if err != nil {
return "", err
}
for _, scheduledPod := range pods {
host := scheduledPod.CurrentState.Host
machineToPods[host] = append(machineToPods[host], scheduledPod)
}
var machineOptions []string
for _, machine := range machines {
podFits := true
for _, scheduledPod := range machineToPods[machine] {
for _, container := range pod.DesiredState.Manifest.Containers {
for _, port := range container.Ports {
if s.containsPort(scheduledPod, port) {
podFits = false
}
}
}
}
if podFits {
machineOptions = append(machineOptions, machine)
}
}
if len(machineOptions) == 0 {
return "", fmt.Errorf("failed to find fit for %#v", pod)
} else {
return machineOptions[s.random.Int()%len(machineOptions)], nil
}
Schedule(api.Pod, MinionLister) (selectedMachine string, err error)
}