diff --git a/pkg/scheduler/algorithm/predicates/guest/storage_predicate.go b/pkg/scheduler/algorithm/predicates/guest/storage_predicate.go index 35c36a2748..89bdbc733c 100644 --- a/pkg/scheduler/algorithm/predicates/guest/storage_predicate.go +++ b/pkg/scheduler/algorithm/predicates/guest/storage_predicate.go @@ -19,10 +19,14 @@ import ( "fmt" "strings" + "yunion.io/x/pkg/errors" "yunion.io/x/pkg/tristate" + "yunion.io/x/pkg/util/sets" "yunion.io/x/pkg/utils" + "yunion.io/x/onecloud/pkg/apis/compute" "yunion.io/x/onecloud/pkg/scheduler/algorithm/predicates" + "yunion.io/x/onecloud/pkg/scheduler/cache/candidate" "yunion.io/x/onecloud/pkg/scheduler/core" ) @@ -48,6 +52,104 @@ func (p *StoragePredicate) PreExecute(ctx context.Context, u *core.Unit, cs []co return true, nil } +type diskSizeRequest struct { + backend string + max int64 + total int64 + disks []*compute.DiskConfig +} + +func newDiskSizeRequest(backend string) *diskSizeRequest { + return &diskSizeRequest{ + backend: backend, + max: -1, + total: 0, + disks: make([]*compute.DiskConfig, 0), + } +} + +func (req *diskSizeRequest) GetMax() int64 { + return req.max +} + +func (req *diskSizeRequest) GetTotal() int64 { + return req.total +} + +func (req *diskSizeRequest) Add(disk *compute.DiskConfig) *diskSizeRequest { + req.total += int64(disk.SizeMb) + if req.max < int64(disk.SizeMb) { + req.max = int64(disk.SizeMb) + } + + req.disks = append(req.disks, disk) + return req +} + +func (req *diskSizeRequest) NewByMediumType(mt string) *diskSizeRequest { + newReq := newDiskSizeRequest(req.backend) + for _, d := range req.disks { + if d.Medium == mt { + newReq.Add(d) + } + } + return newReq +} + +type diskBackendSizeRequest struct { + reqs map[string]*diskSizeRequest + // beMdm is a map that use backend as key and medium types as value + beMdm map[string]sets.String +} + +func newDiskBackendSizeRequest() *diskBackendSizeRequest { + return &diskBackendSizeRequest{ + reqs: make(map[string]*diskSizeRequest), + beMdm: make(map[string]sets.String), + } +} + +func (ds *diskBackendSizeRequest) get(backend string) (*diskSizeRequest, bool) { + req, ok := ds.reqs[backend] + return req, ok +} + +func (ds *diskBackendSizeRequest) set(backend string, req *diskSizeRequest) *diskBackendSizeRequest { + ds.reqs[backend] = req + return ds +} + +func (ds *diskBackendSizeRequest) Add(disk *compute.DiskConfig) *diskBackendSizeRequest { + backend := disk.Backend + req, ok := ds.get(backend) + if !ok { + req = newDiskSizeRequest(backend) + } + req.Add(disk) + ds.set(backend, req) + + mds, ok := ds.beMdm[backend] + if !ok { + mds = sets.NewString() + } + mds.Insert(disk.Medium) + ds.beMdm[backend] = mds + + return ds +} + +func (ds *diskBackendSizeRequest) GetBackendMediumMap() map[string]sets.String { + return ds.beMdm +} + +func (ds *diskBackendSizeRequest) Get(backend string, mediumType string) (*diskSizeRequest, error) { + req, ok := ds.get(backend) + if !ok { + return nil, errors.Errorf("Not found diskBackendSizeRequest by backend %q", backend) + } + return req.NewByMediumType(mediumType), nil +} + func (p *StoragePredicate) Execute(ctx context.Context, u *core.Unit, c core.Candidater) (bool, []core.PredicateFailureReason, error) { h := predicates.NewPredicateHelper(p, u, c) @@ -87,12 +189,15 @@ func (p *StoragePredicate) Execute(ctx context.Context, u *core.Unit, c core.Can } } - getStorageCapacity := func(backend string, reqMaxSize int64, reqTotalSize int64, useRsvd bool) (*storageCapacity, *storageCapacity) { - totalFree, actualFree := getter.GetFreeStorageSizeOfType(backend, useRsvd) + getStorageCapacity := func(backend string, mediumType string, reqMaxSize int64, reqTotalSize int64, useRsvd bool) (*storageCapacity, *storageCapacity, error) { + totalFree, actualFree, err := getter.GetFreeStorageSizeOfType(backend, mediumType, useRsvd, reqMaxSize) + if err != nil { + return nil, nil, err + } reqTotalSize = utils.Max(reqTotalSize, 1) capacity := totalFree / reqTotalSize actualCapacity := actualFree / reqTotalSize - return newStorageCapacity(capacity, totalFree, false), newStorageCapacity(actualCapacity, actualFree, true) + return newStorageCapacity(capacity, totalFree, false), newStorageCapacity(actualCapacity, actualFree, true), nil } getReqSizeStr := func(backend string) string { @@ -106,27 +211,27 @@ func (p *StoragePredicate) Execute(ctx context.Context, u *core.Unit, c core.Can return strings.Join(ss, "+") } - getStorageFreeStr := func(backend string, useRsvd bool, isActual bool) string { + getStorageFreeStr := func(backend string, mediumType string, useRsvd bool, isActual bool) string { ss := []string{} for _, s := range getter.Storages() { - if s.StorageType == backend { + if candidate.IsStorageBackendMediumMatch(s, backend, mediumType) { if isActual { total := s.Capacity free := total - s.ActualCapacityUsed - ss = append(ss, fmt.Sprintf("actual_total:%d - actual_used:%d = free:%d", total, s.ActualCapacityUsed, free)) + ss = append(ss, fmt.Sprintf("storage %q, actual_total:%d - actual_used:%d = free:%d", s.GetName(), total, s.ActualCapacityUsed, free)) } else { total := int64(float32(s.Capacity) * s.Cmtbound) used := s.GetUsedCapacity(tristate.True) waste := s.GetUsedCapacity(tristate.False) free := total - int64(used) - int64(waste) - ss = append(ss, fmt.Sprintf("total:%d - used:%d - waste:%d = free:%d", total, used, waste, free)) + ss = append(ss, fmt.Sprintf("storage %q, total:%d - used:%d - waste:%d = free:%d", s.GetName(), total, used, waste, free)) } } } return strings.Join(ss, " + ") } - sizeRequest := make(map[string]map[string]int64, 0) + sizeRequest := newDiskBackendSizeRequest() storeRequest := make(map[string]int64, 0) for _, disk := range d.Disks { if isMigrate() && !isLocalhostBackend(disk.Backend) { @@ -136,14 +241,7 @@ func (p *StoragePredicate) Execute(ctx context.Context, u *core.Unit, c core.Can storeRequest[disk.Storage] = 1 } else if !isMigrate() || (isMigrate() && isLocalhostBackend(disk.Backend)) { // if migrate, only local storage need check capacity constraint - if _, ok := sizeRequest[disk.Backend]; !ok { - sizeRequest[disk.Backend] = map[string]int64{"max": -1, "total": 0} - } - max := sizeRequest[disk.Backend]["max"] - if max < int64(disk.SizeMb) { - sizeRequest[disk.Backend]["max"] = int64(disk.SizeMb) - } - sizeRequest[disk.Backend]["total"] += int64(disk.SizeMb) + sizeRequest.Add(disk) } } @@ -157,27 +255,38 @@ func (p *StoragePredicate) Execute(ctx context.Context, u *core.Unit, c core.Can useRsvd := h.UseReserved() minCapacity := int64(0xFFFFFFFF) - appendFailMsg := func(backend string, req map[string]int64, useRsvd bool, capacity *storageCapacity) { - reqStr := fmt.Sprintf("no enough %q storage, req=%v(%v)", backend, req["total"], getReqSizeStr(backend)) + appendFailMsg := func(backend string, mediumType string, req *diskSizeRequest, useRsvd bool, capacity *storageCapacity) { + reqStr := fmt.Sprintf("no enough backend %q, mediumType %q storage, req=%v(%v)", backend, mediumType, req.GetTotal(), getReqSizeStr(backend)) freePrex := "free" isActual := capacity.isActual if isActual { freePrex = "actual_free" } - freeStr := fmt.Sprintf("%s=%v(%v)", freePrex, capacity.free, getStorageFreeStr(backend, useRsvd, isActual)) + freeStr := fmt.Sprintf("%s=%v(%v)", freePrex, capacity.free, getStorageFreeStr(backend, mediumType, useRsvd, isActual)) msg := reqStr + ", " + freeStr h.AppendPredicateFailMsg(msg) } - for be, req := range sizeRequest { - capacity, actualCapacity := getStorageCapacity(be, req["max"], req["total"], useRsvd) - tmpCap := utils.Min(capacity.capacity, actualCapacity.capacity) - if capacity.capacity <= 0 { - appendFailMsg(be, req, useRsvd, capacity) - } else if actualCapacity.capacity <= 0 { - appendFailMsg(be, req, useRsvd, actualCapacity) + for be, mds := range sizeRequest.GetBackendMediumMap() { + for _, medium := range mds.List() { + req, err := sizeRequest.Get(be, medium) + if err != nil { + h.Exclude(fmt.Sprintf("get request size by backend %q, medium %q: %v", be, medium, err)) + break + } + capacity, actualCapacity, err := getStorageCapacity(be, medium, req.GetMax(), req.GetTotal(), useRsvd) + if err != nil { + h.Exclude(err.Error()) + continue + } + tmpCap := utils.Min(capacity.capacity, actualCapacity.capacity) + if capacity.capacity <= 0 { + appendFailMsg(be, medium, req, useRsvd, capacity) + } else if actualCapacity.capacity <= 0 { + appendFailMsg(be, medium, req, useRsvd, actualCapacity) + } + minCapacity = utils.Min(minCapacity, tmpCap) } - minCapacity = utils.Min(minCapacity, tmpCap) } h.SetCapacity(minCapacity) diff --git a/pkg/scheduler/cache/candidate/base.go b/pkg/scheduler/cache/candidate/base.go index 803a1948c4..e9b108c0d9 100644 --- a/pkg/scheduler/cache/candidate/base.go +++ b/pkg/scheduler/cache/candidate/base.go @@ -223,16 +223,48 @@ func (b baseHostGetter) TotalMemorySize(_ bool) int64 { return int64(b.h.MemSize) } -func (b baseHostGetter) GetFreeStorageSizeOfType(storageType string, useRsvd bool) (int64, int64) { +func checkStorageSize(s *api.CandidateStorage, reqMaxSize int64, useRsvd bool) error { + storageSize := s.FreeCapacity + if useRsvd { + storageSize += s.GetReserved() + } + minSize := utils.Min(storageSize, s.ActualFreeCapacity) + if minSize >= reqMaxSize { + return nil + } + return errors.Errorf("storage %q free size %d less than max request %d, use reserverd %v, free_capacity(%d), actual_free_capacity(%d)", s.GetName(), minSize, reqMaxSize, useRsvd, storageSize, s.ActualFreeCapacity) +} + +func IsStorageBackendMediumMatch(s *api.CandidateStorage, backend string, mediumType string) bool { + if s.StorageType != backend { + return false + } + if mediumType == "" { + return true + } + return s.MediumType == mediumType +} + +func (b baseHostGetter) GetFreeStorageSizeOfType(storageType string, mediumType string, useRsvd bool, reqMaxSize int64) (int64, int64, error) { var size int64 var actualSize int64 + foundLEReqStore := false + errs := make([]error, 0) for _, s := range b.Storages() { - if s.StorageType == storageType { - size += int64(float32(s.Capacity) * s.Cmtbound) - actualSize += s.Capacity - s.ActualCapacityUsed + if IsStorageBackendMediumMatch(s, storageType, mediumType) { + size += s.FreeCapacity + actualSize += s.ActualFreeCapacity + if err := checkStorageSize(s, reqMaxSize, false); err != nil { + errs = append(errs, err) + } else { + foundLEReqStore = true + } } } - return size, actualSize + if foundLEReqStore { + return size, actualSize, nil + } + return size, actualSize, errors.NewAggregate(errs) } func (b baseHostGetter) GetFreePort(netId string) int { diff --git a/pkg/scheduler/cache/candidate/hosts.go b/pkg/scheduler/cache/candidate/hosts.go index 2aff4cab78..61cf19184d 100644 --- a/pkg/scheduler/cache/candidate/hosts.go +++ b/pkg/scheduler/cache/candidate/hosts.go @@ -23,7 +23,7 @@ import ( "time" "yunion.io/x/log" - "yunion.io/x/pkg/util/errors" + "yunion.io/x/pkg/errors" "yunion.io/x/pkg/util/sets" "yunion.io/x/pkg/util/workqueue" "yunion.io/x/pkg/utils" @@ -85,8 +85,8 @@ func (h *hostGetter) StorageInfo() []*baremetal.BaremetalStorage { return nil } -func (h *hostGetter) GetFreeStorageSizeOfType(storageType string, useRsvd bool) (int64, int64) { - return h.h.GetFreeStorageSizeOfType(storageType, useRsvd) +func (h *hostGetter) GetFreeStorageSizeOfType(storageType string, mediumType string, useRsvd bool, reqMaxSize int64) (int64, int64, error) { + return h.h.GetFreeStorageSizeOfType(storageType, mediumType, useRsvd, reqMaxSize) } func (h *hostGetter) GetFreePort(netId string) int { @@ -317,18 +317,25 @@ func (h *HostDesc) freeStorageSize(onlyLocal, useRsvd bool) int64 { return total } -func (h *HostDesc) GetFreeStorageSizeOfType(sType string, useRsvd bool) (int64, int64) { - return h.freeStorageSizeOfType(sType, useRsvd) +func (h *HostDesc) GetFreeStorageSizeOfType(sType string, mediumType string, useRsvd bool, reqMaxSize int64) (int64, int64, error) { + return h.freeStorageSizeOfType(sType, mediumType, useRsvd, reqMaxSize) } -func (h *HostDesc) freeStorageSizeOfType(storageType string, useRsvd bool) (int64, int64) { +func (h *HostDesc) freeStorageSizeOfType(storageType string, mediumType string, useRsvd bool, reqMaxSize int64) (int64, int64, error) { var total int64 var actualTotal int64 + foundLEReqStore := false + errs := make([]error, 0) for _, storage := range h.Storages { - if storage.StorageType == storageType { + if IsStorageBackendMediumMatch(storage, storageType, mediumType) { total += int64(storage.FreeCapacity) actualTotal += int64(storage.ActualFreeCapacity) + if err := checkStorageSize(storage, reqMaxSize, useRsvd); err != nil { + errs = append(errs, err) + } else { + foundLEReqStore = true + } } } if utils.IsLocalStorage(storageType) { @@ -338,11 +345,15 @@ func (h *HostDesc) freeStorageSizeOfType(storageType string, useRsvd bool) (int6 total += sizeSub } } - if useRsvd { - return reservedResourceAddCal(total, h.GuestReservedStorageSizeFree(), useRsvd), actualTotal + if !foundLEReqStore { + return 0, 0, errors.NewAggregate(errs) } - return total - int64(h.GetPendingUsage().DiskUsage.Get(storageType)), actualTotal + if useRsvd { + return reservedResourceAddCal(total, h.GuestReservedStorageSizeFree(), useRsvd), actualTotal, nil + } + + return total - int64(h.GetPendingUsage().DiskUsage.Get(storageType)), actualTotal, nil } func (h *HostDesc) GetFreePort(netId string) int { diff --git a/pkg/scheduler/core/types.go b/pkg/scheduler/core/types.go index f292507a33..3642d90007 100644 --- a/pkg/scheduler/core/types.go +++ b/pkg/scheduler/core/types.go @@ -100,7 +100,7 @@ type CandidatePropertyGetter interface { FreeMemorySize(useRsvd bool) int64 StorageInfo() []*baremetal.BaremetalStorage - GetFreeStorageSizeOfType(storageType string, useRsvd bool) (int64, int64) + GetFreeStorageSizeOfType(storageType string, mediumType string, useRsvd bool, reqMaxSize int64) (int64, int64, error) GetFreePort(netId string) int diff --git a/pkg/scheduler/test/mock/core.go b/pkg/scheduler/test/mock/core.go index ec72489caf..5acfd82888 100644 --- a/pkg/scheduler/test/mock/core.go +++ b/pkg/scheduler/test/mock/core.go @@ -180,17 +180,17 @@ func (mr *MockCandidatePropertyGetterMockRecorder) GetFreePort(arg0 interface{}) } // GetFreeStorageSizeOfType mocks base method -func (m *MockCandidatePropertyGetter) GetFreeStorageSizeOfType(arg0 string, arg1 bool) (int64, int64) { +func (m *MockCandidatePropertyGetter) GetFreeStorageSizeOfType(arg0 string, arg1 string, arg2 bool, arg3 int64) (int64, int64, error) { m.ctrl.T.Helper() - ret := m.ctrl.Call(m, "GetFreeStorageSizeOfType", arg0, arg1) + ret := m.ctrl.Call(m, "GetFreeStorageSizeOfType", arg0, arg1, arg2, arg3) ret0, _ := ret[0].(int64) - return ret0, math.MaxInt64 + return ret0, math.MaxInt64, nil } // GetFreeStorageSizeOfType indicates an expected call of GetFreeStorageSizeOfType -func (mr *MockCandidatePropertyGetterMockRecorder) GetFreeStorageSizeOfType(arg0, arg1 interface{}) *gomock.Call { +func (mr *MockCandidatePropertyGetterMockRecorder) GetFreeStorageSizeOfType(arg0, arg1, arg2, arg3 interface{}) *gomock.Call { mr.mock.ctrl.T.Helper() - return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "GetFreeStorageSizeOfType", reflect.TypeOf((*MockCandidatePropertyGetter)(nil).GetFreeStorageSizeOfType), arg0, arg1) + return mr.mock.ctrl.RecordCallWithMethodType(mr.mock, "GetFreeStorageSizeOfType", reflect.TypeOf((*MockCandidatePropertyGetter)(nil).GetFreeStorageSizeOfType), arg0, arg1, arg2, arg3) } // GetIpmiInfo mocks base method diff --git a/pkg/scheduler/test/prepare.go b/pkg/scheduler/test/prepare.go index f2f6d5c9ca..198e98bb31 100644 --- a/pkg/scheduler/test/prepare.go +++ b/pkg/scheduler/test/prepare.go @@ -171,7 +171,7 @@ func buildGetter(ctrl *gomock.Controller, param sGetterParams) *mock.MockCandida cg.EXPECT().FreeCPUCount(gomock.Any()).AnyTimes().Return(param.FreeCPUCount) cg.EXPECT().TotalMemorySize(gomock.Any()).AnyTimes().Return(param.TotalMemorySize) cg.EXPECT().FreeMemorySize(gomock.Any()).AnyTimes().Return(param.FreeMemorySize) - cg.EXPECT().GetFreeStorageSizeOfType(gomock.Any(), gomock.Any()).AnyTimes().Return(param.FreeStorageSizeAnyType, int64(0)) + cg.EXPECT().GetFreeStorageSizeOfType(gomock.Any(), gomock.Any(), gomock.Any(), gomock.Any()).AnyTimes().Return(param.FreeStorageSizeAnyType, int64(0), nil) if param.QuotaKeys != nil { cg.EXPECT().GetQuotaKeys(gomock.Any()).AnyTimes().Return(param.QuotaKeys) }