Sitelet https://github.com/kubernetes/kops/pull/18100/files
Skip to content
Merged
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
2 changes: 1 addition & 1 deletion tests/e2e/scenarios/ai-conformance/validators/kube.go
Original file line number Diff line number Diff line change
Expand Up @@ -205,7 +205,7 @@ func (h *ValidatorHarness) TestNamespace() string {
}

// Wait for namespace deletion to complete so that we don't have leftover namespaces consuming resources.
if err := wait.PollUntilContextTimeout(ctx, 2*time.Second, 5*time.Minute, false, func(ctx context.Context) (done bool, err error) {
if err := wait.PollUntilContextTimeout(ctx, 2*time.Second, 10*time.Minute, false, func(ctx context.Context) (done bool, err error) {
if _, err := h.DynamicClient().Resource(namespaceGVR).Get(ctx, ns, metav1.GetOptions{}); err != nil {
if apierrors.IsNotFound(err) {
return true, nil
Expand Down
8 changes: 4 additions & 4 deletions tests/e2e/scenarios/ai-conformance/validators/output.go
Original file line number Diff line number Diff line change
Expand Up @@ -60,16 +60,16 @@ func (h *ValidatorHarness) Log(s string) {
func (h *ValidatorHarness) Fatalf(format string, args ...interface{}) {
s := fmt.Sprintf(format, args...)

h.output.WriteText("FATAL: " + s)
h.t.Fatalf(format, args...)
h.output.WriteText("FAIL: " + s)
h.t.Fatalf("FAIL: "+format, args...)
}

// Errorf is like t.Errorf, but also writes to the sinks.
func (h *ValidatorHarness) Errorf(format string, args ...interface{}) {
s := fmt.Sprintf(format, args...)

h.output.WriteText("ERROR: " + s)
h.t.Errorf(format, args...)
h.output.WriteText("FAIL: " + s)
h.t.Errorf("FAIL: "+format, args...)
}

// Run is like t.Run, but creates a sub-harness that shares the output.
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,181 @@
/*
Copyright The Kubernetes Authors.

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 clusterautoscaling

import (
"fmt"
"strings"
"testing"
"time"

"k8s.io/kops/tests/e2e/scenarios/ai-conformance/validators"
)

// Test_SchedulingOrchestration_ClusterAutoscaling verifies that the cluster autoscaler
// (or equivalent mechanism) can scale up node groups containing GPU accelerators
// based on pending pods requesting those accelerators.
//
// It counts the current number of GPU nodes, deploys N+1 replicas of a simple
// GPU workload (each requesting one GPU), and verifies that the cluster
// scales up to accommodate the additional pod.
func Test_SchedulingOrchestration_ClusterAutoscaling(t *testing.T) {
// Description:
// If the platform provides a cluster autoscaler or an equivalent mechanism,
// it must be able to scale up/down node groups containing specific accelerator types
// based on pending pods requesting those accelerators.

h := validators.NewValidatorHarness(t)

h.Logf("# Cluster Autoscaling for GPU Nodes")

h.Run("cluster-autoscaling-gpu", func(h *validators.ValidatorHarness) {
ns := h.TestNamespace()

// Count the current number of nodes with GPUs by looking at resource slices
// that advertise GPU devices.
h.Logf("## Determine current GPU node count")

listGPUNodes := func() []string {
result := h.ShellExec("kubectl get nodes -l nvidia.com/gpu.present=true -o name")
nodes := strings.Split(strings.TrimSpace(result.Stdout()), "\n")
return nodes
}

initialGPUNodes := listGPUNodes()
h.Logf("Found %d GPU nodes initially (%v)", len(initialGPUNodes), initialGPUNodes)

if len(initialGPUNodes) == 0 {
h.Fatalf("No GPU nodes found in the cluster; cannot test cluster autoscaling for GPUs")
}

// Deploy the GPU probe workload with 1 replica first.
h.Logf("## Deploy GPU probe workload")
h.ApplyManifest(ns, "testdata/cluster-autoscaling-workload.yaml")

// Scale to N+1 replicas to force the autoscaler to add a GPU node.
targetReplicas := len(initialGPUNodes) + 1
h.Logf("## Scale deployment to %d replicas (initial GPU nodes: %d)", targetReplicas, len(initialGPUNodes))
h.ShellExec(fmt.Sprintf("kubectl scale deployment/cluster-autoscaling-workload -n %s --replicas=%d", ns, targetReplicas))

// Wait for at least one pod to be Pending (confirming we need a new node).
h.Logf("### Verify at least one pod is pending")
h.ShellExec(fmt.Sprintf("kubectl get pods -n %s -l app=cluster-autoscaling-workload -o wide", ns))

// Poll for the GPU node count to increase.
h.Logf("## Wait for cluster to scale up")
var scaledUp bool
const maxAttempts = 40 // 40 * 30s = 20 minutes
for attempt := 1; attempt <= maxAttempts; attempt++ {
currentGPUNodes := listGPUNodes()

if len(currentGPUNodes) > len(initialGPUNodes) {
h.Logf("Cluster scaled up: GPU nodes increased from %d to %d on attempt %d", len(initialGPUNodes), len(currentGPUNodes), attempt)
scaledUp = true
break
}

// Periodic diagnostics.
if attempt%5 == 1 {
h.Logf("### Diagnostics at attempt %d", attempt)
h.ShellExec(fmt.Sprintf("kubectl get pods -n %s -l app=cluster-autoscaling-workload -o wide", ns))
h.ShellExec("kubectl get nodes -o wide")
}

if attempt < maxAttempts {
h.Logf("Attempt %d: GPU node count is still %d (need > %d), waiting 30s...", attempt, len(currentGPUNodes), len(initialGPUNodes))
time.Sleep(30 * time.Second)
}
}

if !scaledUp {
// Failure diagnostics.
h.Logf("### Failure diagnostics")
h.ShellExec(fmt.Sprintf("kubectl get pods -n %s -l app=cluster-autoscaling-workload -o wide", ns))
h.ShellExec(fmt.Sprintf("kubectl describe pods -n %s -l app=cluster-autoscaling-workload", ns))
h.ShellExec("kubectl get nodes -o wide")
h.ShellExec("kubectl describe nodes")
h.Errorf("Cluster did not scale up GPU nodes within the expected time (initial: %d)", len(initialGPUNodes))
}

// Verify all replicas eventually become ready.
if scaledUp {
h.Logf("## Wait for all replicas to be ready")
result := h.ShellExec(fmt.Sprintf(
"kubectl rollout status deployment/cluster-autoscaling-workload -n %s --timeout=600s",
ns,
))
if result.Err() != nil {
h.Errorf("Deployment did not become fully ready: %v", result.Err())
}

// Verify GPU pods are actually running nvidia-smi.
h.Logf("### Verify GPU pods are running")
h.ShellExec(fmt.Sprintf("kubectl get pods -n %s -l app=cluster-autoscaling-workload -o wide", ns))

// Check logs from one of the pods to confirm GPU access.
podListResult := h.ShellExec(fmt.Sprintf(
"kubectl get pods -n %s -l app=cluster-autoscaling-workload -o name",
ns,
))
for _, podName := range strings.Split(strings.TrimSpace(podListResult.Stdout()), "\n") {
h.ShellExec(fmt.Sprintf("kubectl logs -n %s %s --tail=5", ns, podName))
}

h.Success("Cluster autoscaler scaled up GPU nodes from %d to accommodate %d GPU pods", len(initialGPUNodes), targetReplicas)
}

// Scale down and verify the cluster scales back down.
h.Logf("## Scale down and verify cluster scale-down")
h.ShellExec(fmt.Sprintf("kubectl scale deployment/cluster-autoscaling-workload -n %s --replicas=0", ns))

h.Logf("Waiting for cluster to scale down (this may take several minutes)...")
var scaledDown bool
const scaleDownMaxAttempts = 40 // 40 * 30s = 20 minutes
for attempt := 1; attempt <= scaleDownMaxAttempts; attempt++ {
currentGPUNodes := listGPUNodes()

if len(currentGPUNodes) <= len(initialGPUNodes) {
h.Logf("Cluster scaled down: GPU nodes decreased to %d on attempt %d", len(currentGPUNodes), attempt)
scaledDown = true
break
}

if attempt%5 == 1 {
h.Logf("### Scale-down diagnostics at attempt %d", attempt)
h.ShellExec("kubectl get nodes -o wide")
}

if attempt < scaleDownMaxAttempts {
h.Logf("Attempt %d: GPU node count is still %d (need <= %d), waiting 30s...", attempt, len(currentGPUNodes), len(initialGPUNodes))
time.Sleep(30 * time.Second)
}
}

if !scaledDown {
h.Logf("### Scale-down failure diagnostics")
h.ShellExec("kubectl get nodes -o wide")
h.ShellExec("kubectl describe nodes")
h.Errorf("Cluster did not scale down GPU nodes within the expected time")
} else {
h.Success("Cluster autoscaler scaled down GPU nodes back to %d", len(initialGPUNodes))
}
})

if h.AllPassed() {
h.RecordConformance("schedulingOrchestration", "cluster_autoscaling")
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,44 @@
# Deployment that requests a GPU via DRA and runs nvidia-smi periodically.
# The test creates this with N+1 replicas (where N is the current GPU node count),
# forcing the cluster autoscaler to provision an additional GPU node.


apiVersion: apps/v1
kind: Deployment
metadata:
name: cluster-autoscaling-workload
labels:
app: cluster-autoscaling-workload
spec:
# replicas is set dynamically by the test via kubectl scale
replicas: 1
selector:
matchLabels:
app: cluster-autoscaling-workload
template:
metadata:
labels:
app: cluster-autoscaling-workload
spec:
terminationGracePeriodSeconds: 5
tolerations:
- key: "nvidia.com/gpu"
operator: "Exists"
effect: "NoSchedule"
containers:
- name: cluster-autoscaling-workload
image: nvcr.io/nvidia/k8s/cuda-sample:vectoradd-cuda12.5.0
command:
- "/bin/sh"
- "-c"
- |
while true; do
echo "$(date -Iseconds) GPU probe alive on $(hostname)"
nvidia-smi --query-gpu=name,temperature.gpu,utilization.gpu,memory.used,memory.total --format=csv,noheader
sleep 60
done
resources:
requests:
nvidia.com/gpu: 1
limits:
nvidia.com/gpu: 1
Loading