main.go 28.4 KB
Newer Older
1
/*
2
 * SPDX-FileCopyrightText: Copyright (c) 2022 Atalaya Tech. Inc
3
 * SPDX-FileCopyrightText: Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved.
4
5
6
7
8
9
10
11
12
13
14
15
16
 * SPDX-License-Identifier: Apache-2.0
 *
 * 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.
17
 * Modifications Copyright (c) 2025-2026 NVIDIA CORPORATION & AFFILIATES
18
19
20
21
22
 */

package main

import (
23
	"context"
24
25
	"crypto/tls"
	"flag"
26
	"fmt"
27
	"net/http"
28
	"os"
29
	"strings"
30
	"time"
31
32
33

	// Import all Kubernetes client auth plugins (e.g. Azure, GCP, OIDC, etc.)
	// to ensure that exec-entrypoint and run can make use of them.
34

35
	admissionregistrationv1 "k8s.io/api/admissionregistration/v1"
36
	corev1 "k8s.io/api/core/v1"
37
	apiextensionsv1 "k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1"
38
39
	"k8s.io/client-go/discovery/cached/memory"
	"k8s.io/client-go/dynamic"
40
41
	"k8s.io/client-go/informers"
	"k8s.io/client-go/kubernetes"
42
	_ "k8s.io/client-go/plugin/pkg/client/auth"
43
44
	"k8s.io/client-go/restmapper"
	"k8s.io/client-go/scale"
45
	k8sCache "k8s.io/client-go/tools/cache"
46
	"k8s.io/utils/ptr"
47
	"sigs.k8s.io/controller-runtime/pkg/cache"
48
	"sigs.k8s.io/controller-runtime/pkg/client"
49

50
51
	k8sruntime "k8s.io/apimachinery/pkg/runtime"
	"k8s.io/apimachinery/pkg/runtime/serializer"
52
53
54
55
56
	utilruntime "k8s.io/apimachinery/pkg/util/runtime"
	clientgoscheme "k8s.io/client-go/kubernetes/scheme"
	ctrl "sigs.k8s.io/controller-runtime"
	"sigs.k8s.io/controller-runtime/pkg/healthz"
	"sigs.k8s.io/controller-runtime/pkg/log/zap"
57
	metricsfilters "sigs.k8s.io/controller-runtime/pkg/metrics/filters"
58
59
60
	metricsserver "sigs.k8s.io/controller-runtime/pkg/metrics/server"
	"sigs.k8s.io/controller-runtime/pkg/webhook"

61
62
63
	lwsscheme "sigs.k8s.io/lws/client-go/clientset/versioned/scheme"
	volcanoscheme "volcano.sh/apis/pkg/client/clientset/versioned/scheme"

64
	semver "github.com/Masterminds/semver/v3"
65
66
	configv1alpha1 "github.com/ai-dynamo/dynamo/deploy/operator/api/config/v1alpha1"
	configvalidation "github.com/ai-dynamo/dynamo/deploy/operator/api/config/validation"
67
	nvidiacomv1alpha1 "github.com/ai-dynamo/dynamo/deploy/operator/api/v1alpha1"
68
	nvidiacomv1beta1 "github.com/ai-dynamo/dynamo/deploy/operator/api/v1beta1"
69
	internalcert "github.com/ai-dynamo/dynamo/deploy/operator/internal/cert"
70
71
	"github.com/ai-dynamo/dynamo/deploy/operator/internal/controller"
	commonController "github.com/ai-dynamo/dynamo/deploy/operator/internal/controller_common"
72
	"github.com/ai-dynamo/dynamo/deploy/operator/internal/gpu"
73
74
	"github.com/ai-dynamo/dynamo/deploy/operator/internal/modelendpoint"
	"github.com/ai-dynamo/dynamo/deploy/operator/internal/namespace_scope"
75
	"github.com/ai-dynamo/dynamo/deploy/operator/internal/observability"
76
77
78
79
	"github.com/ai-dynamo/dynamo/deploy/operator/internal/rbac"
	"github.com/ai-dynamo/dynamo/deploy/operator/internal/secret"
	"github.com/ai-dynamo/dynamo/deploy/operator/internal/secrets"
	internalwebhook "github.com/ai-dynamo/dynamo/deploy/operator/internal/webhook"
80
	webhookdefaulting "github.com/ai-dynamo/dynamo/deploy/operator/internal/webhook/defaulting"
81
	webhookvalidation "github.com/ai-dynamo/dynamo/deploy/operator/internal/webhook/validation"
82
	grovev1alpha1 "github.com/ai-dynamo/grove/operator/api/core/v1alpha1"
83
	istioclientsetscheme "istio.io/client-go/pkg/clientset/versioned/scheme"
84
	gaiev1 "sigs.k8s.io/gateway-api-inference-extension/api/v1"
85
86
87
88
	//+kubebuilder:scaffold:imports
)

