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
16 changes: 16 additions & 0 deletions pkg/api/agent_init_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,16 @@
package api

// Blank-import pkg/agent so its init() runs during pkg/api's test binary,
// which assigns ai.SetClusterContextProviders — called unconditionally in
// NewServer (server.go). Before pkg/api/gpu_utilization_worker_test.go was
// moved to pkg/api/gpuworker, that file transitively imported pkg/agent and
// pulled the init() into every pkg/api test run. This file preserves the
// exact same test-binary side effect without pinning any real dependency.
//
// The underlying hidden coupling (production callers relying on a package
// they never import to have init'd a global var) is tracked separately —
// see kubestellar/console#23735 comment thread.

import (
_ "github.com/kubestellar/console/pkg/agent"
)
24 changes: 12 additions & 12 deletions pkg/api/gpu_utilization_worker.go → pkg/api/gpuworker/worker.go
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
package api
package gpuworker

import (
"context"
Expand Down Expand Up @@ -54,8 +54,8 @@ const (
maxConcurrentReservationCollectors = 10
)

// GPUUtilizationWorker periodically collects GPU utilization data for active reservations
type GPUUtilizationWorker struct {
// Worker periodically collects GPU utilization data for active reservations
type Worker struct {
store store.Store
k8sClient *k8s.MultiClusterClient
interval time.Duration
Expand All @@ -77,8 +77,8 @@ type GPUUtilizationWorker struct {
dcgmService string
}

// NewGPUUtilizationWorker creates a new GPU utilization worker
func NewGPUUtilizationWorker(s store.Store, k8sClient *k8s.MultiClusterClient, notificationService *notifications.Service) *GPUUtilizationWorker {
// New creates a new GPU utilization worker
func New(s store.Store, k8sClient *k8s.MultiClusterClient, notificationService *notifications.Service) *Worker {
intervalMs := defaultUtilPollIntervalMs
if envVal := os.Getenv("GPU_UTIL_POLL_INTERVAL_MS"); envVal != "" {
if parsed, err := strconv.Atoi(envVal); err == nil && parsed > 0 {
Expand Down Expand Up @@ -113,7 +113,7 @@ func NewGPUUtilizationWorker(s store.Store, k8sClient *k8s.MultiClusterClient, n
}

ctx, cancel := context.WithCancel(context.Background())
return &GPUUtilizationWorker{
return &Worker{
store: s,
k8sClient: k8sClient,
interval: time.Duration(intervalMs) * time.Millisecond,
Expand All @@ -131,7 +131,7 @@ func NewGPUUtilizationWorker(s store.Store, k8sClient *k8s.MultiClusterClient, n
}

// Start begins the background polling loop
func (w *GPUUtilizationWorker) Start() {
func (w *Worker) Start() {
w.wg.Add(1)
safego.GoWith("gpu-utilization-worker", func() {
defer w.wg.Done()
Expand All @@ -158,7 +158,7 @@ func (w *GPUUtilizationWorker) Start() {

// Stop signals the worker to stop and waits for the background goroutine to exit.
// It is safe to call multiple times; only the first call actually closes the stop channel.
func (w *GPUUtilizationWorker) Stop() {
func (w *Worker) Stop() {
w.stopOnce.Do(func() {
w.baseCancel() // cancel all in-flight Kubernetes API calls (#6966)
close(w.stopCh)
Expand All @@ -167,7 +167,7 @@ func (w *GPUUtilizationWorker) Stop() {
}

// collectUtilization queries active reservations and records utilization snapshots
func (w *GPUUtilizationWorker) collectUtilization() {
func (w *Worker) collectUtilization() {
if w.k8sClient == nil {
return
}
Expand Down Expand Up @@ -219,7 +219,7 @@ func (w *GPUUtilizationWorker) collectUtilization() {
// per-namespace framebuffer utilization from the nested map. Returns
// nil when DCGM is disabled via env flag — callers handle nil as
// "no DCGM data, use legacy zero fallback".
func (w *GPUUtilizationWorker) scrapeDCGMPerCluster(
func (w *Worker) scrapeDCGMPerCluster(
reservations []models.GPUReservation,
timeout time.Duration,
) map[string]map[string]*gpu.NamespaceMetrics {
Expand Down Expand Up @@ -266,7 +266,7 @@ func (w *GPUUtilizationWorker) scrapeDCGMPerCluster(
// dcgmClusterMetrics is the per-namespace DCGM framebuffer map for the
// reservation's cluster (or nil when DCGM is disabled / unreachable);
// callers pass nil for the legacy zero-memory fallback.
func (w *GPUUtilizationWorker) collectForReservation(
func (w *Worker) collectForReservation(
ctx context.Context,
reservation *models.GPUReservation,
dcgmClusterMetrics map[string]*gpu.NamespaceMetrics,
Expand Down Expand Up @@ -400,7 +400,7 @@ func (w *GPUUtilizationWorker) collectForReservation(
}

// cleanupOldSnapshots removes snapshots older than the retention period
func (w *GPUUtilizationWorker) cleanupOldSnapshots() {
func (w *Worker) cleanupOldSnapshots() {
cutoff := time.Now().AddDate(0, 0, -snapshotRetentionDays)
deleted, err := w.store.DeleteOldUtilizationSnapshots(w.baseCtx, cutoff)
if err != nil {
Expand Down
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
package api
package gpuworker

import (
"context"
Expand All @@ -21,7 +21,7 @@ import (
"k8s.io/client-go/tools/clientcmd/api"
)

func TestGPUUtilizationWorker_MetricAccuracy(t *testing.T) {
func TestWorker_MetricAccuracy(t *testing.T) {
mockStore := new(test.MockStore)
k8sClient, _ := k8s.NewMultiClusterClient("")
fakeClient := k8sfake.NewSimpleClientset()
Expand All @@ -35,7 +35,7 @@ func TestGPUUtilizationWorker_MetricAccuracy(t *testing.T) {
},
})

worker := NewGPUUtilizationWorker(mockStore, k8sClient, nil)
worker := New(mockStore, k8sClient, nil)

t.Run("collectForReservation - Accurate GPU calculation", func(t *testing.T) {
reservation := &models.GPUReservation{
Expand Down Expand Up @@ -172,9 +172,9 @@ func TestGPUUtilizationWorker_MetricAccuracy(t *testing.T) {
})
}

func TestGPUUtilizationWorker_Cleanup(t *testing.T) {
func TestWorker_Cleanup(t *testing.T) {
mockStore := new(test.MockStore)
worker := NewGPUUtilizationWorker(mockStore, nil, nil)
worker := New(mockStore, nil, nil)

t.Run("cleanupOldSnapshots calls store", func(t *testing.T) {
mockStore.On("DeleteOldUtilizationSnapshots", mock.Anything).Return(int64(5), nil).Once()
Expand All @@ -183,16 +183,16 @@ func TestGPUUtilizationWorker_Cleanup(t *testing.T) {
})
}

func TestGPUUtilizationWorker_IntervalFromEnv(t *testing.T) {
func TestWorker_IntervalFromEnv(t *testing.T) {
os.Setenv("GPU_UTIL_POLL_INTERVAL_MS", "5000")
defer os.Unsetenv("GPU_UTIL_POLL_INTERVAL_MS")

worker := NewGPUUtilizationWorker(nil, nil, nil)
worker := New(nil, nil, nil)
assert.Equal(t, 5*time.Second, worker.interval)
}

func TestGPUUtilizationWorker_StopCancel(t *testing.T) {
worker := NewGPUUtilizationWorker(nil, nil, nil)
func TestWorker_StopCancel(t *testing.T) {
worker := New(nil, nil, nil)

ctx := worker.baseCtx
worker.Stop()
Expand All @@ -206,7 +206,7 @@ func TestGPUUtilizationWorker_StopCancel(t *testing.T) {
}
}

func TestGPUUtilizationWorker_ThresholdAlerting(t *testing.T) {
func TestWorker_ThresholdAlerting(t *testing.T) {
mockStore := new(test.MockStore)
notificationService := notifications.NewService()

Expand All @@ -226,7 +226,7 @@ func TestGPUUtilizationWorker_ThresholdAlerting(t *testing.T) {
},
})

worker := NewGPUUtilizationWorker(mockStore, k8sClient, notificationService)
worker := New(mockStore, k8sClient, notificationService)
assert.Equal(t, 80.0, worker.overThreshold)

reservation := &models.GPUReservation{
Expand Down Expand Up @@ -278,7 +278,7 @@ func TestGPUUtilizationWorker_ThresholdAlerting(t *testing.T) {
},
})

worker := NewGPUUtilizationWorker(mockStore, k8sClient, notificationService)
worker := New(mockStore, k8sClient, notificationService)
assert.Equal(t, 50.0, worker.underThreshold)

reservation := &models.GPUReservation{
Expand Down Expand Up @@ -332,7 +332,7 @@ func TestGPUUtilizationWorker_ThresholdAlerting(t *testing.T) {
},
})

worker := NewGPUUtilizationWorker(mockStore, k8sClient, notificationService)
worker := New(mockStore, k8sClient, notificationService)
assert.Equal(t, 90.0, worker.overThreshold)
assert.Equal(t, 10.0, worker.underThreshold)

Expand Down Expand Up @@ -372,13 +372,13 @@ func TestGPUUtilizationWorker_ThresholdAlerting(t *testing.T) {

// Issue 9135 — DCGM GPU memory integration tests.

func TestGPUUtilizationWorker_DCGMDisabled_MemoryZero(t *testing.T) {
func TestWorker_DCGMDisabled_MemoryZero(t *testing.T) {
t.Setenv("GPU_METRICS_DCGM_ENABLED", "")

mockStore := new(test.MockStore)
k8sClient, _ := k8s.NewMultiClusterClient("")
k8sClient.InjectClient("c1", k8sfake.NewSimpleClientset())
worker := NewGPUUtilizationWorker(mockStore, k8sClient, nil)
worker := New(mockStore, k8sClient, nil)

if worker.dcgmEnabled {
t.Fatal("expected dcgmEnabled=false when GPU_METRICS_DCGM_ENABLED is unset")
Expand All @@ -404,14 +404,14 @@ func TestGPUUtilizationWorker_DCGMDisabled_MemoryZero(t *testing.T) {
mockStore.AssertExpectations(t)
}

func TestGPUUtilizationWorker_DCGMEnabled_EnvOverrides(t *testing.T) {
func TestWorker_DCGMEnabled_EnvOverrides(t *testing.T) {
t.Setenv("GPU_METRICS_DCGM_ENABLED", "true")
t.Setenv("GPU_METRICS_DCGM_NAMESPACE", "custom-ns")
t.Setenv("GPU_METRICS_DCGM_SERVICE", "custom-svc")

mockStore := new(test.MockStore)
k8sClient, _ := k8s.NewMultiClusterClient("")
worker := NewGPUUtilizationWorker(mockStore, k8sClient, nil)
worker := New(mockStore, k8sClient, nil)

if !worker.dcgmEnabled {
t.Fatal("expected dcgmEnabled=true when GPU_METRICS_DCGM_ENABLED=true")
Expand All @@ -424,14 +424,14 @@ func TestGPUUtilizationWorker_DCGMEnabled_EnvOverrides(t *testing.T) {
}
}

func TestGPUUtilizationWorker_DCGMEnabled_MemoryFromScraper(t *testing.T) {
func TestWorker_DCGMEnabled_MemoryFromScraper(t *testing.T) {
// Pass DCGM metrics directly to collectForReservation to verify the
// percentage computation and that non-matching namespaces fall back to 0.
mockStore := new(test.MockStore)
k8sClient, _ := k8s.NewMultiClusterClient("")
fakeClient := k8sfake.NewSimpleClientset()
k8sClient.InjectClient("c1", fakeClient)
worker := NewGPUUtilizationWorker(mockStore, k8sClient, nil)
worker := New(mockStore, k8sClient, nil)

// 75% framebuffer utilization: 30720 used out of 30720+10240 total.
const (
Expand All @@ -457,13 +457,13 @@ func TestGPUUtilizationWorker_DCGMEnabled_MemoryFromScraper(t *testing.T) {
mockStore.AssertExpectations(t)
}

func TestGPUUtilizationWorker_DCGMEnabled_NamespaceMiss_Zero(t *testing.T) {
func TestWorker_DCGMEnabled_NamespaceMiss_Zero(t *testing.T) {
// DCGM returned data, but not for this reservation's namespace.
mockStore := new(test.MockStore)
k8sClient, _ := k8s.NewMultiClusterClient("")
fakeClient := k8sfake.NewSimpleClientset()
k8sClient.InjectClient("c1", fakeClient)
worker := NewGPUUtilizationWorker(mockStore, k8sClient, nil)
worker := New(mockStore, k8sClient, nil)

dcgmByNs := map[string]*agent.DCGMNamespaceMetrics{
"other-ns": {FBUsedMiB: 1000, FBFreeMiB: 1000, SampleCount: 1},
Expand Down
3 changes: 2 additions & 1 deletion pkg/api/server.go
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@ import (

"github.com/kubestellar/console/pkg/ai"
"github.com/kubestellar/console/pkg/api/audit"
"github.com/kubestellar/console/pkg/api/gpuworker"
"github.com/kubestellar/console/pkg/api/metrics"
"github.com/kubestellar/console/pkg/api/middleware"
"github.com/kubestellar/console/pkg/api/tracing"
Expand Down Expand Up @@ -269,7 +270,7 @@ func NewServer(cfg Config) (*Server, error) {

// Start GPU utilization background worker (collects hourly snapshots)
if k8sClient != nil {
server.background.gpuUtilWorker = NewGPUUtilizationWorker(db, k8sClient, notificationService)
server.background.gpuUtilWorker = gpuworker.New(db, k8sClient, notificationService)
server.background.gpuUtilWorker.Start()
} else {
slog.Info("[Server] GPU utilization worker skipped — no Kubernetes client available")
Expand Down
3 changes: 2 additions & 1 deletion pkg/api/server_runtime.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import (
"github.com/kubestellar/console/pkg/api/handlers/auth"
"github.com/kubestellar/console/pkg/api/handlers/rewards"
"github.com/kubestellar/console/pkg/api/handlers/workloads"
"github.com/kubestellar/console/pkg/api/gpuworker"
"github.com/kubestellar/console/pkg/api/middleware"
"github.com/kubestellar/console/pkg/k8s"
)
Expand All @@ -36,7 +37,7 @@ type authRuntime struct {
}

type backgroundServices struct {
gpuUtilWorker *GPUUtilizationWorker
gpuUtilWorker *gpuworker.Worker
workloadHandlers *workloads.WorkloadHandlers
rewardsHandler *rewards.RewardsHandler
}
Expand Down
Loading