diff --git a/pkg/apis/compute/api.go b/pkg/apis/compute/api.go index 6b9f029a8e..e7908a9f2d 100644 --- a/pkg/apis/compute/api.go +++ b/pkg/apis/compute/api.go @@ -455,6 +455,9 @@ type ServerConfigs struct { // required: false Schedtags []*SchedtagConfig `json:"schedtags"` + // 宿主机路径调度约束,仅检查 auto_create=false 的 host_path + HostPathRequirements []apis.HostPathRequirement `json:"host_path_requirements,omitempty"` + // 透传设备列表 // required: false IsolatedDevices []*IsolatedDeviceConfig `json:"isolated_devices"` diff --git a/pkg/apis/container.go b/pkg/apis/container.go index 21dbaed2cb..a5952d7f70 100644 --- a/pkg/apis/container.go +++ b/pkg/apis/container.go @@ -378,6 +378,11 @@ const ( CONTAINER_VOLUME_MOUNT_HOST_PATH_TYPE_FILE ContainerVolumeMountHostPathType = "file" ) +type HostPathRequirement struct { + Path string `json:"path"` + Type ContainerVolumeMountHostPathType `json:"type"` +} + type ContainerVolumeMountHostPathAutoCreateConfig struct { Uid uint `json:"uid"` Gid uint `json:"gid"` diff --git a/pkg/apis/host/host_path.go b/pkg/apis/host/host_path.go new file mode 100644 index 0000000000..ba04d0fa99 --- /dev/null +++ b/pkg/apis/host/host_path.go @@ -0,0 +1,38 @@ +// Copyright 2019 Yunion +// +// 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 host + +import "yunion.io/x/onecloud/pkg/apis" + +type HostPathCheckItem struct { + Path string `json:"path"` + Type apis.ContainerVolumeMountHostPathType `json:"type"` +} + +type HostPathCheckInput struct { + Paths []HostPathCheckItem `json:"paths"` +} + +type HostPathCheckResult struct { + Path string `json:"path"` + Type apis.ContainerVolumeMountHostPathType `json:"type"` + Exists bool `json:"exists"` + TypeMatched bool `json:"type_matched"` + Error string `json:"error,omitempty"` +} + +type HostPathCheckOutput struct { + Results []HostPathCheckResult `json:"results"` +} diff --git a/pkg/apis/scheduler/api.go b/pkg/apis/scheduler/api.go index a1f5c6167f..b9b59a491d 100644 --- a/pkg/apis/scheduler/api.go +++ b/pkg/apis/scheduler/api.go @@ -44,6 +44,8 @@ type GroupRelation struct { Scope string `json:"scope"` } +type HostPathRequirement = apis.HostPathRequirement + type ServerConfig struct { *compute.ServerConfigs @@ -103,6 +105,8 @@ type ScheduleInput struct { // we don't need reallocate network ReuseNetwork bool `json:"reuse_network"` + HostPathRequirements []HostPathRequirement `json:"host_path_requirements,omitempty"` + // Change config ChangeConfig bool // guest who change config has isolated device diff --git a/pkg/cloudcommon/cmdline/helper.go b/pkg/cloudcommon/cmdline/helper.go index 25e0676622..49d883b77d 100644 --- a/pkg/cloudcommon/cmdline/helper.go +++ b/pkg/cloudcommon/cmdline/helper.go @@ -21,6 +21,7 @@ import ( "yunion.io/x/pkg/util/fileutils" "yunion.io/x/pkg/util/regutils" + "yunion.io/x/onecloud/pkg/apis" "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/apis/scheduler" ) @@ -293,6 +294,9 @@ func FetchServerConfigsByJSON(obj jsonutils.JSONObject) (*compute.ServerConfigs, if err := obj.Unmarshal(conf); err != nil { return nil, err } + if len(conf.HostPathRequirements) == 0 { + conf.HostPathRequirements = fetchHostPathRequirementsByJSON(obj) + } if instanceType, _ := obj.GetString("sku"); instanceType != "" { conf.InstanceType = instanceType @@ -323,6 +327,7 @@ func FetchServerConfigsByJSON(obj jsonutils.JSONObject) (*compute.ServerConfigs, } func FetchScheduleInputByJSON(obj jsonutils.JSONObject) (*scheduler.ScheduleInput, error) { + rawObj := obj input := new(scheduler.ScheduleInput) err := obj.Unmarshal(input) if err != nil { @@ -343,9 +348,57 @@ func FetchScheduleInputByJSON(obj jsonutils.JSONObject) (*scheduler.ScheduleInpu if input.Project == "" { input.Project = jsonutils.GetAnyString(obj, []string{"project", "project_id"}) } + if len(input.HostPathRequirements) == 0 { + input.HostPathRequirements = fetchHostPathRequirementsByJSON(rawObj) + } return input, nil } +func fetchHostPathRequirementsByJSON(obj jsonutils.JSONObject) []apis.HostPathRequirement { + if obj == nil { + return nil + } + + type podContainer struct { + VolumeMounts []*apis.ContainerVolumeMount `json:"volume_mounts"` + } + type podConfig struct { + Containers []podContainer `json:"containers"` + } + type scheduleHostPathInput struct { + Pod *podConfig `json:"pod"` + } + + input := new(scheduleHostPathInput) + if err := obj.Unmarshal(input); err != nil || input.Pod == nil { + return nil + } + + reqs := make([]apis.HostPathRequirement, 0) + seen := make(map[string]struct{}) + for _, container := range input.Pod.Containers { + for _, volumeMount := range container.VolumeMounts { + if volumeMount == nil || volumeMount.Type != apis.CONTAINER_VOLUME_MOUNT_TYPE_HOST_PATH || volumeMount.HostPath == nil { + continue + } + hostPath := volumeMount.HostPath + if hostPath.AutoCreate || hostPath.Path == "" { + continue + } + key := fmt.Sprintf("%s\x00%s", hostPath.Path, hostPath.Type) + if _, ok := seen[key]; ok { + continue + } + seen[key] = struct{}{} + reqs = append(reqs, apis.HostPathRequirement{ + Path: hostPath.Path, + Type: hostPath.Type, + }) + } + } + return reqs +} + func FetchDeployConfigsByJSON(obj jsonutils.JSONObject) ([]*compute.DeployConfig, error) { deploys := make([]*compute.DeployConfig, 0) if obj.Contains("deploy_configs") { diff --git a/pkg/hostman/hosthandler/handler.go b/pkg/hostman/hosthandler/handler.go index 8ee6a081db..e7c13c1312 100644 --- a/pkg/hostman/hosthandler/handler.go +++ b/pkg/hostman/hosthandler/handler.go @@ -26,9 +26,11 @@ import ( "yunion.io/x/pkg/gotypes" "yunion.io/x/pkg/utils" + hostapi "yunion.io/x/onecloud/pkg/apis/host" "yunion.io/x/onecloud/pkg/appsrv" "yunion.io/x/onecloud/pkg/hostman/hostinfo" "yunion.io/x/onecloud/pkg/hostman/hostinfo/hostconsts" + "yunion.io/x/onecloud/pkg/hostman/hostpath" "yunion.io/x/onecloud/pkg/hostman/hostutils" "yunion.io/x/onecloud/pkg/httperrors" "yunion.io/x/onecloud/pkg/mcclient/auth" @@ -46,6 +48,7 @@ func AddHostHandler(prefix string, app *appsrv.Application) { for action, f := range map[string]actionFunc{ "sync": hostSync, "probe-isolated-devices": hostProbeIsolatedDevices, + "check-host-paths": hostCheckHostPaths, "shutdown-servers-on-host-down": setOnHostDown, "restart-host-agent": hostRestart, } { @@ -106,3 +109,11 @@ func hostRestart(ctx context.Context, hostId string, body jsonutils.JSONObject) func hostProbeIsolatedDevices(ctx context.Context, hostId string, body jsonutils.JSONObject) (interface{}, error) { return hostinfo.Instance().ProbeSyncIsolatedDevices(hostId, body) } + +func hostCheckHostPaths(ctx context.Context, hostId string, body jsonutils.JSONObject) (interface{}, error) { + input := hostapi.HostPathCheckInput{} + if err := body.Unmarshal(&input); err != nil { + return nil, err + } + return hostpath.Check(input) +} diff --git a/pkg/hostman/hostpath/check.go b/pkg/hostman/hostpath/check.go new file mode 100644 index 0000000000..d5aa7ac423 --- /dev/null +++ b/pkg/hostman/hostpath/check.go @@ -0,0 +1,84 @@ +// Copyright 2019 Yunion +// +// 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 hostpath + +import ( + "fmt" + "path/filepath" + "strings" + + "yunion.io/x/onecloud/pkg/apis" + hostapi "yunion.io/x/onecloud/pkg/apis/host" + "yunion.io/x/onecloud/pkg/util/procutils" +) + +var runHostPathCheckCommand = func(name string, args ...string) error { + return procutils.NewRemoteCommandAsFarAsPossible(name, args...).Run() +} + +func Check(input hostapi.HostPathCheckInput) (*hostapi.HostPathCheckOutput, error) { + output := &hostapi.HostPathCheckOutput{ + Results: make([]hostapi.HostPathCheckResult, 0, len(input.Paths)), + } + for _, item := range input.Paths { + if err := validateItem(item); err != nil { + return nil, err + } + output.Results = append(output.Results, checkOne(item)) + } + return output, nil +} + +func validateItem(item hostapi.HostPathCheckItem) error { + if item.Path == "" { + return fmt.Errorf("host path is empty") + } + if strings.ContainsRune(item.Path, 0) { + return fmt.Errorf("host path %q contains NUL byte", item.Path) + } + if !filepath.IsAbs(item.Path) { + return fmt.Errorf("host path %q must be absolute", item.Path) + } + switch item.Type { + case apis.CONTAINER_VOLUME_MOUNT_HOST_PATH_TYPE_DIRECTORY, apis.CONTAINER_VOLUME_MOUNT_HOST_PATH_TYPE_FILE: + return nil + default: + return fmt.Errorf("unsupported host path type %q", item.Type) + } +} + +func checkOne(item hostapi.HostPathCheckItem) hostapi.HostPathCheckResult { + result := hostapi.HostPathCheckResult{ + Path: item.Path, + Type: item.Type, + } + + if err := runHostPathCheckCommand("test", "-e", item.Path); err != nil { + result.Error = fmt.Sprintf("%s does not exist", item.Path) + return result + } + + result.Exists = true + switch item.Type { + case apis.CONTAINER_VOLUME_MOUNT_HOST_PATH_TYPE_DIRECTORY: + result.TypeMatched = runHostPathCheckCommand("test", "-d", item.Path) == nil + case apis.CONTAINER_VOLUME_MOUNT_HOST_PATH_TYPE_FILE: + result.TypeMatched = runHostPathCheckCommand("test", "-f", item.Path) == nil + } + if !result.TypeMatched { + result.Error = fmt.Sprintf("%s is not %s", item.Path, item.Type) + } + return result +} diff --git a/pkg/hostman/hostpath/doc.go b/pkg/hostman/hostpath/doc.go new file mode 100644 index 0000000000..577aad8ea0 --- /dev/null +++ b/pkg/hostman/hostpath/doc.go @@ -0,0 +1 @@ +package hostpath // import "yunion.io/x/onecloud/pkg/hostman/hostpath" diff --git a/pkg/scheduler/algorithm/predicates/host_path_predicate.go b/pkg/scheduler/algorithm/predicates/host_path_predicate.go new file mode 100644 index 0000000000..d574dcb4d6 --- /dev/null +++ b/pkg/scheduler/algorithm/predicates/host_path_predicate.go @@ -0,0 +1,173 @@ +// Copyright 2019 Yunion +// +// 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 predicates + +import ( + "context" + "fmt" + "sync" + + "yunion.io/x/jsonutils" + + "yunion.io/x/onecloud/pkg/apis" + hostapi "yunion.io/x/onecloud/pkg/apis/host" + schedulerapi "yunion.io/x/onecloud/pkg/apis/scheduler" + "yunion.io/x/onecloud/pkg/compute/models" + "yunion.io/x/onecloud/pkg/mcclient" + "yunion.io/x/onecloud/pkg/scheduler/core" +) + +const hostPathPredicateParallelism = 16 + +type hostPathChecker func(context.Context, *models.SHost, []schedulerapi.HostPathRequirement) (*hostapi.HostPathCheckOutput, error) + +type HostPathPredicate struct { + BasePredicate + + checker hostPathChecker + failReasons map[string]string +} + +func NewHostPathPredicate() *HostPathPredicate { + return &HostPathPredicate{} +} + +func (p *HostPathPredicate) Name() string { + return "host_path" +} + +func (p *HostPathPredicate) Clone() core.FitPredicate { + return &HostPathPredicate{ + checker: p.checker, + } +} + +func (p *HostPathPredicate) PreExecute(ctx context.Context, u *core.Unit, cs []core.Candidater) (bool, error) { + reqs := u.SchedData().HostPathRequirements + if len(reqs) == 0 { + return false, nil + } + + checker := p.checker + if checker == nil { + userCred := u.SchedInfo.UserCred + checker = func(ctx context.Context, host *models.SHost, reqs []schedulerapi.HostPathRequirement) (*hostapi.HostPathCheckOutput, error) { + return checkHostPathsOnHost(ctx, userCred, host, reqs) + } + } + + p.failReasons = make(map[string]string) + sem := make(chan struct{}, hostPathPredicateParallelism) + wg := sync.WaitGroup{} + lock := sync.Mutex{} + for _, c := range cs { + candidate := c + wg.Add(1) + go func() { + defer wg.Done() + sem <- struct{}{} + defer func() { + <-sem + }() + + reason := checkHostPathCandidate(ctx, checker, candidate, reqs) + if reason == "" { + return + } + lock.Lock() + p.failReasons[candidate.IndexKey()] = reason + lock.Unlock() + }() + } + wg.Wait() + + return true, nil +} + +func (p *HostPathPredicate) Execute(ctx context.Context, u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) { + h := NewPredicateHelper(p, u, c) + if reason := p.failReasons[c.IndexKey()]; reason != "" { + h.Exclude(reason) + } + return h.GetResult() +} + +func checkHostPathCandidate(ctx context.Context, checker hostPathChecker, candidate core.Candidater, reqs []schedulerapi.HostPathRequirement) string { + host := candidate.Getter().Host() + if host == nil { + return "host path check unavailable: candidate host is nil" + } + output, err := checker(ctx, host, reqs) + return hostPathCheckFailureReason(reqs, output, err) +} + +func checkHostPathsOnHost(ctx context.Context, userCred mcclient.TokenCredential, host *models.SHost, reqs []schedulerapi.HostPathRequirement) (*hostapi.HostPathCheckOutput, error) { + input := hostapi.HostPathCheckInput{ + Paths: make([]hostapi.HostPathCheckItem, 0, len(reqs)), + } + for _, req := range reqs { + input.Paths = append(input.Paths, hostapi.HostPathCheckItem{ + Path: req.Path, + Type: req.Type, + }) + } + + resp, err := host.Request(ctx, userCred, "POST", fmt.Sprintf("/hosts/%s/check-host-paths", host.Id), mcclient.GetTokenHeaders(userCred), jsonutils.Marshal(input)) + if err != nil { + return nil, err + } + output := new(hostapi.HostPathCheckOutput) + if err := resp.Unmarshal(output); err != nil { + return nil, err + } + return output, nil +} + +func hostPathCheckFailureReason(reqs []schedulerapi.HostPathRequirement, output *hostapi.HostPathCheckOutput, err error) string { + if err != nil { + return fmt.Sprintf("host path check unavailable: %v", err) + } + if output == nil { + return "host path check unavailable: empty response" + } + + results := make(map[string]hostapi.HostPathCheckResult, len(output.Results)) + for _, result := range output.Results { + results[hostPathRequirementKey(result.Path, result.Type)] = result + } + for _, req := range reqs { + result, ok := results[hostPathRequirementKey(req.Path, req.Type)] + if !ok { + return fmt.Sprintf("host path %s check result missing", req.Path) + } + if !result.Exists { + if result.Error != "" { + return result.Error + } + return fmt.Sprintf("host path %s does not exist", req.Path) + } + if !result.TypeMatched { + if result.Error != "" { + return result.Error + } + return fmt.Sprintf("host path %s is not %s", req.Path, req.Type) + } + } + return "" +} + +func hostPathRequirementKey(path string, typ apis.ContainerVolumeMountHostPathType) string { + return fmt.Sprintf("%s\x00%s", path, string(typ)) +} diff --git a/pkg/scheduler/algorithmprovider/defaults.go b/pkg/scheduler/algorithmprovider/defaults.go index 2e0867e80f..1a20313f64 100644 --- a/pkg/scheduler/algorithmprovider/defaults.go +++ b/pkg/scheduler/algorithmprovider/defaults.go @@ -31,6 +31,7 @@ func defaultPredicates() sets.String { return sets.NewString( factory.RegisterFitPredicate("a-GuestHostStatusFilter", &predicateguest.StatusPredicate{}), factory.RegisterFitPredicate("b-GuestHypervisorFilter", &predicateguest.HypervisorPredicate{}), + factory.RegisterFitPredicate("c-GuestHostPathFilter", predicates.NewHostPathPredicate()), factory.RegisterFitPredicate("c-GuestHostschedtagFilter", predicates.NewHostSchedtagPredicate()), factory.RegisterFitPredicate("d-GuestMigrateFilter", &predicateguest.MigratePredicate{}), factory.RegisterFitPredicate("e-GuestDomainFilter", &predicates.DomainPredicate{}),