var (
89
90
91
	crdScheme    = k8sruntime.NewScheme()
	setupLog     = ctrl.Log.WithName("setup")
	configScheme = k8sruntime.NewScheme()
92
93
)

94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
// LoadAndValidateOperatorConfig loads the operator configuration from a file,
// applies defaults via the scheme, and validates it.
func LoadAndValidateOperatorConfig(path string) (*configv1alpha1.OperatorConfiguration, error) {
	data, err := os.ReadFile(path)
	if err != nil {
		return nil, fmt.Errorf("failed to read config file %s: %w", path, err)
	}

	codecFactory := serializer.NewCodecFactory(configScheme)
	cfg := &configv1alpha1.OperatorConfiguration{}
	if err := k8sruntime.DecodeInto(codecFactory.UniversalDecoder(), data, cfg); err != nil {
		return nil, fmt.Errorf("failed to decode config file %s: %w", path, err)
	}

	// Validate the configuration
	if errs := configvalidation.ValidateOperatorConfiguration(cfg); len(errs) > 0 {
		return nil, fmt.Errorf("config validation failed: %s", errs.ToAggregate().Error())
	}

	return cfg, nil
}

116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
func createScalesGetter(mgr ctrl.Manager) (scale.ScalesGetter, error) {
	config := mgr.GetConfig()

	// Create kubernetes client for discovery
	kubeClient, err := kubernetes.NewForConfig(config)
	if err != nil {
		return nil, err
	}

	// Create cached discovery client
	cachedDiscovery := memory.NewMemCacheClient(kubeClient.Discovery())

	// Create REST mapper
	restMapper := restmapper.NewDeferredDiscoveryRESTMapper(cachedDiscovery)

	scalesGetter, err := scale.NewForConfig(
		config,
		restMapper,
		dynamic.LegacyAPIPathResolverFunc,
		scale.NewDiscoveryScaleKindResolver(cachedDiscovery),
	)
	if err != nil {
		return nil, err
	}

	return scalesGetter, nil
}

144
145
146
147
func initCRDSchemes() {
	utilruntime.Must(clientgoscheme.AddToScheme(crdScheme))

	utilruntime.Must(nvidiacomv1alpha1.AddToScheme(crdScheme))
148

149
	utilruntime.Must(nvidiacomv1beta1.AddToScheme(crdScheme))
150

151
	utilruntime.Must(lwsscheme.AddToScheme(crdScheme))
152

153
	utilruntime.Must(volcanoscheme.AddToScheme(crdScheme))
154

155
	utilruntime.Must(grovev1alpha1.AddToScheme(crdScheme))
156

157
	utilruntime.Must(apiextensionsv1.AddToScheme(crdScheme))
158

159
160
	utilruntime.Must(admissionregistrationv1.AddToScheme(crdScheme))

161
	utilruntime.Must(istioclientsetscheme.AddToScheme(crdScheme))
162

163
	utilruntime.Must(gaiev1.Install(crdScheme))
164
165
166
	//+kubebuilder:scaffold:scheme
}

167
168
169
170
func initConfigScheme() {
	utilruntime.Must(configv1alpha1.AddToScheme(configScheme))
}

171
172
173
// +kubebuilder:rbac:groups=authentication.k8s.io,resources=tokenreviews,verbs=create
// +kubebuilder:rbac:groups=authorization.k8s.io,resources=subjectaccessreviews,verbs=create

