Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
16 changes: 8 additions & 8 deletions go.mod
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
module sigs.k8s.io/karpenter

go 1.25.7
go 1.25.11

require (
github.com/Pallinder/go-randomdata v1.2.0
Expand All @@ -22,7 +22,7 @@ require (
github.com/samber/lo v1.52.0
go.uber.org/multierr v1.11.0
go.uber.org/zap v1.27.1
golang.org/x/text v0.33.0
golang.org/x/text v0.37.0
golang.org/x/time v0.14.0
k8s.io/api v0.35.0
k8s.io/apiextensions-apiserver v0.35.0
Expand Down Expand Up @@ -81,13 +81,13 @@ require (
github.com/x448/float16 v0.8.4 // indirect
go.yaml.in/yaml/v2 v2.4.3 // indirect
go.yaml.in/yaml/v3 v3.0.4 // indirect
golang.org/x/mod v0.32.0 // indirect
golang.org/x/net v0.49.0 // indirect
golang.org/x/mod v0.35.0 // indirect
golang.org/x/net v0.55.0 // indirect
golang.org/x/oauth2 v0.34.0 // indirect
golang.org/x/sync v0.19.0 // indirect
golang.org/x/sys v0.40.0 // indirect
golang.org/x/term v0.39.0 // indirect
golang.org/x/tools v0.41.0 // indirect
golang.org/x/sync v0.20.0 // indirect
golang.org/x/sys v0.45.0 // indirect
golang.org/x/term v0.43.0 // indirect
golang.org/x/tools v0.44.0 // indirect
gomodules.xyz/jsonpatch/v2 v2.5.0 // indirect
google.golang.org/protobuf v1.36.11 // indirect
gopkg.in/evanphx/json-patch.v4 v4.13.0 // indirect
Expand Down
28 changes: 14 additions & 14 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -182,24 +182,24 @@ go.yaml.in/yaml/v2 v2.4.3 h1:6gvOSjQoTB3vt1l+CU+tSyi/HOjfOjRLJ4YwYZGwRO0=
go.yaml.in/yaml/v2 v2.4.3/go.mod h1:zSxWcmIDjOzPXpjlTTbAsKokqkDNAVtZO0WOMiT90s8=
go.yaml.in/yaml/v3 v3.0.4 h1:tfq32ie2Jv2UxXFdLJdh3jXuOzWiL1fo0bu/FbuKpbc=
go.yaml.in/yaml/v3 v3.0.4/go.mod h1:DhzuOOF2ATzADvBadXxruRBLzYTpT36CKvDb3+aBEFg=
golang.org/x/mod v0.32.0 h1:9F4d3PHLljb6x//jOyokMv3eX+YDeepZSEo3mFJy93c=
golang.org/x/mod v0.32.0/go.mod h1:SgipZ/3h2Ci89DlEtEXWUk/HteuRin+HHhN+WbNhguU=
golang.org/x/net v0.49.0 h1:eeHFmOGUTtaaPSGNmjBKpbng9MulQsJURQUAfUwY++o=
golang.org/x/net v0.49.0/go.mod h1:/ysNB2EvaqvesRkuLAyjI1ycPZlQHM3q01F02UY/MV8=
golang.org/x/mod v0.35.0 h1:Ww1D637e6Pg+Zb2KrWfHQUnH2dQRLBQyAtpr/haaJeM=
golang.org/x/mod v0.35.0/go.mod h1:+GwiRhIInF8wPm+4AoT6L0FA1QWAad3OMdTRx4tFYlU=
golang.org/x/net v0.55.0 h1:bcvxaJn3e1U6InsFWt1JUq1aSjnRxLzT2rtD2KfkDF8=
golang.org/x/net v0.55.0/go.mod h1:L5U2KuzuOe1lY7Z+aWVIKK6qEeJXnXV9yzGA+WCHJww=
golang.org/x/oauth2 v0.34.0 h1:hqK/t4AKgbqWkdkcAeI8XLmbK+4m4G5YeQRrmiotGlw=
golang.org/x/oauth2 v0.34.0/go.mod h1:lzm5WQJQwKZ3nwavOZ3IS5Aulzxi68dUSgRHujetwEA=
golang.org/x/sync v0.19.0 h1:vV+1eWNmZ5geRlYjzm2adRgW2/mcpevXNg50YZtPCE4=
golang.org/x/sync v0.19.0/go.mod h1:9KTHXmSnoGruLpwFjVSX0lNNA75CykiMECbovNTZqGI=
golang.org/x/sys v0.40.0 h1:DBZZqJ2Rkml6QMQsZywtnjnnGvHza6BTfYFWY9kjEWQ=
golang.org/x/sys v0.40.0/go.mod h1:OgkHotnGiDImocRcuBABYBEXf8A9a87e/uXjp9XT3ks=
golang.org/x/term v0.39.0 h1:RclSuaJf32jOqZz74CkPA9qFuVTX7vhLlpfj/IGWlqY=
golang.org/x/term v0.39.0/go.mod h1:yxzUCTP/U+FzoxfdKmLaA0RV1WgE0VY7hXBwKtY/4ww=
golang.org/x/text v0.33.0 h1:B3njUFyqtHDUI5jMn1YIr5B0IE2U0qck04r6d4KPAxE=
golang.org/x/text v0.33.0/go.mod h1:LuMebE6+rBincTi9+xWTY8TztLzKHc/9C1uBCG27+q8=
golang.org/x/sync v0.20.0 h1:e0PTpb7pjO8GAtTs2dQ6jYa5BWYlMuX047Dco/pItO4=
golang.org/x/sync v0.20.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0=
golang.org/x/sys v0.45.0 h1:dO4czNzziLiiXplLQgBCEpCvXQ3dnkn0SdaZSYdQ+FY=
golang.org/x/sys v0.45.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw=
golang.org/x/term v0.43.0 h1:S4RLU2sB31O/NCl+zFN9Aru9A/Cq2aqKpTZJ6B+DwT4=
golang.org/x/term v0.43.0/go.mod h1:lrhlHNdQJHO+1qVYiHfFKVuVioJIheAc3fBSMFYEIsk=
golang.org/x/text v0.37.0 h1:Cqjiwd9eSg8e0QAkyCaQTNHFIIzWtidPahFWR83rTrc=
golang.org/x/text v0.37.0/go.mod h1:a5sjxXGs9hsn/AJVwuElvCAo9v8QYLzvavO5z2PiM38=
golang.org/x/time v0.14.0 h1:MRx4UaLrDotUKUdCIqzPC48t1Y9hANFKIRpNx+Te8PI=
golang.org/x/time v0.14.0/go.mod h1:eL/Oa2bBBK0TkX57Fyni+NgnyQQN4LitPmob2Hjnqw4=
golang.org/x/tools v0.41.0 h1:a9b8iMweWG+S0OBnlU36rzLp20z1Rp10w+IY2czHTQc=
golang.org/x/tools v0.41.0/go.mod h1:XSY6eDqxVNiYgezAVqqCeihT4j1U2CCsqvH3WhQpnlg=
golang.org/x/tools v0.44.0 h1:UP4ajHPIcuMjT1GqzDWRlalUEoY+uzoZKnhOjbIPD2c=
golang.org/x/tools v0.44.0/go.mod h1:KA0AfVErSdxRZIsOVipbv3rQhVXTnlU6UhKxHd1seDI=
gomodules.xyz/jsonpatch/v2 v2.5.0 h1:JELs8RLM12qJGXU4u/TO3V25KW8GreMKl9pdkk14RM0=
gomodules.xyz/jsonpatch/v2 v2.5.0/go.mod h1:AH3dM2RI6uoBZxn3LVrfvJ3E0/9dG4cSrbuBJT4moAY=
google.golang.org/protobuf v1.36.11 h1:fV6ZwhNocDyBLK0dj+fg8ektcVegBBuEolpbTQyBNVE=
Expand Down
102 changes: 96 additions & 6 deletions pkg/controllers/disruption/consolidation.go
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,7 @@ import (
"context"
"errors"
"fmt"
"math"
"sort"
"time"

Expand Down Expand Up @@ -136,7 +137,7 @@ func (c *consolidation) sortCandidates(candidates []*Candidate) []*Candidate {
func (c *consolidation) computeConsolidation(ctx context.Context, candidates ...*Candidate) (Command, error) {
var err error
// Run scheduling simulation to compute consolidation option
results, err := SimulateScheduling(ctx, c.kubeClient, c.cluster, c.provisioner, candidates...)
results, err := simulateScheduling(ctx, c.kubeClient, c.cluster, c.provisioner, []pscheduling.Options{pscheduling.DisableReservedCapacityFallback}, candidates...)
if err != nil {
// if a candidate node is now deleting, just retry
if errors.Is(err, errCandidateDeleting) {
Expand All @@ -154,26 +155,37 @@ func (c *consolidation) computeConsolidation(ctx context.Context, candidates ...
return Command{}, nil
}

// get the current node price based on the offering
// fallback if we can't find the specific zonal pricing data
candidatePrice := getCandidatePrices(candidates)

// were we able to schedule all the pods on the inflight candidates?
if len(results.NewNodeClaims) == 0 {
if hasNonEmptyReservedCandidatesMovingToNonReservedCapacity(candidates, results) {
if len(candidates) == 1 {
c.recorder.Publish(disruptionevents.Unconsolidatable(candidates[0].Node, candidates[0].NodeClaim, "Can't consolidate non-empty reserved node onto non-reserved capacity")...)
}
return Command{}, nil
}
return Command{
Candidates: candidates,
Results: results,
}, nil
}

// we're not going to turn a single node into multiple candidates
if len(results.NewNodeClaims) != 1 {
if len(candidates) == 1 {
cmd := c.computeReservedMultiNodeReplacement(candidates, results, candidatePrice)
if cmd.Decision() != NoOpDecision {
return cmd, nil
}
}
if len(candidates) == 1 {
c.recorder.Publish(disruptionevents.Unconsolidatable(candidates[0].Node, candidates[0].NodeClaim, fmt.Sprintf("Can't remove without creating %d candidates", len(results.NewNodeClaims)))...)
}
return Command{}, nil
}

// get the current node price based on the offering
// fallback if we can't find the specific zonal pricing data
candidatePrice := getCandidatePrices(candidates)

allExistingAreSpot := true
for _, cn := range candidates {
if cn.capacityType != v1.CapacityTypeSpot {
Expand Down Expand Up @@ -228,6 +240,84 @@ func (c *consolidation) computeConsolidation(ctx context.Context, candidates ...
return cmd, nil
}

func (c *consolidation) computeReservedMultiNodeReplacement(candidates []*Candidate, results pscheduling.Results, candidatePrice float64) Command {
if !allNodeClaimsLaunchReservedOnly(results.NewNodeClaims) {
return Command{}
}
var err error
for i := range results.NewNodeClaims {
results.NewNodeClaims[i].InstanceTypeOptions = results.NewNodeClaims[i].InstanceTypeOptions.OrderByPrice(results.NewNodeClaims[i].Requirements)
results.NewNodeClaims[i], err = results.NewNodeClaims[i].RemoveInstanceTypeOptionsByPriceAndMinValues(results.NewNodeClaims[i].Requirements, candidatePrice)
if err != nil {
c.recorder.Publish(disruptionevents.Unconsolidatable(candidates[0].Node, candidates[0].NodeClaim, fmt.Sprintf("Filtering by price: %v", err))...)
return Command{}
}
if len(results.NewNodeClaims[i].InstanceTypeOptions) == 0 {
c.recorder.Publish(disruptionevents.Unconsolidatable(candidates[0].Node, candidates[0].NodeClaim, "Can't replace with a cheaper node")...)
return Command{}
}
}
if price, ok := replacementNodeClaimsPrice(results.NewNodeClaims); !ok || price >= candidatePrice {
c.recorder.Publish(disruptionevents.Unconsolidatable(candidates[0].Node, candidates[0].NodeClaim, "Can't replace with cheaper reserved nodes")...)
return Command{}
}
cmd := Command{
Candidates: candidates,
Replacements: replacementsFromNodeClaims(results.NewNodeClaims...),
Results: results,
}
cmd.EmitCandidateEvents(c.recorder)
return cmd
}

func hasNonEmptyReservedCandidatesMovingToNonReservedCapacity(candidates []*Candidate, results pscheduling.Results) bool {
for _, candidate := range candidates {
if candidate.capacityType != v1.CapacityTypeReserved || len(candidate.reschedulablePods) == 0 {
continue
}
for _, p := range candidate.reschedulablePods {
existingNode, ok := lo.Find(results.ExistingNodes, func(n *pscheduling.ExistingNode) bool {
return lo.ContainsBy(n.Pods, func(scheduledPod *corev1.Pod) bool {
return podsEqual(p, scheduledPod)
})
})
if !ok || existingNode.Labels()[v1.CapacityTypeLabelKey] != v1.CapacityTypeReserved {
return true
}
}
}
return false
}

func allNodeClaimsLaunchReservedOnly(nodeClaims []*pscheduling.NodeClaim) bool {
return len(nodeClaims) > 0 && lo.EveryBy(nodeClaims, func(nc *pscheduling.NodeClaim) bool {
capacityTypeRequirement := nc.Requirements.Get(v1.CapacityTypeLabelKey)
return capacityTypeRequirement.Len() == 1 && capacityTypeRequirement.Has(v1.CapacityTypeReserved)
})
}

func replacementNodeClaimsPrice(nodeClaims []*pscheduling.NodeClaim) (float64, bool) {
var price float64
for _, nodeClaim := range nodeClaims {
if len(nodeClaim.InstanceTypeOptions) == 0 {
return 0.0, false
}
launchPrice := nodeClaim.InstanceTypeOptions[0].Offerings.Available().WorstLaunchPrice(nodeClaim.Requirements)
if launchPrice == math.MaxFloat64 {
return 0.0, false
}
price += launchPrice
}
return price, true
}

func podsEqual(a, b *corev1.Pod) bool {
if a.UID != "" && b.UID != "" {
return a.UID == b.UID
}
return client.ObjectKeyFromObject(a) == client.ObjectKeyFromObject(b)
}

// Compute command to execute spot-to-spot consolidation if:
// 1. The SpotToSpotConsolidation feature flag is set to true.
// 2. For single-node consolidation:
Expand Down
155 changes: 155 additions & 0 deletions pkg/controllers/disruption/consolidation_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -4616,6 +4616,161 @@ var _ = Describe("Consolidation", func() {
Entry("from on-demand", v1.CapacityTypeOnDemand),
Entry("from spot", v1.CapacityTypeSpot),
)
It("does not delete a non-empty reserved node onto existing non-reserved capacity", func() {
rs := test.ReplicaSet()
ExpectApplied(ctx, env.Client, rs)
Expect(env.Client.Get(ctx, client.ObjectKeyFromObject(rs), rs)).To(Succeed())
pod := test.Pod(test.PodOptions{
ObjectMeta: metav1.ObjectMeta{Labels: labels,
OwnerReferences: []metav1.OwnerReference{
{
APIVersion: "apps/v1",
Kind: "ReplicaSet",
Name: rs.Name,
UID: rs.UID,
Controller: lo.ToPtr(true),
BlockOwnerDeletion: lo.ToPtr(true),
},
}},
})
ExpectApplied(ctx, env.Client, rs, pod, reservedNode, reservedNodeClaim, node, nodeClaim, nodePool)
ExpectManualBinding(ctx, env.Client, pod, reservedNode)
ExpectMakeNodesAndNodeClaimsInitializedAndStateUpdated(ctx, env.Client, nodeStateController, nodeClaimStateController, []*corev1.Node{reservedNode, node}, []*v1.NodeClaim{reservedNodeClaim, nodeClaim})

c := disruption.MakeConsolidation(fakeClock, cluster, env.Client, prov, cloudProvider, recorder, queue)
singleNodeConsolidation := disruption.NewSingleNodeConsolidation(c, disruption.WithValidator(NopValidator{}))
budgets, err := disruption.BuildDisruptionBudgetMapping(ctx, cluster, fakeClock, env.Client, cloudProvider, recorder, singleNodeConsolidation.Reason())
Expect(err).To(Succeed())
candidates, err := disruption.GetCandidates(ctx, cluster, env.Client, recorder, fakeClock, cloudProvider, singleNodeConsolidation.ShouldDisrupt, singleNodeConsolidation.Class(), queue)
Expect(err).To(Succeed())
reservedCandidate, ok := lo.Find(candidates, func(candidate *disruption.Candidate) bool {
return candidate.NodeClaim.Name == reservedNodeClaim.Name
})
Expect(ok).To(BeTrue())

cmds, err := singleNodeConsolidation.ComputeCommands(ctx, budgets, reservedCandidate)
Expect(err).To(Succeed())
Expect(cmds).To(BeEmpty())
})
It("can replace one on-demand node with multiple reserved nodeclaims", func() {
reservationID := "r-small-reserved"
largeOnDemand := fake.NewInstanceType(fake.InstanceTypeOptions{
Name: "large-on-demand",
Resources: corev1.ResourceList{
corev1.ResourceCPU: resource.MustParse("8"),
corev1.ResourceMemory: resource.MustParse("8Gi"),
corev1.ResourcePods: resource.MustParse("10"),
},
Offerings: cloudprovider.Offerings{{
Available: true,
Price: 100.0,
Requirements: scheduling.NewLabelRequirements(map[string]string{
v1.CapacityTypeLabelKey: v1.CapacityTypeOnDemand,
corev1.LabelTopologyZone: "test-zone-1",
}),
}},
})
smallReserved := fake.NewInstanceType(fake.InstanceTypeOptions{
Name: "small-reserved",
Resources: corev1.ResourceList{
corev1.ResourceCPU: resource.MustParse("4"),
corev1.ResourceMemory: resource.MustParse("4Gi"),
corev1.ResourcePods: resource.MustParse("10"),
},
Offerings: cloudprovider.Offerings{{
Available: true,
Price: 0.00001,
ReservationCapacity: 2,
Requirements: scheduling.NewLabelRequirements(map[string]string{
v1.CapacityTypeLabelKey: v1.CapacityTypeReserved,
corev1.LabelTopologyZone: "test-zone-1",
v1alpha1.LabelReservationID: reservationID,
}),
}},
})
cloudProvider.InstanceTypes = []*cloudprovider.InstanceType{largeOnDemand, smallReserved}

nodePool.Spec.Template.Spec.Requirements = []v1.NodeSelectorRequirementWithMinValues{
{
Key: v1.CapacityTypeLabelKey,
Operator: corev1.NodeSelectorOpIn,
Values: []string{v1.CapacityTypeOnDemand, v1.CapacityTypeReserved},
},
}
onDemandNodeClaim, onDemandNode := test.NodeClaimAndNode(v1.NodeClaim{
ObjectMeta: metav1.ObjectMeta{
Labels: map[string]string{
v1.NodePoolLabelKey: nodePool.Name,
corev1.LabelInstanceTypeStable: largeOnDemand.Name,
v1.CapacityTypeLabelKey: v1.CapacityTypeOnDemand,
corev1.LabelTopologyZone: "test-zone-1",
},
},
Status: v1.NodeClaimStatus{
Allocatable: map[corev1.ResourceName]resource.Quantity{
corev1.ResourceCPU: resource.MustParse("8"),
corev1.ResourceMemory: resource.MustParse("8Gi"),
corev1.ResourcePods: resource.MustParse("10"),
},
},
})
onDemandNodeClaim.StatusConditions().SetTrue(v1.ConditionTypeConsolidatable)
rs := test.ReplicaSet()
ExpectApplied(ctx, env.Client, rs)
Expect(env.Client.Get(ctx, client.ObjectKeyFromObject(rs), rs)).To(Succeed())
pods := test.Pods(2, test.PodOptions{
ResourceRequirements: corev1.ResourceRequirements{
Requests: corev1.ResourceList{
corev1.ResourceCPU: resource.MustParse("3"),
corev1.ResourceMemory: resource.MustParse("100Mi"),
},
},
ObjectMeta: metav1.ObjectMeta{Labels: labels,
OwnerReferences: []metav1.OwnerReference{
{
APIVersion: "apps/v1",
Kind: "ReplicaSet",
Name: rs.Name,
UID: rs.UID,
Controller: lo.ToPtr(true),
BlockOwnerDeletion: lo.ToPtr(true),
},
}},
})
ExpectApplied(ctx, env.Client, rs, pods[0], pods[1], onDemandNode, onDemandNodeClaim, nodePool)
ExpectManualBinding(ctx, env.Client, pods[0], onDemandNode)
ExpectManualBinding(ctx, env.Client, pods[1], onDemandNode)
ExpectMakeNodesAndNodeClaimsInitializedAndStateUpdated(ctx, env.Client, nodeStateController, nodeClaimStateController, []*corev1.Node{onDemandNode}, []*v1.NodeClaim{onDemandNodeClaim})

c := disruption.MakeConsolidation(fakeClock, cluster, env.Client, prov, cloudProvider, recorder, queue)
singleNodeConsolidation := disruption.NewSingleNodeConsolidation(c, disruption.WithValidator(NopValidator{}))
budgets, err := disruption.BuildDisruptionBudgetMapping(ctx, cluster, fakeClock, env.Client, cloudProvider, recorder, singleNodeConsolidation.Reason())
Expect(err).To(Succeed())
candidates, err := disruption.GetCandidates(ctx, cluster, env.Client, recorder, fakeClock, cloudProvider, singleNodeConsolidation.ShouldDisrupt, singleNodeConsolidation.Class(), queue)
Expect(err).To(Succeed())
onDemandCandidate, ok := lo.Find(candidates, func(candidate *disruption.Candidate) bool {
return candidate.NodeClaim.Name == onDemandNodeClaim.Name
})
Expect(ok).To(BeTrue())

cmds, err := singleNodeConsolidation.ComputeCommands(ctx, budgets, onDemandCandidate)
Expect(err).To(Succeed())
Expect(cmds).To(HaveLen(1))
Expect(cmds[0].Replacements).To(HaveLen(2))
for _, replacement := range cmds[0].Replacements {
Expect(replacement.Requirements.Get(v1.CapacityTypeLabelKey).Values()).To(ConsistOf(v1.CapacityTypeReserved))
Expect(replacement.Requirements.Get(cloudprovider.ReservationIDLabel).Values()).To(ConsistOf(reservationID))
}

validator := disruption.NewSingleConsolidationValidator(c)
_, err = validator.Validate(ctx, cmds[0], 0)
Expect(err).ToNot(HaveOccurred())

cmds[0].Replacements[0].Requirements[v1.CapacityTypeLabelKey] = scheduling.NewRequirement(v1.CapacityTypeLabelKey, corev1.NodeSelectorOpIn, v1.CapacityTypeOnDemand)
_, err = validator.Validate(ctx, cmds[0], 0)
Expect(err).To(HaveOccurred())
Expect(disruption.IsValidationError(err)).To(BeTrue())
})
})
Context("Preferences", func() {
It("should consolidate a node through deletion when ignoring preferences", func() {
Expand Down
Loading
Loading