main.go 26.8 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
	"time"
30
31
32

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

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

49
50
	k8sruntime "k8s.io/apimachinery/pkg/runtime"
	"k8s.io/apimachinery/pkg/runtime/serializer"
51
52
53
54
55
	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"
56
	metricsfilters "sigs.k8s.io/controller-runtime/pkg/metrics/filters"
57
58
59
	metricsserver "sigs.k8s.io/controller-runtime/pkg/metrics/server"
	"sigs.k8s.io/controller-runtime/pkg/webhook"

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

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

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

93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
// 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
}

115
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
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
}

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

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

203
204
205
206
	// 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())
207
208
		os.Exit(1)
	}
209
	setupLog.Info("Operator version configured", "version", operatorVersion)
210

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

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

	// 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){}
228
	if !operatorCfg.Security.EnableHTTP2 {
229
230
231
232
		tlsOpts = append(tlsOpts, disableHTTP2)
	}

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

239
240
241
242
243
	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,
	)

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

	restrictedNamespace := operatorCfg.Namespace.Restricted
260
261
262
263
	if restrictedNamespace != "" {
		mgrOpts.Cache.DefaultNamespaces = map[string]cache.Config{
			restrictedNamespace: {},
		}
264
265
266
		setupLog.Info("Restricted namespace configured, launching in restricted mode", "namespace", restrictedNamespace)
	} else {
		setupLog.Info("No restricted namespace configured, launching in cluster-wide mode")
267
268
269
270
271
272
273
	}
	mgr, err := ctrl.NewManager(ctrl.GetConfigOrDie(), mgrOpts)
	if err != nil {
		setupLog.Error(err, "unable to start manager")
		os.Exit(1)
	}

274
275
276
277
	// Initialize observability metrics
	setupLog.Info("Initializing observability metrics")
	observability.InitMetrics()

278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
	// 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)
	}

295
296
297
298
299
300
301
302
	// 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,
303
304
			"leaseDuration", operatorCfg.Namespace.Scope.LeaseDuration.Duration,
			"renewInterval", operatorCfg.Namespace.Scope.LeaseRenewInterval.Duration)
305
306
307
308
309

		leaseManager, err = namespace_scope.NewLeaseManager(
			mgr.GetConfig(),
			restrictedNamespace,
			operatorVersion,
310
311
			operatorCfg.Namespace.Scope.LeaseDuration.Duration,
			operatorCfg.Namespace.Scope.LeaseRenewInterval.Duration,
312
313
314
315
316
317
318
319
320
321
322
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
		)
		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")
364

365
366
		// Pass leaseWatcher to runtime config for namespace exclusion filtering
		runtimeConfig.ExcludedNamespaces = leaseWatcher
367
368
	}

369
370
	// Start resource counter background goroutine (after ExcludedNamespaces is set)
	setupLog.Info("Starting resource counter")
371
	go observability.StartResourceCounter(mainCtx, mgr.GetClient(), runtimeConfig.ExcludedNamespaces)
372

373
374
375
376
377
	// 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)
378
	setupLog.Info("Detecting Grove availability...")
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
	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
	}

394
	setupLog.Info("Detecting LWS availability...")
395
	lwsDetected := commonController.DetectLWSAvailability(mainCtx, mgr)
396
	setupLog.Info("Detecting Volcano availability...")
397
	volcanoDetected := commonController.DetectVolcanoAvailability(mainCtx, mgr)
398
	// LWS for multinode deployment usage depends on both LWS and Volcano availability
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
	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
	}

417
418
	// Detect Kai-scheduler availability using discovery client
	setupLog.Info("Detecting Kai-scheduler availability...")
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
	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
	}
435

436
	setupLog.Info("Detected orchestrators availability",
437
438
439
440
		"grove", runtimeConfig.GroveEnabled,
		"lws", runtimeConfig.LWSEnabled,
		"volcano", volcanoDetected,
		"kai-scheduler", runtimeConfig.KaiSchedulerEnabled,
441
442
	)

443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
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
	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")
506
507
		os.Exit(1)
	}
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
	// 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")
			}
		}
	}()
529

530
	sshKeyManager := secret.NewSSHKeyManager(mgr.GetClient(), operatorCfg.MPI)
531

532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
	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,
553
		dockerSecretRetriever, sshKeyManager,
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
	); 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,
595
	sshKeyManager *secret.SSHKeyManager,