174
//nolint:gocyclo
175
func main() {
176
177
178
179
	initCRDSchemes()
	initConfigScheme()

	var configFile string
180
	var operatorVersion string
181
	flag.StringVar(&configFile, "config", "", "Path to operator configuration file (required)")
182
183
	flag.StringVar(&operatorVersion, "operator-version", "unknown",
		"Version of the operator (used in lease holder identity)")
184
185
186
187
188
	opts := zap.Options{
		Development: true,
	}
	opts.BindFlags(flag.CommandLine)
	flag.Parse()
189
	ctrl.SetLogger(zap.New(zap.UseFlagOptions(&opts)))
190

191
192
	if configFile == "" {
		setupLog.Error(nil, "--config flag is required")
193
194
195
		os.Exit(1)
	}

196
197
198
199
	// Load, default, and validate operator configuration
	operatorCfg, err := LoadAndValidateOperatorConfig(configFile)
	if err != nil {
		setupLog.Error(err, "failed to load operator configuration", "configFile", configFile)
200
201
		os.Exit(1)
	}
202
	setupLog.Info("Operator configuration loaded successfully", "configFile", configFile)
203

204
205
206
207
	// Validate and normalize operator version to semver
	if _, err := semver.NewVersion(operatorVersion); err != nil {
		setupLog.Error(err, "operator-version is not valid semver",
			"provided", operatorVersion, "error", err.Error())
208
209
		os.Exit(1)
	}
210
	setupLog.Info("Operator version configured", "version", operatorVersion)
211

212
213
	// Initialize runtime config (will be populated after detection)
	runtimeConfig := &commonController.RuntimeConfig{}
214

215
	mainCtx := ctrl.SetupSignalHandler()
216
217
218
219
220
221
222
223
224
225
226
227
228

	// if the enable-http2 flag is false (the default), http/2 should be disabled
	// due to its vulnerabilities. More specifically, disabling http/2 will
	// prevent from being vulnerable to the HTTP/2 Stream Cancellation and
	// Rapid Reset CVEs. For more information see:
	// - https://github.com/advisories/GHSA-qppj-fm5r-hxr3
	// - https://github.com/advisories/GHSA-4374-p667-p6c8
	disableHTTP2 := func(c *tls.Config) {
		setupLog.Info("disabling http/2")
		c.NextProtos = []string{"http/1.1"}
	}

	tlsOpts := []func(*tls.Config){}
229
	if !operatorCfg.Security.EnableHTTP2 {
230
231
232
233
		tlsOpts = append(tlsOpts, disableHTTP2)
	}

	webhookServer := webhook.NewServer(webhook.Options{
234
235
236
		Host:    operatorCfg.Server.Webhook.Host,
		Port:    operatorCfg.Server.Webhook.Port,
		CertDir: operatorCfg.Server.Webhook.CertDir,
237
238
239
		TLSOpts: tlsOpts,
	})

240
241
242
243
244
	metricsBindAddr := fmt.Sprintf("%s:%d", operatorCfg.Server.Metrics.BindAddress, operatorCfg.Server.Metrics.Port)
	healthProbeAddr := fmt.Sprintf(
		"%s:%d", operatorCfg.Server.HealthProbe.BindAddress, operatorCfg.Server.HealthProbe.Port,
	)

245
	mgrOpts := ctrl.Options{
246
		Scheme: crdScheme,
247
		Metrics: metricsserver.Options{
248
249
250
251
			BindAddress:    metricsBindAddr,
			SecureServing:  ptr.Deref(operatorCfg.Server.Metrics.Secure, true),
			FilterProvider: metricsfilters.WithAuthenticationAndAuthorization,
			TLSOpts:        tlsOpts,
252
		},
253
		WebhookServer:           webhookServer,
254
255
256
257
		HealthProbeBindAddress:  healthProbeAddr,
		LeaderElection:          operatorCfg.LeaderElection.Enabled,
		LeaderElectionID:        operatorCfg.LeaderElection.ID,
		LeaderElectionNamespace: operatorCfg.LeaderElection.Namespace,
258
	}
259
260

	restrictedNamespace := operatorCfg.Namespace.Restricted
261
262
263
264
	if restrictedNamespace != "" {
		mgrOpts.Cache.DefaultNamespaces = map[string]cache.Config{
			restrictedNamespace: {},
		}
265
		setupLog.Info("Restricted namespace configured, launching in restricted mode", "namespace", restrictedNamespace)
266
267
268
269
270
271
272
273
274
275

		banner := strings.Repeat("=", 80)
		setupLog.Error(nil, banner)
		setupLog.Error(nil, "DEPRECATION WARNING: Namespace-restricted mode is deprecated "+
			"and will be removed in a future release.")
		setupLog.Error(nil, "The operator is running in namespace-restricted mode",
			"namespace", restrictedNamespace)
		setupLog.Error(nil, "Please migrate to cluster-wide mode "+
			"by removing the namespaceRestriction configuration.")
		setupLog.Error(nil, banner)
276
277
	} else {
		setupLog.Info("No restricted namespace configured, launching in cluster-wide mode")
278
279
280
281
282
283
284
	}
	mgr, err := ctrl.NewManager(ctrl.GetConfigOrDie(), mgrOpts)
	if err != nil {
		setupLog.Error(err, "unable to start manager")
		os.Exit(1)
	}

285
286
287
288
	// Initialize observability metrics
	setupLog.Info("Initializing observability metrics")
	observability.InitMetrics()

289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
	// Set up webhook certificate management.
	// A direct (non-cached) client is needed because the manager's cache isn't started yet.
	directClient, err := client.New(mgr.GetConfig(), client.Options{Scheme: crdScheme})
	if err != nil {
		setupLog.Error(err, "unable to create direct client for cert management")
		os.Exit(1)
	}
	certMgr, err := internalcert.NewCertManager(directClient, &operatorCfg.Server.Webhook)
	if err != nil {
		setupLog.Error(err, "unable to create cert manager")
		os.Exit(1)
	}
	if err = certMgr.Setup(mainCtx, mgr); err != nil {
		setupLog.Error(err, "failed to setup webhook certificate management")
		os.Exit(1)
	}

306
307
308
309
310
311
312
313
	// Initialize namespace scope mechanism
	var leaseManager *namespace_scope.LeaseManager
	var leaseWatcher *namespace_scope.LeaseWatcher

	if restrictedNamespace != "" {
		// Namespace-restricted mode: Create and maintain namespace scope marker lease
		setupLog.Info("Creating namespace scope marker lease manager",
			"namespace", restrictedNamespace,
314
315
			"leaseDuration", operatorCfg.Namespace.Scope.LeaseDuration.Duration,
			"renewInterval", operatorCfg.Namespace.Scope.LeaseRenewInterval.Duration)
316
317
318
319
320

		leaseManager, err = namespace_scope.NewLeaseManager(
			mgr.GetConfig(),
			restrictedNamespace,
			operatorVersion,
321
322
			operatorCfg.Namespace.Scope.LeaseDuration.Duration,
			operatorCfg.Namespace.Scope.LeaseRenewInterval.Duration,
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
		)
		if err != nil {
			setupLog.Error(err, "unable to create namespace scope marker lease manager")
			os.Exit(1)
		}

		// Start the lease manager
		if err = leaseManager.Start(mainCtx); err != nil {
			setupLog.Error(err, "unable to start namespace scope marker lease manager")
			os.Exit(1)
		}

		// Monitor for fatal lease errors
		// If lease renewal fails repeatedly, we must exit to prevent split-brain
		go func() {
			select {
			case err := <-leaseManager.Errors():
				setupLog.Error(err, "FATAL: Lease manager encountered unrecoverable error, shutting down to prevent split-brain")
				os.Exit(1)
			case <-mainCtx.Done():
				// Normal shutdown, error channel monitoring no longer needed
				return
			}
		}()

		// Ensure lease is released on shutdown
		defer func() {
			shutdownCtx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
			defer cancel()
			if err := leaseManager.Stop(shutdownCtx); err != nil {
				setupLog.Error(err, "failed to stop lease manager cleanly")
			}
		}()

		setupLog.Info("Namespace scope marker lease manager started successfully")
	} else {
		// Cluster-wide mode: Watch for namespace scope marker leases
		setupLog.Info("Setting up namespace scope marker lease watcher for cluster-wide mode")

		leaseWatcher, err = namespace_scope.NewLeaseWatcher(mgr.GetConfig())
		if err != nil {
			setupLog.Error(err, "unable to create namespace scope marker lease watcher")
			os.Exit(1)
		}

		// Start the lease watcher
		if err = leaseWatcher.Start(mainCtx); err != nil {
			setupLog.Error(err, "unable to start namespace scope marker lease watcher")
			os.Exit(1)
		}

		setupLog.Info("Namespace scope marker lease watcher started successfully")
375

376
377
		// Pass leaseWatcher to runtime config for namespace exclusion filtering
		runtimeConfig.ExcludedNamespaces = leaseWatcher
378
379
	}

380
381
	// Start resource counter background goroutine (after ExcludedNamespaces is set)
	setupLog.Info("Starting resource counter")
382
	go observability.StartResourceCounter(mainCtx, mgr.GetClient(), runtimeConfig.ExcludedNamespaces)
383

384
385
386
387
388
	// Detect orchestrators availability using discovery client.
	// Config overrides (*bool) take precedence over auto-detection:
	//   nil   = auto-detect (backward compatible default)
	//   false = forcibly disabled regardless of API availability
	//   true  = forcibly enabled; hard exit if API is not available (misconfiguration)
389
	setupLog.Info("Detecting Grove availability...")
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
	groveDetected := commonController.DetectGroveAvailability(mainCtx, mgr)
	switch {
	case operatorCfg.Orchestrators.Grove.Enabled == nil:
		runtimeConfig.GroveEnabled = groveDetected
	case *operatorCfg.Orchestrators.Grove.Enabled:
		if !groveDetected {
			setupLog.Error(nil, "Grove is explicitly enabled in config but the Grove API group was not detected in the cluster")
			os.Exit(1)
		}
		runtimeConfig.GroveEnabled = true
	default:
		setupLog.Info("Grove is explicitly disabled via config override")
		runtimeConfig.GroveEnabled = false
	}

405
	setupLog.Info("Detecting LWS availability...")
406
	lwsDetected := commonController.DetectLWSAvailability(mainCtx, mgr)
407
	setupLog.Info("Detecting Volcano availability...")
408
	volcanoDetected := commonController.DetectVolcanoAvailability(mainCtx, mgr)
409
	// LWS for multinode deployment usage depends on both LWS and Volcano availability
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
	switch {
	case operatorCfg.Orchestrators.LWS.Enabled == nil:
		runtimeConfig.LWSEnabled = lwsDetected && volcanoDetected
	case *operatorCfg.Orchestrators.LWS.Enabled:
		if !lwsDetected {
			setupLog.Error(nil, "LWS is explicitly enabled in config but the LWS API group was not detected in the cluster")
			os.Exit(1)
		}
		if !volcanoDetected {
			setupLog.Error(nil, "LWS is explicitly enabled in config but the Volcano API group was not detected in the cluster")
			os.Exit(1)
		}
		runtimeConfig.LWSEnabled = true
	default:
		setupLog.Info("LWS is explicitly disabled via config override")
		runtimeConfig.LWSEnabled = false
	}

428
429
	// Detect Kai-scheduler availability using discovery client
	setupLog.Info("Detecting Kai-scheduler availability...")
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
	kaiSchedulerDetected := commonController.DetectKaiSchedulerAvailability(mainCtx, mgr)
	switch {
	case operatorCfg.Orchestrators.KaiScheduler.Enabled == nil:
		runtimeConfig.KaiSchedulerEnabled = kaiSchedulerDetected
	case *operatorCfg.Orchestrators.KaiScheduler.Enabled:
		if !kaiSchedulerDetected {
			setupLog.Error(nil,
				"Kai-scheduler is explicitly enabled in config but the scheduling.run.ai API group was not detected in the cluster",
			)
			os.Exit(1)
		}
		runtimeConfig.KaiSchedulerEnabled = true
	default:
		setupLog.Info("Kai-scheduler is explicitly disabled via config override")
		runtimeConfig.KaiSchedulerEnabled = false
	}
446

447
	setupLog.Info("Detecting DRA (Dynamic Resource Allocation) availability...")
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
	draDetected := commonController.DetectDRAAvailability(mainCtx, mgr)
	switch {
	case operatorCfg.DRA.Enabled == nil:
		runtimeConfig.DRAEnabled = draDetected
	case *operatorCfg.DRA.Enabled:
		if !draDetected {
			setupLog.Error(nil,
				"DRA is explicitly enabled in config but the resource.k8s.io API group"+
					" was not detected in the cluster (requires Kubernetes 1.32+)",
			)
			os.Exit(1)
		}
		runtimeConfig.DRAEnabled = true
	default:
		setupLog.Info("DRA is explicitly disabled via config override")
		runtimeConfig.DRAEnabled = false
	}
465

466
467
468
	setupLog.Info("Detecting Istio availability...")
	runtimeConfig.IstioAvailable = commonController.DetectIstioAvailability(mainCtx, mgr)

469
	setupLog.Info("Detected orchestrators availability",
470
471
472
473
		"grove", runtimeConfig.GroveEnabled,
		"lws", runtimeConfig.LWSEnabled,
		"volcano", volcanoDetected,
		"kai-scheduler", runtimeConfig.KaiSchedulerEnabled,
474
		"dra", runtimeConfig.DRAEnabled,
475
		"istio", runtimeConfig.IstioAvailable,
476
477
	)

478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
	dockerSecretRetriever := secrets.NewDockerSecretIndexer(mgr.GetClient())
	// refresh whenever a secret is created/deleted/updated
	// Set up informer
	var factory informers.SharedInformerFactory
	if restrictedNamespace == "" {
		factory = informers.NewSharedInformerFactory(kubernetes.NewForConfigOrDie(mgr.GetConfig()), time.Hour*24)
	} else {
		factory = informers.NewFilteredSharedInformerFactory(
			kubernetes.NewForConfigOrDie(mgr.GetConfig()),
			time.Hour*24,
			restrictedNamespace,
			nil,
		)
	}
	secretInformer := factory.Core().V1().Secrets().Informer()
	// Start the informer factory
	go factory.Start(mainCtx.Done())
	// Wait for the initial sync
	if !k8sCache.WaitForCacheSync(mainCtx.Done(), secretInformer.HasSynced) {
		setupLog.Error(nil, "Failed to sync informer cache")
		os.Exit(1)
	}
	setupLog.Info("Secret informer cache synced and ready")
	_, err = secretInformer.AddEventHandler(k8sCache.ResourceEventHandlerFuncs{
		AddFunc: func(obj interface{}) {
			secret := obj.(*corev1.Secret)
			if secret.Type == corev1.SecretTypeDockerConfigJson {
				setupLog.Info("refreshing docker secrets index after secret creation...")
				err := dockerSecretRetriever.RefreshIndex(context.Background())
				if err != nil {
					setupLog.Error(err, "unable to refresh docker secrets index after secret creation")
				} else {
					setupLog.Info("docker secrets index refreshed after secret creation")
				}
			}
		},
		UpdateFunc: func(old, new interface{}) {
			newSecret := new.(*corev1.Secret)
			if newSecret.Type == corev1.SecretTypeDockerConfigJson {
				setupLog.Info("refreshing docker secrets index after secret update...")
				err := dockerSecretRetriever.RefreshIndex(context.Background())
				if err != nil {
					setupLog.Error(err, "unable to refresh docker secrets index after secret update")
				} else {
					setupLog.Info("docker secrets index refreshed after secret update")
				}
			}
		},
		DeleteFunc: func(obj interface{}) {
			secret := obj.(*corev1.Secret)
			if secret.Type == corev1.SecretTypeDockerConfigJson {
				setupLog.Info("refreshing docker secrets index after secret deletion...")
				err := dockerSecretRetriever.RefreshIndex(context.Background())
				if err != nil {
					setupLog.Error(err, "unable to refresh docker secrets index after secret deletion")
				} else {
					setupLog.Info("docker secrets index refreshed after secret deletion")
				}
			}
		},
	})
	if err != nil {
		setupLog.Error(err, "unable to add event handler to secret informer")
541
542
		os.Exit(1)
	}
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
	// launch a goroutine to refresh the docker secret indexer in any case every minute
	go func() {
		// Initial refresh
		if err := dockerSecretRetriever.RefreshIndex(context.Background()); err != nil {
			setupLog.Error(err, "initial docker secrets index refresh failed")
		}
		ticker := time.NewTicker(60 * time.Second)
		defer ticker.Stop()
		for {
			select {
			case <-mainCtx.Done():
				return
			case <-ticker.C:
				setupLog.Info("refreshing docker secrets index...")
				if err := dockerSecretRetriever.RefreshIndex(mainCtx); err != nil {
					setupLog.Error(err, "unable to refresh docker secrets index")
				}
				setupLog.Info("docker secrets index refreshed")
			}
		}
	}()
564

565
	sshKeyManager := secret.NewSSHKeyManager(mgr.GetClient(), operatorCfg.MPI)
566

567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
	if err := mgr.AddHealthzCheck("healthz", healthz.Ping); err != nil {
		setupLog.Error(err, "unable to set up health check")
		os.Exit(1)
	}
	webhooksReady := make(chan struct{})
	if err := mgr.AddReadyzCheck("readyz", func(req *http.Request) error {
		select {
		case <-webhooksReady:
			return nil
		default:
			return fmt.Errorf("webhook handlers not yet registered")
		}
	}); err != nil {
		setupLog.Error(err, "unable to set up ready check")
		os.Exit(1)
	}

	// Register controllers synchronously before mgr.Start().
	// Controllers don't depend on TLS certificates.
	if err := registerControllers(
		mgr, operatorCfg, runtimeConfig,
588
		dockerSecretRetriever, sshKeyManager,
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
	); err != nil {
		setupLog.Error(err, "failed to register controllers")
		os.Exit(1)
	}

	// Webhooks require TLS certificates to serve HTTPS. Register them in a
	// goroutine that blocks until the cert-controller has written the certs.
	go func() {
		certMgr.WaitReady()

		if operatorCfg.Server.Webhook.CertProvisionMode == configv1alpha1.CertProvisionModeAuto {
			injector, err := internalcert.NewCABundleInjector(mgr.GetClient(), operatorCfg)
			if err != nil {
				setupLog.Error(err, "unable to create CA bundle injector")
				os.Exit(1)
			}
			if err := injector.InjectAll(mainCtx); err != nil {
				setupLog.Error(err, "failed to inject CA bundles into webhook configurations")
				os.Exit(1)
			}
		}

		if err := registerWebhooks(mgr, operatorCfg, runtimeConfig, operatorVersion); err != nil {
			setupLog.Error(err, "failed to register webhooks")
			os.Exit(1)
		}
		close(webhooksReady)
	}()

	setupLog.Info("starting manager")
	if err := mgr.Start(mainCtx); err != nil {
		setupLog.Error(err, "problem running manager")
		os.Exit(1)
	}
}

