diff --git a/staging/src/k8s.io/dynamic-resource-allocation/structured/internal/allocatortesting/allocator_testing.go b/staging/src/k8s.io/dynamic-resource-allocation/structured/internal/allocatortesting/allocator_testing.go index 305b1c3ced5..01ddd6b0ba0 100644 --- a/staging/src/k8s.io/dynamic-resource-allocation/structured/internal/allocatortesting/allocator_testing.go +++ b/staging/src/k8s.io/dynamic-resource-allocation/structured/internal/allocatortesting/allocator_testing.go @@ -637,7 +637,11 @@ func TestAllocator(t *testing.T, node *v1.Node expectResults []any - expectError types.GomegaMatcher // can be used to check for no error or match specific error types + expectError types.GomegaMatcher // can be used to check for no error or match specific error + + // Test case setting expectNumAllocateOneInvocations do not run against the "stable" variant of the allocator, + // which doesn't provide the stats and also falls over with excessive runtime for them. + expectNumAllocateOneInvocations int64 }{ "empty": {}, @@ -3487,6 +3491,67 @@ func TestAllocator(t *testing.T, deviceAllocationResult(req1, driverA, pool1, device1, false), )}, }, + "prioritized-list-max-allocation-limit-request": { + features: Features{ + PrioritizedList: true, + }, + claimsToAllocate: objects( + claimWithRequests(claim0, nil, + requestWithPrioritizedList(req0, + subRequest(subReq0, classA, resourceapi.AllocationResultsMaxSize), + subRequest(subReq1, classA, resourceapi.AllocationResultsMaxSize/2), + ), + requestWithPrioritizedList(req1, + subRequest(subReq0, classA, resourceapi.AllocationResultsMaxSize), + subRequest(subReq1, classA, resourceapi.AllocationResultsMaxSize/2), + ), + ), + ), + classes: objects(class(classA, driverA)), + slices: unwrap(sliceWithMultipleDevices(slice1, node1, pool1, driverA, resourceapi.AllocationResultsMaxSize*2)), + node: node(node1, region1), + + expectResults: []any{allocationResult( + localNodeSelector(node1), + slices.Concat( + multipleDeviceAllocationResults(req0SubReq1, driverA, pool1, resourceapi.AllocationResultsMaxSize/2, 0), + multipleDeviceAllocationResults(req1SubReq1, driverA, pool1, resourceapi.AllocationResultsMaxSize/2, resourceapi.AllocationResultsMaxSize/2), + )..., + )}, + expectNumAllocateOneInvocations: 75, + }, + "prioritized-list-max-allocation-allocation-mode-all": { + features: Features{ + PrioritizedList: true, + }, + claimsToAllocate: objects( + claimWithRequests(claim0, nil, + requestWithPrioritizedList(req0, + resourceapi.DeviceSubRequest{ + Name: subReq0, + AllocationMode: resourceapi.DeviceAllocationModeAll, + DeviceClassName: classA, + }, + subRequest(subReq1, classA, 2), + ), + request(req1, classB, 2), + ), + ), + classes: objects(class(classA, driverA), class(classB, driverB)), + slices: unwrap( + sliceWithMultipleDevices(slice1, node1, pool1, driverA, resourceapi.AllocationResultsMaxSize-1), + sliceWithMultipleDevices(slice2, node1, pool2, driverB, 2), + ), + node: node(node1, region1), + expectResults: []any{allocationResult( + localNodeSelector(node1), + deviceAllocationResult(req0SubReq1, driverA, pool1, "device-0", false), + deviceAllocationResult(req0SubReq1, driverA, pool1, "device-1", false), + deviceAllocationResult(req1, driverB, pool2, "device-0", false), + deviceAllocationResult(req1, driverB, pool2, "device-1", false), + )}, + expectNumAllocateOneInvocations: 42, + }, } for name, tc := range testcases { @@ -3517,6 +3582,10 @@ func TestAllocator(t *testing.T, allocator, err := newAllocator(ctx, tc.features, unwrap(claimsToAllocate...), sets.New(allocatedDevices...), classLister, slices, cel.NewCache(1)) g.Expect(err).ToNot(gomega.HaveOccurred()) + if _, ok := allocator.(internal.AllocatorExtended); tc.expectNumAllocateOneInvocations > 0 && !ok { + t.Skipf("%T does not support the AllocatorStats interface", allocator) + } + results, err := allocator.Allocate(ctx, tc.node) matchError := tc.expectError if matchError == nil { @@ -3530,6 +3599,11 @@ func TestAllocator(t *testing.T, g.Expect(allocatedDevices).To(gomega.HaveExactElements(tc.allocatedDevices)) g.Expect(slices).To(gomega.ConsistOf(tc.slices)) g.Expect(classLister.objs).To(gomega.ConsistOf(tc.classes)) + + if tc.expectNumAllocateOneInvocations > 0 { + stats := allocator.(internal.AllocatorExtended).GetStats() + g.Expect(stats.NumAllocateOneInvocations).To(gomega.Equal(tc.expectNumAllocateOneInvocations)) + } }) } } diff --git a/staging/src/k8s.io/dynamic-resource-allocation/structured/internal/incubating/allocator_incubating.go b/staging/src/k8s.io/dynamic-resource-allocation/structured/internal/incubating/allocator_incubating.go index 11060a15b89..28c8d93c777 100644 --- a/staging/src/k8s.io/dynamic-resource-allocation/structured/internal/incubating/allocator_incubating.go +++ b/staging/src/k8s.io/dynamic-resource-allocation/structured/internal/incubating/allocator_incubating.go @@ -24,6 +24,7 @@ import ( "slices" "strings" "sync" + "sync/atomic" v1 "k8s.io/api/core/v1" resourceapi "k8s.io/api/resource/v1beta1" @@ -39,6 +40,7 @@ import ( type DeviceClassLister = internal.DeviceClassLister type Features = internal.Features type DeviceID = internal.DeviceID +type Stats = internal.Stats func MakeDeviceID(driver, pool, device string) DeviceID { return internal.MakeDeviceID(driver, pool, device) @@ -72,8 +74,15 @@ type Allocator struct { // access to this map must be synchronized. availableCounters map[string]counterSets mutex sync.RWMutex + // numAllocateOneInvocations counts the number of times the allocateOne + // function is called for the allocator. This is a measurement of the + // amount of work the allocator had to do to allocate devices + // for the claims. + numAllocateOneInvocations atomic.Int64 } +var _ internal.AllocatorExtended = &Allocator{} + // NewAllocator returns an allocator for a certain set of claims or an error if // some problem was detected which makes it impossible to allocate claims. // @@ -362,6 +371,13 @@ func (a *Allocator) Allocate(ctx context.Context, node *v1.Node) (finalResult [] return result, nil } +func (a *Allocator) GetStats() Stats { + s := Stats{ + NumAllocateOneInvocations: a.numAllocateOneInvocations.Load(), + } + return s +} + func (alloc *allocator) validateDeviceRequest(request requestAccessor, parentRequest requestAccessor, requestKey requestIndices, pools []*Pool) (requestData, error) { claim := alloc.claimsToAllocate[requestKey.claimIndex] requestData := requestData{ @@ -446,6 +462,13 @@ func (alloc *allocator) validateDeviceRequest(request requestAccessor, parentReq // that allocation cannot succeed. var errStop = errors.New("stop allocation") +// errAllocationResultMaxSizeExceeded is a special error that gets return by +// allocatedOne when the number of allocated devices exceeds the max number +// allowed. This is checked by earlier invocations in the recursion and used +// to do more aggressive backtracking and avoid attempting allocations that +// we know can not succeed. +var errAllocationResultMaxSizeExceeded = errors.New("allocation max size exceeded") + // allocator is used while an [Allocator.Allocate] is running. Only a single // goroutine works with it, so there is no need for locking. type allocator struct { @@ -679,6 +702,7 @@ func lookupAttribute(device *draapi.BasicDevice, deviceID DeviceID, attributeNam // This allows the logic for subrequests to call allocateOne with the same // device index without causing infinite recursion. func (alloc *allocator) allocateOne(r deviceIndices, allocateSubRequest bool) (bool, error) { + alloc.numAllocateOneInvocations.Add(1) if r.claimIndex >= len(alloc.claimsToAllocate) { // Done! If we were doing scoring, we would compare the current allocation result // against the previous one, keep the best, and continue. Without scoring, we stop @@ -690,7 +714,17 @@ func (alloc *allocator) allocateOne(r deviceIndices, allocateSubRequest bool) (b claim := alloc.claimsToAllocate[r.claimIndex] if r.requestIndex >= len(claim.Spec.Devices.Requests) { // Done with the claim, continue with the next one. - return alloc.allocateOne(deviceIndices{claimIndex: r.claimIndex + 1}, false) + success, err := alloc.allocateOne(deviceIndices{claimIndex: r.claimIndex + 1}, false) + if errors.Is(err, errAllocationResultMaxSizeExceeded) { + // We don't need to propagate this further because + // this is not a fatal error. Retrying the claim under + // different circumstances may succeed if it uses + // subrequests and changing the allocation of some + // prior claim enables allocating a subrequest here + // which needs fewer devices. + return false, nil + } + return success, err } // r.subRequestIndex is zero unless the for loop below is in the @@ -703,16 +737,40 @@ func (alloc *allocator) allocateOne(r deviceIndices, allocateSubRequest bool) (b // hitting the first subrequest, but not if we are already working on a // specific subrequest. if !allocateSubRequest && requestData.parentRequest != nil { + // Keep track of whether all attempts to do allocation with the + // subrequests results in the allocation result limit exceeded. + // If so, there is no need to make attempts with other devices + // in the previous request (if any), except when + // it is a firstAvailable request where some sub-requests + // need less devices than others. + allAllocationExceeded := true for subRequestIndex := 0; ; subRequestIndex++ { nextSubRequestKey := requestKey nextSubRequestKey.subRequestIndex = subRequestIndex if _, ok := alloc.requestData[nextSubRequestKey]; !ok { // Past the end of the subrequests without finding a solution -> give up. + // + // Return errAllocationResultMaxSizeExceeded if all + // attempts for the subrequests failed to due to reaching + // the max size limit. This would mean that there are no + // solution that involves the previous request (if any). + if allAllocationExceeded { + return false, errAllocationResultMaxSizeExceeded + } return false, nil } r.subRequestIndex = subRequestIndex success, err := alloc.allocateOne(r, true /* prevent infinite recusion */) + // If we reached the allocation result limit, we can try + // with the next subrequest if there is one. It might request + // fewer devices, so it might succeed. + if errors.Is(err, errAllocationResultMaxSizeExceeded) { + continue + } + // If we get here, at least one of the subrequests failed for a + // different reason than errAllocationResultMaxSizeExceeded. + allAllocationExceeded = false if err != nil { return false, err } @@ -744,17 +802,24 @@ func (alloc *allocator) allocateOne(r deviceIndices, allocateSubRequest bool) (b // Done with request, continue with next one. We have completed the work for // the request or subrequest, so we can no longer be allocating devices for // a subrequest. - return alloc.allocateOne(deviceIndices{claimIndex: r.claimIndex, requestIndex: r.requestIndex + 1}, false) + success, err := alloc.allocateOne(deviceIndices{claimIndex: r.claimIndex, requestIndex: r.requestIndex + 1}, false) + // We want to propagate any errAllocationResultMaxSizeExceeded to the caller. If + // that error is returned here, it means none of the requests/subrequests after this one + // could be allocated while staying within the limit on the number of devices, so there + // are no solution in the current request/subrequest that would work. + return success, err } + // Before trying to allocate devices, check if allocating the devices + // in the current request will put us over the threshold. // We can calculate this by adding the number of already allocated devices with the number // of devices in the current request, and then finally subtract the deviceIndex since we // don't want to double count any devices already allocated for the current request. numDevicesAfterAlloc := len(alloc.result[r.claimIndex].devices) + requestData.numDevices - r.deviceIndex if numDevicesAfterAlloc > resourceapi.AllocationResultsMaxSize { - // Don't return an error here since we want to keep searching for - // a solution that works. - return false, nil + // Return a special error so we can identify this situation in the + // callers and do more aggressive backtracking. + return false, errAllocationResultMaxSizeExceeded } alloc.logger.V(6).Info("Allocating one device", "currentClaim", r.claimIndex, "totalClaims", len(alloc.claimsToAllocate), "currentRequest", r.requestIndex, "currentSubRequest", r.subRequestIndex, "totalRequestsPerClaim", len(claim.Spec.Devices.Requests), "currentDevice", r.deviceIndex, "devicesPerRequest", requestData.numDevices, "allDevices", doAllDevices, "adminAccess", request.adminAccess()) @@ -773,13 +838,12 @@ func (alloc *allocator) allocateOne(r deviceIndices, allocateSubRequest bool) (b return false, nil } done, err := alloc.allocateOne(deviceIndices{claimIndex: r.claimIndex, requestIndex: r.requestIndex, deviceIndex: r.deviceIndex + 1}, allocateSubRequest) - if err != nil { - return false, err - } - if !done { - // Backtrack. + if err != nil || !done { + // If we get an error or didn't complete, we need to backtrack. Depending + // on the situation we might be able to retry, so we make sure we + // deallocate. deallocate() - return false, nil + return false, err } return done, nil } @@ -839,17 +903,22 @@ func (alloc *allocator) allocateOne(r deviceIndices, allocateSubRequest bool) (b deviceIndex: r.deviceIndex + 1, } done, err := alloc.allocateOne(deviceKey, allocateSubRequest) - if err != nil { - return false, err - } - - // If we found a solution, then we can stop. - if done { + // If we found a solution, we can stop. + if err == nil && done { return done, nil } - // Otherwise try some other device after rolling back. + // Otherwise we didn't find a solution, and we need to deallocate + // so the temporary allocation is correct for trying other devices. deallocate() + + if err != nil { + // If we hit an error, we return. This might be that we reached + // the allocation size limit, and if so, it will be caught further + // up the stack and other subrequests will be attempted if there + // are any. + return false, err + } } } } @@ -970,6 +1039,9 @@ func (alloc *allocator) selectorsMatch(r requestIndices, device *draapi.BasicDev // as if that candidate had been allocated. If allocation cannot continue later // and must try something else, then the rollback function can be invoked to // restore the previous state. +// +// The rollback function is only provided in case of a successful allocation +// (true and no error). func (alloc *allocator) allocateDevice(r deviceIndices, device deviceWithID, must bool) (bool, func(), error) { claim := alloc.claimsToAllocate[r.claimIndex] requestKey := requestIndices{claimIndex: r.claimIndex, requestIndex: r.requestIndex, subRequestIndex: r.subRequestIndex} diff --git a/staging/src/k8s.io/dynamic-resource-allocation/structured/internal/stable/allocator_stable.go b/staging/src/k8s.io/dynamic-resource-allocation/structured/internal/stable/allocator_stable.go index 59974caa846..a445c8b6aed 100644 --- a/staging/src/k8s.io/dynamic-resource-allocation/structured/internal/stable/allocator_stable.go +++ b/staging/src/k8s.io/dynamic-resource-allocation/structured/internal/stable/allocator_stable.go @@ -39,6 +39,7 @@ import ( type DeviceClassLister = internal.DeviceClassLister type Features = internal.Features type DeviceID = internal.DeviceID +type Stats = internal.Stats func MakeDeviceID(driver, pool, device string) DeviceID { return internal.MakeDeviceID(driver, pool, device) @@ -361,6 +362,11 @@ func (a *Allocator) Allocate(ctx context.Context, node *v1.Node) (finalResult [] return result, nil } +func (a *Allocator) GetStats() *Stats { + // Not implemented. + return nil +} + func (alloc *allocator) validateDeviceRequest(request requestAccessor, parentRequest requestAccessor, requestKey requestIndices, pools []*Pool) (requestData, error) { claim := alloc.claimsToAllocate[requestKey.claimIndex] requestData := requestData{ diff --git a/staging/src/k8s.io/dynamic-resource-allocation/structured/internal/types.go b/staging/src/k8s.io/dynamic-resource-allocation/structured/internal/types.go index 9f85f1ad853..d10ffc58811 100644 --- a/staging/src/k8s.io/dynamic-resource-allocation/structured/internal/types.go +++ b/staging/src/k8s.io/dynamic-resource-allocation/structured/internal/types.go @@ -33,11 +33,27 @@ type DeviceClassLister interface { } // Allocator is intentionally not documented here. See the main package for docs. +// +// This interface is also broader than the public one. type Allocator interface { ClaimsToAllocate() []*resourceapi.ResourceClaim Allocate(ctx context.Context, node *v1.Node) (finalResult []resourceapi.AllocationResult, finalErr error) } +// AllocatorExtended is an optional interface. Not all variants implement it. +type AllocatorExtended interface { + // Stats shows statistics from the allocation process. + // May return nil if not implemented. + GetStats() Stats +} + +// Stats shows statistics from the allocation process. +type Stats struct { + // NumAllocateOneInvocations counts the number of times the allocateOne function + // got called. + NumAllocateOneInvocations int64 +} + // Features contains all feature gates that may influence the behavior of ResourceClaim allocation. type Features struct { // Sorted alphabetically. When adding a new entry, also extend Set and FeaturesAll.