596
597
) error {
	if err := (&controller.DynamoComponentDeploymentReconciler{
598
599
		Client:                mgr.GetClient(),
		Recorder:              mgr.GetEventRecorderFor("dynamocomponentdeployment"),
600
601
		Config:                operatorCfg,
		RuntimeConfig:         runtimeConfig,
602
		DockerSecretRetriever: dockerSecretRetriever,
603
	}).SetupWithManager(mgr); err != nil {
604
		return fmt.Errorf("unable to create DynamoComponentDeployment controller: %w", err)
605
	}
606

607
608
	scaleClient, err := createScalesGetter(mgr)
	if err != nil {
609
		return fmt.Errorf("unable to create scale client: %w", err)
610
611
	}

612
613
	rbacManager := rbac.NewManager(mgr.GetClient())

614
	if err = (&controller.DynamoGraphDeploymentReconciler{
615
616
		Client:                mgr.GetClient(),
		Recorder:              mgr.GetEventRecorderFor("dynamographdeployment"),
617
618
		Config:                operatorCfg,
		RuntimeConfig:         runtimeConfig,
619
		DockerSecretRetriever: dockerSecretRetriever,
620
		ScaleClient:           scaleClient,
621
		SSHKeyManager:         sshKeyManager,
622
		RBACManager:           rbacManager,
623
	}).SetupWithManager(mgr); err != nil {
624
		return fmt.Errorf("unable to create DynamoGraphDeployment controller: %w", err)
625
	}
626

627
	if err = (&controller.DynamoGraphDeploymentScalingAdapterReconciler{
628
629
630
631
632
		Client:        mgr.GetClient(),
		Scheme:        mgr.GetScheme(),
		Recorder:      mgr.GetEventRecorderFor("dgdscalingadapter"),
		Config:        operatorCfg,
		RuntimeConfig: runtimeConfig,
633
	}).SetupWithManager(mgr); err != nil {
634
		return fmt.Errorf("unable to create DGDScalingAdapter controller: %w", err)
635
636
	}

637
	if err = (&controller.DynamoGraphDeploymentRequestReconciler{
638
639
640
641
642
643
644
645
		Client:            mgr.GetClient(),
		APIReader:         mgr.GetAPIReader(),
		Recorder:          mgr.GetEventRecorderFor("dynamographdeploymentrequest"),
		Config:            operatorCfg,
		RuntimeConfig:     runtimeConfig,
		GPUDiscoveryCache: gpu.NewGPUDiscoveryCache(),
		GPUDiscovery:      gpu.NewGPUDiscovery(gpu.ScrapeMetricsEndpoint),
		RBACManager:       rbacManager,
646
	}).SetupWithManager(mgr); err != nil {
647
		return fmt.Errorf("unable to create DynamoGraphDeploymentRequest controller: %w", err)
648
	}
649
650
651
652
653

	if err = (&controller.DynamoModelReconciler{
		Client:         mgr.GetClient(),
		Recorder:       mgr.GetEventRecorderFor("dynamomodel"),
		EndpointClient: modelendpoint.NewClient(),
654
655
		Config:         operatorCfg,
		RuntimeConfig:  runtimeConfig,
656
	}).SetupWithManager(mgr); err != nil {
657
		return fmt.Errorf("unable to create DynamoModel controller: %w", err)
658
	}
659

660
	if err = (&controller.CheckpointReconciler{
661
662
663
664
		Client:        mgr.GetClient(),
		Config:        operatorCfg,
		RuntimeConfig: runtimeConfig,
		Recorder:      mgr.GetEventRecorderFor("checkpoint"),
665
	}).SetupWithManager(mgr); err != nil {
666
		return fmt.Errorf("unable to create DynamoCheckpoint controller: %w", err)
667
668
	}

669
670
671
672
673
674
675
676
677
678
	setupLog.Info("Controllers registered successfully")
	return nil
}

func registerWebhooks(
	mgr ctrl.Manager,
	operatorCfg *configv1alpha1.OperatorConfiguration,
	runtimeConfig *commonController.RuntimeConfig,
	operatorVersion string,
) error {
679
	isClusterWide := operatorCfg.Namespace.Restricted == ""
680
681
	if isClusterWide {
		setupLog.Info("Configuring webhooks with lease-based namespace exclusion for cluster-wide mode")
682
		internalwebhook.SetExcludedNamespaces(runtimeConfig.ExcludedNamespaces)
683
684
	} else {
		setupLog.Info("Configuring webhooks for namespace-restricted mode (no lease checking)",
685
			"restrictedNamespace", operatorCfg.Namespace.Restricted)
686
687
		internalwebhook.SetExcludedNamespaces(nil)
	}
688

689
690
691
692
693
694
695
696
	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")
	}

697
	setupLog.Info("Registering validation webhooks")
698

699
	dcdHandler := webhookvalidation.NewDynamoComponentDeploymentHandler()
700
701
	if err := dcdHandler.RegisterWithManager(mgr); err != nil {
		return fmt.Errorf("unable to register DynamoComponentDeployment webhook: %w", err)
702
	}
703

704
	dgdHandler := webhookvalidation.NewDynamoGraphDeploymentHandler(mgr, operatorPrincipal)
705
706
	if err := dgdHandler.RegisterWithManager(mgr); err != nil {
		return fmt.Errorf("unable to register DynamoGraphDeployment webhook: %w", err)
707
	}
708

709
	dmHandler := webhookvalidation.NewDynamoModelHandler()
710
711
	if err := dmHandler.RegisterWithManager(mgr); err != nil {
		return fmt.Errorf("unable to register DynamoModel webhook: %w", err)
712
	}
713

714
715
716
	dgdrHandler := webhookvalidation.NewDynamoGraphDeploymentRequestHandler(
		isClusterWide, ptr.Deref(operatorCfg.GPU.DiscoveryEnabled, true),
	)
717
718
	if err := dgdrHandler.RegisterWithManager(mgr); err != nil {
		return fmt.Errorf("unable to register DynamoGraphDeploymentRequest webhook: %w", err)
719
	}
720

721
	if err := ctrl.NewWebhookManagedBy(mgr).
722
		For(&nvidiacomv1beta1.DynamoGraphDeploymentRequest{}).
723
		Complete(); err != nil {
724
		return fmt.Errorf("unable to register DynamoGraphDeploymentRequest conversion webhook: %w", err)
725
	}
726

727
	setupLog.Info("Registering defaulting webhooks")
728

729
	dgdDefaulter := webhookdefaulting.NewDGDDefaulter(operatorVersion)
730
731
	if err := dgdDefaulter.RegisterWithManager(mgr); err != nil {
		return fmt.Errorf("unable to register DynamoGraphDeployment defaulting webhook: %w", err)
732
	}
733

734
	dgdrDefaulter := webhookdefaulting.NewDGDRDefaulter(operatorVersion)
735
736
	if err := dgdrDefaulter.RegisterWithManager(mgr); err != nil {
		return fmt.Errorf("unable to register DynamoGraphDeploymentRequest defaulting webhook: %w", err)
737
738
	}

739
740
	setupLog.Info("Webhooks registered successfully")
	return nil
741
}