func registerControllers(
	mgr ctrl.Manager,
	operatorCfg *configv1alpha1.OperatorConfiguration,
	runtimeConfig *commonController.RuntimeConfig,
	dockerSecretRetriever *secrets.DockerSecretIndexer,
630
	sshKeyManager *secret.SSHKeyManager,
631
632
) error {
	if err := (&controller.DynamoComponentDeploymentReconciler{
633
634
		Client:                mgr.GetClient(),
		Recorder:              mgr.GetEventRecorderFor("dynamocomponentdeployment"),
635
636
		Config:                operatorCfg,
		RuntimeConfig:         runtimeConfig,
637
		DockerSecretRetriever: dockerSecretRetriever,
638
	}).SetupWithManager(mgr); err != nil {
639
		return fmt.Errorf("unable to create DynamoComponentDeployment controller: %w", err)
640
	}
641

642
643
	scaleClient, err := createScalesGetter(mgr)
	if err != nil {
644
		return fmt.Errorf("unable to create scale client: %w", err)
645
646
	}

647
648
	rbacManager := rbac.NewManager(mgr.GetClient())

649
	if err = (&controller.DynamoGraphDeploymentReconciler{
650
651
		Client:                mgr.GetClient(),
		Recorder:              mgr.GetEventRecorderFor("dynamographdeployment"),
652
653
		Config:                operatorCfg,
		RuntimeConfig:         runtimeConfig,
654
		DockerSecretRetriever: dockerSecretRetriever,
655
		ScaleClient:           scaleClient,
656
		SSHKeyManager:         sshKeyManager,
657
		RBACManager:           rbacManager,
658
	}).SetupWithManager(mgr); err != nil {
659
		return fmt.Errorf("unable to create DynamoGraphDeployment controller: %w", err)
660
	}
661

662
	if err = (&controller.DynamoGraphDeploymentScalingAdapterReconciler{
663
664
665
666
667
		Client:        mgr.GetClient(),
		Scheme:        mgr.GetScheme(),
		Recorder:      mgr.GetEventRecorderFor("dgdscalingadapter"),
		Config:        operatorCfg,
		RuntimeConfig: runtimeConfig,
668
	}).SetupWithManager(mgr); err != nil {
669
		return fmt.Errorf("unable to create DGDScalingAdapter controller: %w", err)
670
671
	}

672
	if err = (&controller.DynamoGraphDeploymentRequestReconciler{
673
674
675
676
677
678
679
680
		Client:            mgr.GetClient(),
		APIReader:         mgr.GetAPIReader(),
		Recorder:          mgr.GetEventRecorderFor("dynamographdeploymentrequest"),
		Config:            operatorCfg,
		RuntimeConfig:     runtimeConfig,
		GPUDiscoveryCache: gpu.NewGPUDiscoveryCache(),
		GPUDiscovery:      gpu.NewGPUDiscovery(gpu.ScrapeMetricsEndpoint),
		RBACManager:       rbacManager,
681
	}).SetupWithManager(mgr); err != nil {
682
		return fmt.Errorf("unable to create DynamoGraphDeploymentRequest controller: %w", err)
683
	}
684
685
686
687
688

	if err = (&controller.DynamoModelReconciler{
		Client:         mgr.GetClient(),
		Recorder:       mgr.GetEventRecorderFor("dynamomodel"),
		EndpointClient: modelendpoint.NewClient(),
689
690
		Config:         operatorCfg,
		RuntimeConfig:  runtimeConfig,
691
	}).SetupWithManager(mgr); err != nil {
692
		return fmt.Errorf("unable to create DynamoModel controller: %w", err)
693
	}
694

695
	if err = (&controller.CheckpointReconciler{
696
697
698
699
		Client:        mgr.GetClient(),
		Config:        operatorCfg,
		RuntimeConfig: runtimeConfig,
		Recorder:      mgr.GetEventRecorderFor("checkpoint"),
700
	}).SetupWithManager(mgr); err != nil {
701
		return fmt.Errorf("unable to create DynamoCheckpoint controller: %w", err)
702
703
	}

704
705
706
707
708
709
710
711
712
	if runtimeConfig.GroveEnabled {
		if err = controller.NewFailoverCascadeReconciler(
			mgr.GetClient(),
			mgr.GetEventRecorderFor("gms-failover-cascade"),
		).SetupWithManager(mgr); err != nil {
			return fmt.Errorf("unable to create GMS FailoverCascade controller: %w", err)
		}
	}

713
714
715
716
717
718
719
720
721
722
	setupLog.Info("Controllers registered successfully")
	return nil
}

