diff --git a/pkg/api/agent_init_test.go b/pkg/api/agent_init_test.go new file mode 100644 index 0000000000..419ae05bb8 --- /dev/null +++ b/pkg/api/agent_init_test.go @@ -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" +) diff --git a/pkg/api/gpu_utilization_worker.go b/pkg/api/gpuworker/worker.go similarity index 95% rename from pkg/api/gpu_utilization_worker.go rename to pkg/api/gpuworker/worker.go index 9b7ade760d..089abfa3c4 100644 --- a/pkg/api/gpu_utilization_worker.go +++ b/pkg/api/gpuworker/worker.go @@ -1,4 +1,4 @@ -package api +package gpuworker import ( "context" @@ -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 @@ -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 { @@ -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, @@ -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() @@ -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) @@ -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 } @@ -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 { @@ -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, @@ -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 { diff --git a/pkg/api/gpu_utilization_worker_test.go b/pkg/api/gpuworker/worker_test.go similarity index 91% rename from pkg/api/gpu_utilization_worker_test.go rename to pkg/api/gpuworker/worker_test.go index a1494cbbc6..be4d319463 100644 --- a/pkg/api/gpu_utilization_worker_test.go +++ b/pkg/api/gpuworker/worker_test.go @@ -1,4 +1,4 @@ -package api +package gpuworker import ( "context" @@ -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() @@ -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{ @@ -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() @@ -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() @@ -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() @@ -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{ @@ -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{ @@ -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) @@ -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") @@ -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") @@ -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 ( @@ -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}, diff --git a/pkg/api/server.go b/pkg/api/server.go index 31a8e3b06a..8fb4cd3bed 100644 --- a/pkg/api/server.go +++ b/pkg/api/server.go @@ -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" @@ -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") diff --git a/pkg/api/server_runtime.go b/pkg/api/server_runtime.go index 42deb7620d..f93aecda70 100644 --- a/pkg/api/server_runtime.go +++ b/pkg/api/server_runtime.go @@ -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" ) @@ -36,7 +37,7 @@ type authRuntime struct { } type backgroundServices struct { - gpuUtilWorker *GPUUtilizationWorker + gpuUtilWorker *gpuworker.Worker workloadHandlers *workloads.WorkloadHandlers rewardsHandler *rewards.RewardsHandler }