func registerWebhooks(
	mgr ctrl.Manager,
	operatorCfg *configv1alpha1.OperatorConfiguration,
	runtimeConfig *commonController.RuntimeConfig,
	operatorVersion string,
) error {
723
	isClusterWide := operatorCfg.Namespace.Restricted == ""
724
725
	if isClusterWide {
		setupLog.Info("Configuring webhooks with lease-based namespace exclusion for cluster-wide mode")
726
		internalwebhook.SetExcludedNamespaces(runtimeConfig.ExcludedNamespaces)
727
728
	} else {
		setupLog.Info("Configuring webhooks for namespace-restricted mode (no lease checking)",
729
			"restrictedNamespace", operatorCfg.Namespace.Restricted)
730
731
		internalwebhook.SetExcludedNamespaces(nil)
	}
732

733
734
735
736
737
738
739
740
	var operatorPrincipal string
	if sa, ns := os.Getenv("POD_SERVICE_ACCOUNT"), os.Getenv("POD_NAMESPACE"); sa != "" && ns != "" {
		operatorPrincipal = fmt.Sprintf("system:serviceaccount:%s:%s", ns, sa)
		setupLog.Info("Detected operator principal from downward API", "principal", operatorPrincipal)
	} else {
		setupLog.Info("POD_SERVICE_ACCOUNT/POD_NAMESPACE not set; operator SA self-identification disabled")
	}

741
	setupLog.Info("Registering validation webhooks")
742

743
	dcdHandler := webhookvalidation.NewDynamoComponentDeploymentHandler()
744
745
	if err := dcdHandler.RegisterWithManager(mgr); err != nil {
		return fmt.Errorf("unable to register DynamoComponentDeployment webhook: %w", err)
746
	}
747

748
	dgdHandler := webhookvalidation.NewDynamoGraphDeploymentHandler(mgr, operatorPrincipal, runtimeConfig.GroveEnabled)
749
750
	if err := dgdHandler.RegisterWithManager(mgr); err != nil {
		return fmt.Errorf("unable to register DynamoGraphDeployment webhook: %w", err)
751
	}
752

753
	dmHandler := webhookvalidation.NewDynamoModelHandler()
754
755
	if err := dmHandler.RegisterWithManager(mgr); err != nil {
		return fmt.Errorf("unable to register DynamoModel webhook: %w", err)
756
	}
757

758
759
760
	dgdrHandler := webhookvalidation.NewDynamoGraphDeploymentRequestHandler(
		isClusterWide, ptr.Deref(operatorCfg.GPU.DiscoveryEnabled, true),
	)
761
762
	if err := dgdrHandler.RegisterWithManager(mgr); err != nil {
		return fmt.Errorf("unable to register DynamoGraphDeploymentRequest webhook: %w", err)
763
	}
764

765
	if err := ctrl.NewWebhookManagedBy(mgr).
766
		For(&nvidiacomv1beta1.DynamoGraphDeploymentRequest{}).
767
		Complete(); err != nil {
768
		return fmt.Errorf("unable to register DynamoGraphDeploymentRequest conversion webhook: %w", err)
769
	}
770

771
	setupLog.Info("Registering defaulting webhooks")
772

773
	dgdDefaulter := webhookdefaulting.NewDGDDefaulter(operatorVersion)
774
775
	if err := dgdDefaulter.RegisterWithManager(mgr); err != nil {
		return fmt.Errorf("unable to register DynamoGraphDeployment defaulting webhook: %w", err)
776
	}
777

778
	dgdrDefaulter := webhookdefaulting.NewDGDRDefaulter(operatorVersion)
779
780
	if err := dgdrDefaulter.RegisterWithManager(mgr); err != nil {
		return fmt.Errorf("unable to register DynamoGraphDeploymentRequest defaulting webhook: %w", err)
781
782
	}

783
784
	setupLog.Info("Webhooks registered successfully")
	return nil
785
}