From efec27f9cf3fa7c30287c653a0f01cfae0000efd Mon Sep 17 00:00:00 2001 From: Thomas Montfort <61255722+tmonty12@users.noreply.github.com> Date: Mon, 10 Nov 2025 16:03:49 -0500 Subject: [PATCH] feat: etcd-less operator updates (#4214) --- .../operator/templates/_validation.tpl | 10 + .../operator/templates/deployment.yaml | 4 + deploy/cloud/helm/platform/values.yaml | 2 + deploy/cloud/operator/cmd/main.go | 15 + .../cloud/operator/internal/consts/consts.go | 7 + .../dynamocomponentdeployment_controller.go | 55 ++- ...namocomponentdeployment_controller_test.go | 26 +- .../dynamographdeployment_controller.go | 62 ++- .../internal/controller_common/predicate.go | 17 + .../operator/internal/discovery/resource.go | 96 +++++ .../internal/dynamo/component_common.go | 33 +- .../internal/dynamo/component_planner_test.go | 48 ++- .../cloud/operator/internal/dynamo/graph.go | 71 +++- .../operator/internal/dynamo/graph_test.go | 394 +++++++++++++++++- 14 files changed, 758 insertions(+), 82 deletions(-) create mode 100644 deploy/cloud/operator/internal/discovery/resource.go diff --git a/deploy/cloud/helm/platform/components/operator/templates/_validation.tpl b/deploy/cloud/helm/platform/components/operator/templates/_validation.tpl index 99de02d9d..acb764566 100644 --- a/deploy/cloud/helm/platform/components/operator/templates/_validation.tpl +++ b/deploy/cloud/helm/platform/components/operator/templates/_validation.tpl @@ -102,3 +102,13 @@ Validation for configuration consistency {{- end -}} {{- end -}} {{- end -}} + +{{/* +Validation for discoverBackend configuration +*/}} +{{- define "dynamo-operator.validateDiscoveryBackend" -}} +{{- $discoveryBackend := .Values.discoveryBackend -}} +{{- if and (ne $discoveryBackend "") (ne $discoveryBackend "kubernetes") -}} + {{- fail (printf "VALIDATION ERROR: discoveryBackend must be either an empty string (defaults to ETCD) or 'kubernetes'. Got: '%s'" $discoveryBackend) -}} +{{- end -}} +{{- end -}} diff --git a/deploy/cloud/helm/platform/components/operator/templates/deployment.yaml b/deploy/cloud/helm/platform/components/operator/templates/deployment.yaml index 41343bab8..f20f15e6a 100644 --- a/deploy/cloud/helm/platform/components/operator/templates/deployment.yaml +++ b/deploy/cloud/helm/platform/components/operator/templates/deployment.yaml @@ -16,6 +16,7 @@ {{/* Validate installation to prevent conflicts */}} {{- include "dynamo-operator.validateClusterWideInstallation" . -}} {{- include "dynamo-operator.validateConfiguration" . -}} +{{- include "dynamo-operator.validateDiscoveryBackend" . -}} --- apiVersion: apps/v1 @@ -131,6 +132,9 @@ spec: - --dgdr-profiling-cluster-role-name={{ include "dynamo-operator.fullname" . }}-dgdr-profiling - --planner-cluster-role-name={{ include "dynamo-operator.fullname" . }}-planner {{- end }} + {{- if .Values.discoveryBackend }} + - --discovery-backend={{ .Values.discoveryBackend }} + {{- end }} {{- if .Values.namespaceRestriction.enabled }} {{- if .Values.namespaceRestriction.lease }} - --namespace-scope-lease-duration={{ .Values.namespaceRestriction.lease.duration }} diff --git a/deploy/cloud/helm/platform/values.yaml b/deploy/cloud/helm/platform/values.yaml index e5519676e..f60a29ef2 100644 --- a/deploy/cloud/helm/platform/values.yaml +++ b/deploy/cloud/helm/platform/values.yaml @@ -42,6 +42,8 @@ dynamo-operator: # Interval for renewing the namespace scope marker lease (namespace-restricted mode only). The namespace-restricted operator renews its lease at this interval to signal it's still running. renewInterval: 10s + # -- The Dynamo discovery backend to use. By default, will rely on ETCD for discovery. Can be set to "kubernetes" to use Kubernetes API for service discovery. -- + discoveryBackend: "" # Controller manager configuration controllerManager: diff --git a/deploy/cloud/operator/cmd/main.go b/deploy/cloud/operator/cmd/main.go index 7047656ea..19680e53f 100644 --- a/deploy/cloud/operator/cmd/main.go +++ b/deploy/cloud/operator/cmd/main.go @@ -146,6 +146,7 @@ func main() { var namespaceScopeLeaseDuration time.Duration var namespaceScopeLeaseRenewInterval time.Duration var operatorVersion string + var discoveryBackend string flag.StringVar(&metricsAddr, "metrics-bind-address", ":8080", "The address the metric endpoint binds to.") flag.StringVar(&probeAddr, "health-probe-bind-address", ":8081", "The address the probe endpoint binds to.") flag.BoolVar(&enableLeaderElection, "leader-elect", false, @@ -194,6 +195,8 @@ func main() { "Interval for renewing namespace scope marker lease (namespace-restricted mode only)") flag.StringVar(&operatorVersion, "operator-version", "unknown", "Version of the operator (used in lease holder identity)") + flag.StringVar(&discoveryBackend, "discovery-backend", "", + "Discovery backend to use: empty string (default, uses ETCD) or 'kubernetes' (uses Kubernetes API)") opts := zap.Options{ Development: true, } @@ -205,6 +208,17 @@ func main() { os.Exit(1) } + // Validate discoverBackend value + if discoveryBackend != "" && discoveryBackend != "kubernetes" { + setupLog.Error(nil, "invalid discover-backend value, must be empty string or 'kubernetes'", "value", discoveryBackend) + os.Exit(1) + } + if discoveryBackend != "" { + setupLog.Info("Discovery backend configured", "backend", discoveryBackend) + } else { + setupLog.Info("Discovery backend configured", "backend", "etcd (default)") + } + // Validate modelExpressURL if provided if modelExpressURL != "" { if _, err := url.Parse(modelExpressURL); err != nil { @@ -253,6 +267,7 @@ func main() { PlannerClusterRoleName: plannerClusterRoleName, DGDRProfilingClusterRoleName: dgdrProfilingClusterRoleName, }, + DiscoveryBackend: discoveryBackend, } mainCtx := ctrl.SetupSignalHandler() diff --git a/deploy/cloud/operator/internal/consts/consts.go b/deploy/cloud/operator/internal/consts/consts.go index e8de73125..da16238ab 100644 --- a/deploy/cloud/operator/internal/consts/consts.go +++ b/deploy/cloud/operator/internal/consts/consts.go @@ -31,6 +31,7 @@ const ( KubeAnnotationEnableGrove = "nvidia.com/enable-grove" KubeAnnotationDisableImagePullSecretDiscovery = "nvidia.com/disable-image-pull-secret-discovery" + KubeAnnotationDynamoDiscoveryBackend = "nvidia.com/dynamo-discovery-backend" KubeLabelDynamoGraphDeploymentName = "nvidia.com/dynamo-graph-deployment-name" KubeLabelDynamoComponent = "nvidia.com/dynamo-component" @@ -41,6 +42,7 @@ const ( KubeLabelDynamoBaseModel = "nvidia.com/dynamo-base-model" KubeLabelDynamoBaseModelHash = "nvidia.com/dynamo-base-model-hash" KubeAnnotationDynamoBaseModel = "nvidia.com/dynamo-base-model" + KubeLabelDynamoDiscoveryBackend = "nvidia.com/dynamo-discovery-backend" KubeLabelValueFalse = "false" KubeLabelValueTrue = "true" @@ -50,12 +52,17 @@ const ( KubeResourceGPUNvidia = "nvidia.com/gpu" DynamoDeploymentConfigEnvVar = "DYN_DEPLOYMENT_CONFIG" + DynamoNamespaceEnvVar = "DYN_NAMESPACE" + DynamoComponentEnvVar = "DYN_COMPONENT" + DynamoDiscoveryBackendEnvVar = "DYN_DISCOVERY_BACKEND" GlobalDynamoNamespace = "dynamo" ComponentTypePlanner = "planner" ComponentTypeFrontend = "frontend" ComponentTypeWorker = "worker" + ComponentTypePrefill = "prefill" + ComponentTypeDecode = "decode" ComponentTypeDefault = "default" PlannerServiceAccountName = "planner-serviceaccount" diff --git a/deploy/cloud/operator/internal/controller/dynamocomponentdeployment_controller.go b/deploy/cloud/operator/internal/controller/dynamocomponentdeployment_controller.go index 4cac62834..0b9d50a4a 100644 --- a/deploy/cloud/operator/internal/controller/dynamocomponentdeployment_controller.go +++ b/deploy/cloud/operator/internal/controller/dynamocomponentdeployment_controller.go @@ -63,14 +63,7 @@ const ( KubeAnnotationEnableStealingTrafficDebugMode = "nvidia.com/enable-stealing-traffic-debug-mode" KubeAnnotationEnableDebugMode = "nvidia.com/enable-debug-mode" KubeAnnotationEnableDebugPodReceiveProductionTraffic = "nvidia.com/enable-debug-pod-receive-production-traffic" - DeploymentTargetTypeProduction = "production" DeploymentTargetTypeDebug = "debug" - HeaderNameDebug = "X-Nvidia-Debug" - KubernetesDeploymentStrategy = "kubernetes" - - DeploymentTypeStandard = "standard" - DeploymentTypeMultinodeGrove = "multinode-grove" - ComponentTypePlanner = "Planner" ) // DynamoComponentDeploymentReconciler reconciles a DynamoComponentDeployment object @@ -1276,40 +1269,56 @@ func (r *DynamoComponentDeploymentReconciler) generateService(opt generateResour }, } - if !opt.dynamoComponentDeployment.IsFrontendComponent() || (!opt.isGenericService && !opt.containsStealingTrafficDebugModeEnabled) { + isK8sDiscovery := r.Config.IsK8sDiscoveryEnabled(opt.dynamoComponentDeployment.Spec.Annotations) + + // if discovery backend is k8s we want to create a service for each component + // else, only create for the frontend component + if !opt.isGenericService && !opt.containsStealingTrafficDebugModeEnabled && !(isK8sDiscovery || opt.dynamoComponentDeployment.IsFrontendComponent()) { // if it's not the main component or if it's not a generic service and not contains stealing traffic debug mode enabled, we don't need to create the service return kubeService, true, nil } labels := r.getKubeLabels(opt.dynamoComponentDeployment) - selector := make(map[string]string) - - for k, v := range labels { - selector[k] = v + if opt.dynamoComponentDeployment.Spec.DynamoNamespace == nil { + return nil, false, fmt.Errorf("expected DynamoComponentDeployment %s to have a dynamoNamespace", opt.dynamoComponentDeployment.Name) } - // If using LeaderWorkerSet, modify selector to only target leaders + selector := map[string]string{ + commonconsts.KubeLabelDynamoComponentType: opt.dynamoComponentDeployment.Spec.ComponentType, + commonconsts.KubeLabelDynamoNamespace: *opt.dynamoComponentDeployment.Spec.DynamoNamespace, + } + // // If using LeaderWorkerSet, modify selector to only target leaders if opt.dynamoComponentDeployment.IsMultinode() { selector["role"] = "leader" } - if opt.isStealingTrafficDebugModeEnabled { selector[commonconsts.KubeLabelDynamoDeploymentTargetType] = DeploymentTargetTypeDebug } + if isK8sDiscovery { + labels[commonconsts.KubeLabelDynamoDiscoveryBackend] = "kubernetes" + } - targetPort := intstr.FromString(commonconsts.DynamoContainerPortName) + var servicePort corev1.ServicePort + if opt.dynamoComponentDeployment.IsFrontendComponent() { + servicePort = corev1.ServicePort{ + Name: commonconsts.DynamoServicePortName, + Port: commonconsts.DynamoServicePort, + TargetPort: intstr.FromString(commonconsts.DynamoContainerPortName), + Protocol: corev1.ProtocolTCP, + } + } else { // TODO: only for worker components + servicePort = corev1.ServicePort{ + Name: commonconsts.DynamoSystemPortName, + Port: commonconsts.DynamoSystemPort, + TargetPort: intstr.FromString(commonconsts.DynamoSystemPortName), + Protocol: corev1.ProtocolTCP, + } + } spec := corev1.ServiceSpec{ Selector: selector, - Ports: []corev1.ServicePort{ - { - Name: commonconsts.DynamoServicePortName, - Port: commonconsts.DynamoServicePort, - TargetPort: targetPort, - Protocol: corev1.ProtocolTCP, - }, - }, + Ports: []corev1.ServicePort{servicePort}, } annotations := r.getKubeAnnotations(opt.dynamoComponentDeployment) diff --git a/deploy/cloud/operator/internal/controller/dynamocomponentdeployment_controller_test.go b/deploy/cloud/operator/internal/controller/dynamocomponentdeployment_controller_test.go index 4c5d7e939..2469891ef 100644 --- a/deploy/cloud/operator/internal/controller/dynamocomponentdeployment_controller_test.go +++ b/deploy/cloud/operator/internal/controller/dynamocomponentdeployment_controller_test.go @@ -824,11 +824,22 @@ func TestDynamoComponentDeploymentReconciler_generateLeaderWorkerSet(t *testing. Command: []string{"/bin/sh", "-c"}, Args: []string{"ray start --head --port=6379 && some dynamo command --tensor-parallel-size 4 --pipeline-parallel-size 1"}, Env: []corev1.EnvVar{ - {Name: "DYN_NAMESPACE", Value: "default"}, + {Name: commonconsts.DynamoComponentEnvVar, Value: commonconsts.ComponentTypeWorker}, + {Name: commonconsts.DynamoNamespaceEnvVar, Value: "default"}, {Name: "DYN_PARENT_DGD_K8S_NAME", Value: "test-lws-deploy"}, {Name: "DYN_PARENT_DGD_K8S_NAMESPACE", Value: "default"}, {Name: "DYN_SYSTEM_PORT", Value: "9090"}, {Name: "DYN_SYSTEM_USE_ENDPOINT_HEALTH_STATUS", Value: "[\"generate\"]"}, + {Name: "POD_NAME", ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.name", + }, + }}, + {Name: "POD_NAMESPACE", ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.namespace", + }, + }}, {Name: "TEST_ENV_FROM_DYNAMO_COMPONENT_DEPLOYMENT_SPEC", Value: "test_value_from_dynamo_component_deployment_spec"}, {Name: "TEST_ENV_FROM_EXTRA_POD_SPEC", Value: "test_value_from_extra_pod_spec"}, }, @@ -937,11 +948,22 @@ func TestDynamoComponentDeploymentReconciler_generateLeaderWorkerSet(t *testing. Command: []string{"/bin/sh", "-c"}, Args: []string{"ray start --address=$(LWS_LEADER_ADDRESS):6379 --block"}, Env: []corev1.EnvVar{ - {Name: "DYN_NAMESPACE", Value: "default"}, + {Name: commonconsts.DynamoComponentEnvVar, Value: commonconsts.ComponentTypeWorker}, + {Name: commonconsts.DynamoNamespaceEnvVar, Value: "default"}, {Name: "DYN_PARENT_DGD_K8S_NAME", Value: "test-lws-deploy"}, {Name: "DYN_PARENT_DGD_K8S_NAMESPACE", Value: "default"}, {Name: "DYN_SYSTEM_PORT", Value: "9090"}, {Name: "DYN_SYSTEM_USE_ENDPOINT_HEALTH_STATUS", Value: "[\"generate\"]"}, + {Name: "POD_NAME", ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.name", + }, + }}, + {Name: "POD_NAMESPACE", ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.namespace", + }, + }}, {Name: "TEST_ENV_FROM_DYNAMO_COMPONENT_DEPLOYMENT_SPEC", Value: "test_value_from_dynamo_component_deployment_spec"}, {Name: "TEST_ENV_FROM_EXTRA_POD_SPEC", Value: "test_value_from_extra_pod_spec"}, }, diff --git a/deploy/cloud/operator/internal/controller/dynamographdeployment_controller.go b/deploy/cloud/operator/internal/controller/dynamographdeployment_controller.go index e36663190..7a8af3d9b 100644 --- a/deploy/cloud/operator/internal/controller/dynamographdeployment_controller.go +++ b/deploy/cloud/operator/internal/controller/dynamographdeployment_controller.go @@ -25,6 +25,7 @@ import ( grovev1alpha1 "github.com/NVIDIA/grove/operator/api/core/v1alpha1" "k8s.io/apimachinery/pkg/api/errors" + "github.com/ai-dynamo/dynamo/deploy/cloud/operator/internal/discovery" "github.com/ai-dynamo/dynamo/deploy/cloud/operator/internal/secret" networkingv1beta1 "istio.io/client-go/pkg/apis/networking/v1beta1" @@ -48,6 +49,7 @@ import ( "github.com/ai-dynamo/dynamo/deploy/cloud/operator/internal/consts" commonController "github.com/ai-dynamo/dynamo/deploy/cloud/operator/internal/controller_common" "github.com/ai-dynamo/dynamo/deploy/cloud/operator/internal/dynamo" + rbacv1 "k8s.io/api/rbac/v1" ) type State string @@ -200,6 +202,13 @@ func (r *DynamoGraphDeploymentReconciler) reconcileResources(ctx context.Context return "", "", "", fmt.Errorf("failed to reconcile top-level PVCs: %w", err) } + // Reconcile the SA, Role and RoleBinding if k8s discovery is enabled + err = r.reconcileK8sDiscoveryResources(ctx, dynamoDeployment) + if err != nil { + logger.Error(err, "Failed to reconcile K8s discovery resources") + return "", "", "", fmt.Errorf("failed to reconcile K8s discovery resources: %w", err) + } + // Orchestrator selection via single boolean annotation: nvidia.com/enable-grove // Unset or not "false": Grove if available; else component mode // "false": component mode (multinode -> LWS; single-node -> standard) @@ -310,6 +319,7 @@ func (r *DynamoGraphDeploymentReconciler) reconcileGroveScaling(ctx context.Cont func (r *DynamoGraphDeploymentReconciler) reconcileGroveResources(ctx context.Context, dynamoDeployment *nvidiacomv1alpha1.DynamoGraphDeployment) (State, Reason, Message, error) { logger := log.FromContext(ctx) + // generate the dynamoComponentsDeployments from the config groveGangSet, err := dynamo.GenerateGrovePodCliqueSet(ctx, dynamoDeployment, r.Config, r.DockerSecretRetriever) if err != nil { @@ -355,9 +365,11 @@ func (r *DynamoGraphDeploymentReconciler) reconcileGroveResources(ctx context.Co resources := []Resource{groveGangSetAsResource} for componentName, component := range dynamoDeployment.Spec.Services { - if component.ComponentType == consts.ComponentTypeFrontend { - // generate the main component service - mainComponentService, err := dynamo.GenerateComponentService(ctx, dynamo.GetDynamoComponentName(dynamoDeployment, componentName), dynamoDeployment.Namespace) + + // if k8s discovery is enabled, create a service for each component + // else, only create for the frontend component + if r.Config.IsK8sDiscoveryEnabled(dynamoDeployment.Annotations) || component.ComponentType == consts.ComponentTypeFrontend { + mainComponentService, err := dynamo.GenerateComponentService(ctx, dynamoDeployment, component, componentName) if err != nil { logger.Error(err, "failed to generate the main component service") return "", "", "", fmt.Errorf("failed to generate the main component service: %w", err) @@ -374,6 +386,9 @@ func (r *DynamoGraphDeploymentReconciler) reconcileGroveResources(ctx context.Co return true, "" }) resources = append(resources, mainComponentServiceAsResource) + } + + if component.ComponentType == consts.ComponentTypeFrontend { // generate the main component ingress ingressSpec := dynamo.GenerateDefaultIngressSpec(dynamoDeployment, r.Config.IngressConfig) if component.Ingress != nil { @@ -501,6 +516,47 @@ func (r *DynamoGraphDeploymentReconciler) reconcilePVC(ctx context.Context, dyna return pvc, nil } +func (r *DynamoGraphDeploymentReconciler) reconcileK8sDiscoveryResources(ctx context.Context, dynamoDeployment *nvidiacomv1alpha1.DynamoGraphDeployment) error { + logger := log.FromContext(ctx) + + if !r.Config.IsK8sDiscoveryEnabled(dynamoDeployment.Annotations) { + logger.Info("K8s discovery is not enabled") + return nil + } else { + logger.Info("K8s discovery is enabled") + } + + serviceAccount := discovery.GetK8sDiscoveryServiceAccount(dynamoDeployment.Name, dynamoDeployment.Namespace) + _, _, err := commonController.SyncResource(ctx, r, dynamoDeployment, func(ctx context.Context) (*corev1.ServiceAccount, bool, error) { + return serviceAccount, false, nil + }) + if err != nil { + logger.Error(err, "failed to sync the k8s discovery service account") + return fmt.Errorf("failed to sync the k8s discovery service account: %w", err) + } + + role := discovery.GetK8sDiscoveryRole(dynamoDeployment.Name, dynamoDeployment.Namespace) + _, _, err = commonController.SyncResource(ctx, r, dynamoDeployment, func(ctx context.Context) (*rbacv1.Role, bool, error) { + return role, false, nil + }) + if err != nil { + logger.Error(err, "failed to sync the k8s discovery role") + return fmt.Errorf("failed to sync the k8s discovery role: %w", err) + } + + roleBinding := discovery.GetK8sDiscoveryRoleBinding(dynamoDeployment.Name, dynamoDeployment.Namespace) + _, _, err = commonController.SyncResource(ctx, r, dynamoDeployment, func(ctx context.Context) (*rbacv1.RoleBinding, bool, error) { + return roleBinding, false, nil + }) + if err != nil { + logger.Error(err, "failed to sync the k8s discovery role binding") + return fmt.Errorf("failed to sync the k8s discovery role binding: %w", err) + } + + return nil + +} + // reconcilePVCs reconciles all top-level PVCs defined in the DynamoGraphDeployment spec func (r *DynamoGraphDeploymentReconciler) reconcilePVCs(ctx context.Context, dynamoDeployment *nvidiacomv1alpha1.DynamoGraphDeployment) error { logger := log.FromContext(ctx) diff --git a/deploy/cloud/operator/internal/controller_common/predicate.go b/deploy/cloud/operator/internal/controller_common/predicate.go index fb6154c09..c7510c5cb 100644 --- a/deploy/cloud/operator/internal/controller_common/predicate.go +++ b/deploy/cloud/operator/internal/controller_common/predicate.go @@ -22,6 +22,7 @@ import ( "strings" "time" + commonconsts "github.com/ai-dynamo/dynamo/deploy/cloud/operator/internal/consts" "k8s.io/apimachinery/pkg/api/meta" "k8s.io/client-go/discovery" ctrl "sigs.k8s.io/controller-runtime" @@ -75,6 +76,9 @@ type Config struct { RBAC RBACConfig // ExcludedNamespaces is a thread-safe set of namespaces to exclude (cluster-wide mode only) ExcludedNamespaces ExcludedNamespacesInterface + + // DiscoveryBackend is the discovery backend to use. By default, will rely on ETCD for discovery. Can be set to "kubernetes" to use Kubernetes API for service discovery. + DiscoveryBackend string } // RBACConfig holds configuration for RBAC management @@ -153,6 +157,19 @@ func detectAPIGroupAvailability(ctx context.Context, mgr ctrl.Manager, groupName return false } +// For DGD, pass in the meta annotations +// For DCD, pass in the spec annotations +func (c Config) IsK8sDiscoveryEnabled(annotations map[string]string) bool { + return c.GetDiscoveryBackend(annotations) == "kubernetes" +} + +func (c Config) GetDiscoveryBackend(annotations map[string]string) string { + if dgdDiscoveryBackend, exists := annotations[commonconsts.KubeAnnotationDynamoDiscoveryBackend]; exists { + return dgdDiscoveryBackend + } + return c.DiscoveryBackend +} + func EphemeralDeploymentEventFilter(config Config) predicate.Predicate { return predicate.NewPredicateFuncs(func(o client.Object) bool { l := log.FromContext(context.Background()) diff --git a/deploy/cloud/operator/internal/discovery/resource.go b/deploy/cloud/operator/internal/discovery/resource.go new file mode 100644 index 000000000..0d0588f0d --- /dev/null +++ b/deploy/cloud/operator/internal/discovery/resource.go @@ -0,0 +1,96 @@ +/* + * SPDX-FileCopyrightText: Copyright (c) 2025 NVIDIA CORPORATION & AFFILIATES. All rights reserved. + * SPDX-License-Identifier: Apache-2.0 + */ + +package discovery + +import ( + "fmt" + + corev1 "k8s.io/api/core/v1" + rbacv1 "k8s.io/api/rbac/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" +) + +const ( + kindServiceAccount = "ServiceAccount" + apiGroupRBAC = "rbac.authorization.k8s.io" + apiGroupCore = "" +) + +func GetK8sDiscoveryServiceAccountName(dgdName string) string { + return fmt.Sprintf("%s-k8s-service-discovery", dgdName) +} + +func GetK8sDiscoveryServiceAccount(dgdName string, namespace string) *corev1.ServiceAccount { + name := GetK8sDiscoveryServiceAccountName(dgdName) + return &corev1.ServiceAccount{ + ObjectMeta: metav1.ObjectMeta{ + Name: name, + Namespace: namespace, + Labels: map[string]string{ + "app.kubernetes.io/managed-by": "dynamo-operator", + "app.kubernetes.io/component": "rbac", + "app.kubernetes.io/name": name, + }, + }, + } +} + +func GetK8sDiscoveryRole(dgdName string, namespace string) *rbacv1.Role { + name := GetK8sDiscoveryServiceAccountName(dgdName) + roleName := name + "-role" + return &rbacv1.Role{ + ObjectMeta: metav1.ObjectMeta{ + Name: roleName, + Namespace: namespace, + Labels: map[string]string{ + "app.kubernetes.io/managed-by": "dynamo-operator", + "app.kubernetes.io/component": "rbac", + "app.kubernetes.io/name": name, + }, + }, + Rules: []rbacv1.PolicyRule{ + { + APIGroups: []string{apiGroupCore}, + Resources: []string{"endpoints"}, + Verbs: []string{"get", "list", "watch"}, + }, + { + APIGroups: []string{"discovery.k8s.io"}, + Resources: []string{"endpointslices"}, + Verbs: []string{"get", "list", "watch"}, + }, + }, + } +} + +func GetK8sDiscoveryRoleBinding(dgdName, namespace string) *rbacv1.RoleBinding { + name := GetK8sDiscoveryServiceAccountName(dgdName) + roleName := name + "-role" + bindingName := name + "-binding" + return &rbacv1.RoleBinding{ + ObjectMeta: metav1.ObjectMeta{ + Name: bindingName, + Namespace: namespace, + Labels: map[string]string{ + "app.kubernetes.io/managed-by": "dynamo-operator", + "app.kubernetes.io/component": "rbac", + "app.kubernetes.io/name": name, + }, + }, + Subjects: []rbacv1.Subject{ + { + Kind: kindServiceAccount, + Name: name, + Namespace: namespace, + }, + }, + RoleRef: rbacv1.RoleRef{ + APIGroup: apiGroupRBAC, + Kind: "Role", + Name: roleName, + }, + } +} diff --git a/deploy/cloud/operator/internal/dynamo/component_common.go b/deploy/cloud/operator/internal/dynamo/component_common.go index c14bbd9bd..cae254352 100644 --- a/deploy/cloud/operator/internal/dynamo/component_common.go +++ b/deploy/cloud/operator/internal/dynamo/component_common.go @@ -27,7 +27,7 @@ func ComponentDefaultsFactory(componentType string) ComponentDefaults { switch componentType { case commonconsts.ComponentTypeFrontend: return NewFrontendDefaults() - case commonconsts.ComponentTypeWorker: + case commonconsts.ComponentTypeWorker, commonconsts.ComponentTypePrefill, commonconsts.ComponentTypeDecode: return NewWorkerDefaults() case commonconsts.ComponentTypePlanner: return NewPlannerDefaults() @@ -42,8 +42,10 @@ type BaseComponentDefaults struct{} type ComponentContext struct { numberOfNodes int32 DynamoNamespace string + ComponentType string ParentGraphDeploymentName string ParentGraphDeploymentNamespace string + DiscoveryBackend string } func (b *BaseComponentDefaults) GetBaseContainer(context ComponentContext) (corev1.Container, error) { @@ -71,9 +73,13 @@ func (b *BaseComponentDefaults) getCommonContainer(context ComponentContext) cor } container.Env = []corev1.EnvVar{ { - Name: "DYN_NAMESPACE", + Name: commonconsts.DynamoNamespaceEnvVar, Value: context.DynamoNamespace, }, + { + Name: commonconsts.DynamoComponentEnvVar, + Value: context.ComponentType, + }, { Name: "DYN_PARENT_DGD_K8S_NAME", Value: context.ParentGraphDeploymentName, @@ -82,6 +88,29 @@ func (b *BaseComponentDefaults) getCommonContainer(context ComponentContext) cor Name: "DYN_PARENT_DGD_K8S_NAMESPACE", Value: context.ParentGraphDeploymentNamespace, }, + { + Name: "POD_NAME", + ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.name", + }, + }, + }, + { + Name: "POD_NAMESPACE", + ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.namespace", + }, + }, + }, + } + + if context.DiscoveryBackend != "" { + container.Env = append(container.Env, corev1.EnvVar{ + Name: commonconsts.DynamoDiscoveryBackendEnvVar, + Value: context.DiscoveryBackend, + }) } return container diff --git a/deploy/cloud/operator/internal/dynamo/component_planner_test.go b/deploy/cloud/operator/internal/dynamo/component_planner_test.go index a56a1c10a..a290e8955 100644 --- a/deploy/cloud/operator/internal/dynamo/component_planner_test.go +++ b/deploy/cloud/operator/internal/dynamo/component_planner_test.go @@ -18,29 +18,24 @@ func TestPlannerDefaults_GetBaseContainer(t *testing.T) { type fields struct { BaseComponentDefaults *BaseComponentDefaults } - type args struct { - numberOfNodes int32 - parentGraphDeploymentName string - parentGraphDeploymentNamespace string - dynamoNamespace string - } tests := []struct { - name string - fields fields - args args - want corev1.Container - wantErr bool + name string + fields fields + componentContext ComponentContext + want corev1.Container + wantErr bool }{ { name: "test", fields: fields{ BaseComponentDefaults: &BaseComponentDefaults{}, }, - args: args{ + componentContext: ComponentContext{ numberOfNodes: 1, - parentGraphDeploymentName: "name", - parentGraphDeploymentNamespace: "namespace", - dynamoNamespace: "dynamo-namespace", + ParentGraphDeploymentName: "name", + ParentGraphDeploymentNamespace: "namespace", + DynamoNamespace: "dynamo-namespace", + ComponentType: commonconsts.ComponentTypePlanner, }, want: corev1.Container{ Name: commonconsts.MainContainerName, @@ -52,9 +47,23 @@ func TestPlannerDefaults_GetBaseContainer(t *testing.T) { {Name: commonconsts.DynamoMetricsPortName, ContainerPort: commonconsts.DynamoPlannerMetricsPort, Protocol: corev1.ProtocolTCP}, }, Env: []corev1.EnvVar{ - {Name: "DYN_NAMESPACE", Value: "dynamo-namespace"}, + {Name: commonconsts.DynamoNamespaceEnvVar, Value: "dynamo-namespace"}, + {Name: commonconsts.DynamoComponentEnvVar, Value: commonconsts.ComponentTypePlanner}, {Name: "DYN_PARENT_DGD_K8S_NAME", Value: "name"}, {Name: "DYN_PARENT_DGD_K8S_NAMESPACE", Value: "namespace"}, + { + Name: "POD_NAME", + ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.name", + }, + }, + }, + {Name: "POD_NAMESPACE", ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.namespace", + }, + }}, {Name: "PLANNER_PROMETHEUS_PORT", Value: fmt.Sprintf("%d", commonconsts.DynamoPlannerMetricsPort)}, }, }, @@ -65,12 +74,7 @@ func TestPlannerDefaults_GetBaseContainer(t *testing.T) { p := &PlannerDefaults{ BaseComponentDefaults: tt.fields.BaseComponentDefaults, } - got, err := p.GetBaseContainer(ComponentContext{ - numberOfNodes: tt.args.numberOfNodes, - ParentGraphDeploymentName: tt.args.parentGraphDeploymentName, - ParentGraphDeploymentNamespace: tt.args.parentGraphDeploymentNamespace, - DynamoNamespace: tt.args.dynamoNamespace, - }) + got, err := p.GetBaseContainer(tt.componentContext) if (err != nil) != tt.wantErr { t.Errorf("PlannerDefaults.GetBaseContainer() error = %v, wantErr %v", err, tt.wantErr) return diff --git a/deploy/cloud/operator/internal/dynamo/graph.go b/deploy/cloud/operator/internal/dynamo/graph.go index 8ff97895d..c5436df62 100644 --- a/deploy/cloud/operator/internal/dynamo/graph.go +++ b/deploy/cloud/operator/internal/dynamo/graph.go @@ -38,6 +38,7 @@ import ( "github.com/ai-dynamo/dynamo/deploy/cloud/operator/api/v1alpha1" commonconsts "github.com/ai-dynamo/dynamo/deploy/cloud/operator/internal/consts" "github.com/ai-dynamo/dynamo/deploy/cloud/operator/internal/controller_common" + "github.com/ai-dynamo/dynamo/deploy/cloud/operator/internal/discovery" "github.com/imdario/mergo" networkingv1beta1 "istio.io/client-go/pkg/apis/networking/v1beta1" corev1 "k8s.io/api/core/v1" @@ -135,12 +136,15 @@ func GenerateDynamoComponentsDeployments(ctx context.Context, parentDynamoGraphD // Propagate metrics annotation from parent deployment if present if parentDynamoGraphDeployment.Annotations != nil { + if deployment.Spec.Annotations == nil { + deployment.Spec.Annotations = make(map[string]string) + } if val, exists := parentDynamoGraphDeployment.Annotations[commonconsts.KubeAnnotationEnableMetrics]; exists { - if deployment.Spec.Annotations == nil { - deployment.Spec.Annotations = make(map[string]string) - } deployment.Spec.Annotations[commonconsts.KubeAnnotationEnableMetrics] = val } + if val, exists := parentDynamoGraphDeployment.Annotations[commonconsts.KubeAnnotationDynamoDiscoveryBackend]; exists { + deployment.Spec.Annotations[commonconsts.KubeAnnotationDynamoDiscoveryBackend] = val + } } if component.ComponentType == commonconsts.ComponentTypePlanner { @@ -392,24 +396,39 @@ func getCliqueStartupDependencies( return nil } -func GenerateComponentService(ctx context.Context, componentName, componentNamespace string) (*corev1.Service, error) { +func GenerateComponentService(ctx context.Context, dynamoDeployment *v1alpha1.DynamoGraphDeployment, component *v1alpha1.DynamoComponentDeploymentSharedSpec, componentName string) (*corev1.Service, error) { + if component.DynamoNamespace == nil { + return nil, fmt.Errorf("expected DynamoComponentDeployment %s to have a dynamoNamespace", componentName) + } + componentName = GetDynamoComponentName(dynamoDeployment, componentName) + + var servicePort corev1.ServicePort + if component.ComponentType == commonconsts.ComponentTypeFrontend { + servicePort = corev1.ServicePort{ + Name: commonconsts.DynamoServicePortName, + Port: commonconsts.DynamoServicePort, + TargetPort: intstr.FromString(commonconsts.DynamoContainerPortName), + Protocol: corev1.ProtocolTCP, + } + } else { + servicePort = corev1.ServicePort{ + Name: commonconsts.DynamoSystemPortName, + Port: commonconsts.DynamoSystemPort, + TargetPort: intstr.FromString(commonconsts.DynamoSystemPortName), + Protocol: corev1.ProtocolTCP, + } + } service := &corev1.Service{ ObjectMeta: metav1.ObjectMeta{ Name: componentName, - Namespace: componentNamespace, + Namespace: dynamoDeployment.Namespace, }, Spec: corev1.ServiceSpec{ Selector: map[string]string{ - commonconsts.KubeLabelDynamoSelector: componentName, - }, - Ports: []corev1.ServicePort{ - { - Name: commonconsts.DynamoServicePortName, - Port: commonconsts.DynamoServicePort, - TargetPort: intstr.FromString(commonconsts.DynamoContainerPortName), - Protocol: corev1.ProtocolTCP, - }, + commonconsts.KubeLabelDynamoComponentType: component.ComponentType, + commonconsts.KubeLabelDynamoNamespace: *component.DynamoNamespace, }, + Ports: []corev1.ServicePort{servicePort}, }, } return service, nil @@ -631,7 +650,9 @@ func MultinodeDeployerFactory(multinodeDeploymentType commonconsts.MultinodeDepl // isWorkerComponent checks if a component is a worker that needs backend framework detection func isWorkerComponent(componentType string) bool { - return componentType == commonconsts.ComponentTypeWorker + return componentType == commonconsts.ComponentTypeWorker || + componentType == commonconsts.ComponentTypePrefill || + componentType == commonconsts.ComponentTypeDecode } // addStandardEnvVars adds the standard environment variables that are common to both Grove and Controller @@ -685,7 +706,7 @@ func GenerateBasePodSpec( serviceName string, ) (*corev1.PodSpec, error) { // Start with base container generated per component type - componentContext := generateComponentContext(component, parentGraphDeploymentName, namespace, numberOfNodes) + componentContext := generateComponentContext(component, parentGraphDeploymentName, namespace, numberOfNodes, controllerConfig.GetDiscoveryBackend(component.Annotations)) componentDefaults := ComponentDefaultsFactory(component.ComponentType) container, err := componentDefaults.GetBaseContainer(componentContext) if err != nil { @@ -833,6 +854,11 @@ func GenerateBasePodSpec( return nil, fmt.Errorf("failed to merge extraPodSpec: %w", err) } } + + if controllerConfig.IsK8sDiscoveryEnabled(component.Annotations) { + podSpec.ServiceAccountName = discovery.GetK8sDiscoveryServiceAccountName(parentGraphDeploymentName) + } + podSpec.Containers = append(podSpec.Containers, container) podSpec.Volumes = append(podSpec.Volumes, volumes...) podSpec.ImagePullSecrets = append(podSpec.ImagePullSecrets, imagePullSecrets...) @@ -852,11 +878,13 @@ func setMetricsLabels(labels map[string]string, dynamoGraphDeployment *v1alpha1. labels[commonconsts.KubeLabelMetricsEnabled] = commonconsts.KubeLabelValueTrue } -func generateComponentContext(component *v1alpha1.DynamoComponentDeploymentSharedSpec, parentGraphDeploymentName string, namespace string, numberOfNodes int32) ComponentContext { +func generateComponentContext(component *v1alpha1.DynamoComponentDeploymentSharedSpec, parentGraphDeploymentName string, namespace string, numberOfNodes int32, discoveryBackend string) ComponentContext { componentContext := ComponentContext{ numberOfNodes: numberOfNodes, + ComponentType: component.ComponentType, ParentGraphDeploymentName: parentGraphDeploymentName, ParentGraphDeploymentNamespace: namespace, + DiscoveryBackend: discoveryBackend, } if component.DynamoNamespace != nil { componentContext.DynamoNamespace = *component.DynamoNamespace @@ -915,6 +943,8 @@ func GenerateGrovePodCliqueSet( } } + discoveryBackend := controllerConfig.GetDiscoveryBackend(dynamoDeployment.Annotations) + var scalingGroups []grovev1alpha1.PodCliqueScalingGroupConfig for serviceName, component := range dynamoDeployment.Spec.Services { dynamoNamespace := getDynamoNamespace(dynamoDeployment, component) @@ -925,6 +955,13 @@ func GenerateGrovePodCliqueSet( return nil, fmt.Errorf("failed to determine backend framework for service %s: %w", serviceName, err) } + if discoveryBackend != "" { + if component.Annotations == nil { + component.Annotations = make(map[string]string) + } + component.Annotations[commonconsts.KubeAnnotationDynamoDiscoveryBackend] = discoveryBackend + } + numberOfNodes := component.GetNumberOfNodes() isMultinode := numberOfNodes > 1 roles := expandRolesForService(serviceName, component.Replicas, numberOfNodes) diff --git a/deploy/cloud/operator/internal/dynamo/graph_test.go b/deploy/cloud/operator/internal/dynamo/graph_test.go index 553006d76..7afc60a5e 100644 --- a/deploy/cloud/operator/internal/dynamo/graph_test.go +++ b/deploy/cloud/operator/internal/dynamo/graph_test.go @@ -658,6 +658,78 @@ func TestGenerateDynamoComponentsDeployments(t *testing.T) { }, wantErr: false, }, + { + name: "Test GenerateDynamoComponentsDeployments with Discover Backend and Metrics Annotatitions", + args: args{ + parentDynamoGraphDeployment: &v1alpha1.DynamoGraphDeployment{ + ObjectMeta: metav1.ObjectMeta{ + Name: "test-dynamographdeployment", + Namespace: "default", + Annotations: map[string]string{ + commonconsts.KubeAnnotationEnableMetrics: "false", + commonconsts.KubeAnnotationDynamoDiscoveryBackend: "test", + }, + }, + Spec: v1alpha1.DynamoGraphDeploymentSpec{ + BackendFramework: string(BackendFrameworkSGLang), + Services: map[string]*v1alpha1.DynamoComponentDeploymentSharedSpec{ + "service1": { + DynamoNamespace: &[]string{"default-test-dynamographdeployment"}[0], + ComponentType: "frontend", + Replicas: &[]int32{3}[0], + Resources: &common.Resources{ + Requests: &common.ResourceItem{ + CPU: "1", + Memory: "1Gi", + GPU: "0", + Custom: map[string]string{}, + }, + }, + }, + }, + }, + }, + }, + want: map[string]*v1alpha1.DynamoComponentDeployment{ + "service1": { + ObjectMeta: metav1.ObjectMeta{ + Name: "test-dynamographdeployment-service1", + Namespace: "default", + Labels: map[string]string{ + commonconsts.KubeLabelDynamoComponent: "service1", + commonconsts.KubeLabelDynamoNamespace: "default-test-dynamographdeployment", + commonconsts.KubeLabelDynamoGraphDeploymentName: "test-dynamographdeployment", + }, + }, + Spec: v1alpha1.DynamoComponentDeploymentSpec{ + BackendFramework: string(BackendFrameworkSGLang), + DynamoComponentDeploymentSharedSpec: v1alpha1.DynamoComponentDeploymentSharedSpec{ + Annotations: map[string]string{ + commonconsts.KubeAnnotationEnableMetrics: "false", + commonconsts.KubeAnnotationDynamoDiscoveryBackend: "test", + }, + ServiceName: "service1", + DynamoNamespace: &[]string{"default-test-dynamographdeployment"}[0], + ComponentType: "frontend", + Replicas: &[]int32{3}[0], + Resources: &common.Resources{ + Requests: &common.ResourceItem{ + CPU: "1", + Memory: "1Gi", + GPU: "0", + Custom: map[string]string{}, + }, + }, + Labels: map[string]string{ + commonconsts.KubeLabelDynamoComponent: "service1", + commonconsts.KubeLabelDynamoNamespace: "default-test-dynamographdeployment", + commonconsts.KubeLabelDynamoGraphDeploymentName: "test-dynamographdeployment", + }, + }, + }, + }, + }, + }, } for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { @@ -1330,14 +1402,34 @@ func TestGenerateGrovePodCliqueSet(t *testing.T) { Name: "NATS_SERVER", Value: "nats-address", }, + { + Name: "POD_NAME", + ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.name", + }, + }, + }, + { + Name: "POD_NAMESPACE", + ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.namespace", + }, + }, + }, { Name: "ETCD_ENDPOINTS", Value: "etcd-address", }, { - Name: "DYN_NAMESPACE", + Name: commonconsts.DynamoNamespaceEnvVar, Value: "test-namespace-test-dynamo-graph-deployment", }, + { + Name: commonconsts.DynamoComponentEnvVar, + Value: commonconsts.ComponentTypeFrontend, + }, { Name: "DYN_PARENT_DGD_K8S_NAME", Value: "test-dynamo-graph-deployment", @@ -1478,9 +1570,13 @@ func TestGenerateGrovePodCliqueSet(t *testing.T) { Value: "etcd-address", }, { - Name: "DYN_NAMESPACE", + Name: commonconsts.DynamoNamespaceEnvVar, Value: "test-namespace-test-dynamo-graph-deployment", }, + { + Name: commonconsts.DynamoComponentEnvVar, + Value: commonconsts.ComponentTypePlanner, + }, { Name: "DYN_PARENT_DGD_K8S_NAME", Value: "test-dynamo-graph-deployment", @@ -1501,6 +1597,22 @@ func TestGenerateGrovePodCliqueSet(t *testing.T) { Name: "PLANNER_PROMETHEUS_PORT", Value: fmt.Sprintf("%d", commonconsts.DynamoPlannerMetricsPort), }, + { + Name: "POD_NAME", + ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.name", + }, + }, + }, + { + Name: "POD_NAMESPACE", + ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.namespace", + }, + }, + }, }, Resources: corev1.ResourceRequirements{ Requests: corev1.ResourceList{ @@ -1839,9 +1951,13 @@ func TestGenerateGrovePodCliqueSet(t *testing.T) { Value: "etcd-address", }, { - Name: "DYN_NAMESPACE", + Name: commonconsts.DynamoNamespaceEnvVar, Value: "test-namespace-test-dynamo-graph-deployment", }, + { + Name: commonconsts.DynamoComponentEnvVar, + Value: commonconsts.ComponentTypeWorker, + }, { Name: "DYN_PARENT_DGD_K8S_NAME", Value: "test-dynamo-graph-deployment", @@ -1850,6 +1966,22 @@ func TestGenerateGrovePodCliqueSet(t *testing.T) { Name: "DYN_PARENT_DGD_K8S_NAMESPACE", Value: "test-namespace", }, + { + Name: "POD_NAME", + ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.name", + }, + }, + }, + { + Name: "POD_NAMESPACE", + ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.namespace", + }, + }, + }, }, Resources: corev1.ResourceRequirements{ Requests: corev1.ResourceList{ @@ -1988,9 +2120,13 @@ func TestGenerateGrovePodCliqueSet(t *testing.T) { Value: "etcd-address", }, { - Name: "DYN_NAMESPACE", + Name: commonconsts.DynamoNamespaceEnvVar, Value: "test-namespace-test-dynamo-graph-deployment", }, + { + Name: commonconsts.DynamoComponentEnvVar, + Value: commonconsts.ComponentTypeWorker, + }, { Name: "DYN_PARENT_DGD_K8S_NAME", Value: "test-dynamo-graph-deployment", @@ -1999,6 +2135,22 @@ func TestGenerateGrovePodCliqueSet(t *testing.T) { Name: "DYN_PARENT_DGD_K8S_NAMESPACE", Value: "test-namespace", }, + { + Name: "POD_NAME", + ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.name", + }, + }, + }, + { + Name: "POD_NAMESPACE", + ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.namespace", + }, + }, + }, }, Resources: corev1.ResourceRequirements{ Requests: corev1.ResourceList{ @@ -2119,9 +2271,13 @@ func TestGenerateGrovePodCliqueSet(t *testing.T) { Value: "etcd-address", }, { - Name: "DYN_NAMESPACE", + Name: commonconsts.DynamoNamespaceEnvVar, Value: "test-namespace-test-dynamo-graph-deployment", }, + { + Name: commonconsts.DynamoComponentEnvVar, + Value: commonconsts.ComponentTypeFrontend, + }, { Name: "DYN_PARENT_DGD_K8S_NAME", Value: "test-dynamo-graph-deployment", @@ -2130,6 +2286,22 @@ func TestGenerateGrovePodCliqueSet(t *testing.T) { Name: "DYN_PARENT_DGD_K8S_NAMESPACE", Value: "test-namespace", }, + { + Name: "POD_NAME", + ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.name", + }, + }, + }, + { + Name: "POD_NAMESPACE", + ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.namespace", + }, + }, + }, }, Resources: corev1.ResourceRequirements{ Requests: corev1.ResourceList{ @@ -2253,9 +2425,13 @@ func TestGenerateGrovePodCliqueSet(t *testing.T) { Value: "etcd-address", }, { - Name: "DYN_NAMESPACE", + Name: commonconsts.DynamoNamespaceEnvVar, Value: "test-namespace-test-dynamo-graph-deployment", }, + { + Name: commonconsts.DynamoComponentEnvVar, + Value: commonconsts.ComponentTypePlanner, + }, { Name: "DYN_PARENT_DGD_K8S_NAME", Value: "test-dynamo-graph-deployment", @@ -2268,6 +2444,22 @@ func TestGenerateGrovePodCliqueSet(t *testing.T) { Name: "PLANNER_PROMETHEUS_PORT", Value: fmt.Sprintf("%d", commonconsts.DynamoPlannerMetricsPort), }, + { + Name: "POD_NAME", + ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.name", + }, + }, + }, + { + Name: "POD_NAMESPACE", + ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.namespace", + }, + }, + }, }, Resources: corev1.ResourceRequirements{ Requests: corev1.ResourceList{ @@ -2636,9 +2828,13 @@ func TestGenerateGrovePodCliqueSet(t *testing.T) { Value: "etcd-address", }, { - Name: "DYN_NAMESPACE", + Name: commonconsts.DynamoNamespaceEnvVar, Value: "test-namespace-test-dynamo-graph-deployment", }, + { + Name: commonconsts.DynamoComponentEnvVar, + Value: commonconsts.ComponentTypeWorker, + }, { Name: "DYN_PARENT_DGD_K8S_NAME", Value: "test-dynamo-graph-deployment", @@ -2647,6 +2843,22 @@ func TestGenerateGrovePodCliqueSet(t *testing.T) { Name: "DYN_PARENT_DGD_K8S_NAMESPACE", Value: "test-namespace", }, + { + Name: "POD_NAME", + ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.name", + }, + }, + }, + { + Name: "POD_NAMESPACE", + ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.namespace", + }, + }, + }, }, Resources: corev1.ResourceRequirements{ Requests: corev1.ResourceList{ @@ -2772,9 +2984,13 @@ func TestGenerateGrovePodCliqueSet(t *testing.T) { Value: "etcd-address", }, { - Name: "DYN_NAMESPACE", + Name: commonconsts.DynamoNamespaceEnvVar, Value: "test-namespace-test-dynamo-graph-deployment", }, + { + Name: commonconsts.DynamoComponentEnvVar, + Value: commonconsts.ComponentTypeWorker, + }, { Name: "DYN_PARENT_DGD_K8S_NAME", Value: "test-dynamo-graph-deployment", @@ -2783,6 +2999,22 @@ func TestGenerateGrovePodCliqueSet(t *testing.T) { Name: "DYN_PARENT_DGD_K8S_NAMESPACE", Value: "test-namespace", }, + { + Name: "POD_NAME", + ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.name", + }, + }, + }, + { + Name: "POD_NAMESPACE", + ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.namespace", + }, + }, + }, }, Resources: corev1.ResourceRequirements{ Requests: corev1.ResourceList{ @@ -2903,9 +3135,13 @@ func TestGenerateGrovePodCliqueSet(t *testing.T) { Value: "etcd-address", }, { - Name: "DYN_NAMESPACE", + Name: commonconsts.DynamoNamespaceEnvVar, Value: "test-namespace-test-dynamo-graph-deployment", }, + { + Name: commonconsts.DynamoComponentEnvVar, + Value: commonconsts.ComponentTypeFrontend, + }, { Name: "DYN_PARENT_DGD_K8S_NAME", Value: "test-dynamo-graph-deployment", @@ -2914,6 +3150,22 @@ func TestGenerateGrovePodCliqueSet(t *testing.T) { Name: "DYN_PARENT_DGD_K8S_NAMESPACE", Value: "test-namespace", }, + { + Name: "POD_NAME", + ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.name", + }, + }, + }, + { + Name: "POD_NAMESPACE", + ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.namespace", + }, + }, + }, }, Resources: corev1.ResourceRequirements{ Requests: corev1.ResourceList{ @@ -3044,9 +3296,13 @@ func TestGenerateGrovePodCliqueSet(t *testing.T) { Value: "etcd-address", }, { - Name: "DYN_NAMESPACE", + Name: commonconsts.DynamoNamespaceEnvVar, Value: "test-namespace-test-dynamo-graph-deployment", }, + { + Name: commonconsts.DynamoComponentEnvVar, + Value: commonconsts.ComponentTypePlanner, + }, { Name: "DYN_PARENT_DGD_K8S_NAME", Value: "test-dynamo-graph-deployment", @@ -3059,6 +3315,22 @@ func TestGenerateGrovePodCliqueSet(t *testing.T) { Name: "PLANNER_PROMETHEUS_PORT", Value: fmt.Sprintf("%d", commonconsts.DynamoPlannerMetricsPort), }, + { + Name: "POD_NAME", + ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.name", + }, + }, + }, + { + Name: "POD_NAMESPACE", + ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.namespace", + }, + }, + }, }, Resources: corev1.ResourceRequirements{ Requests: corev1.ResourceList{ @@ -3175,13 +3447,13 @@ func assertDYNNamespace(t *testing.T, podSpec corev1.PodSpec, expectedNamespace if assert.Len(t, podSpec.Containers, 1) { foundDYNNamespace := false for _, env := range podSpec.Containers[0].Env { - if env.Name == "DYN_NAMESPACE" { + if env.Name == commonconsts.DynamoNamespaceEnvVar { assert.Equal(t, expectedNamespace, env.Value) foundDYNNamespace = true break } } - assert.True(t, foundDYNNamespace, "DYN_NAMESPACE not found in container environment variables") + assert.True(t, foundDYNNamespace, fmt.Sprintf("%s not found in container environment variables", commonconsts.DynamoNamespaceEnvVar)) } } @@ -4541,6 +4813,91 @@ func TestGenerateBasePodSpec_DisableImagePullSecretDiscovery(t *testing.T) { } } +func TestGenerateBasePodSpec_DiscoverBackend(t *testing.T) { + tests := []struct { + name string + component *v1alpha1.DynamoComponentDeploymentSharedSpec + controllerConfig controller_common.Config + wantEnvVar string + }{ + { + name: "Discover backend should be set", + component: &v1alpha1.DynamoComponentDeploymentSharedSpec{ + Annotations: map[string]string{ + commonconsts.KubeAnnotationDynamoDiscoveryBackend: "kubernetes", + }, + }, + wantEnvVar: "kubernetes", + }, + { + name: "Discover backend should override the controller config", + component: &v1alpha1.DynamoComponentDeploymentSharedSpec{ + Annotations: map[string]string{ + commonconsts.KubeAnnotationDynamoDiscoveryBackend: "test", + }, + }, + controllerConfig: controller_common.Config{ + DiscoveryBackend: "etcd", + }, + wantEnvVar: "test", + }, + { + name: "Discover backend should be set by the controller config", + component: &v1alpha1.DynamoComponentDeploymentSharedSpec{ + Annotations: map[string]string{}, + }, + controllerConfig: controller_common.Config{ + DiscoveryBackend: "etcd", + }, + wantEnvVar: "etcd", + }, + { + name: "Discover backend empty string", + component: &v1alpha1.DynamoComponentDeploymentSharedSpec{ + Annotations: map[string]string{ + commonconsts.KubeAnnotationDynamoDiscoveryBackend: "", + }, + }, + controllerConfig: controller_common.Config{ + DiscoveryBackend: "", + }, + }, + { + name: "Discover backend not set", + component: &v1alpha1.DynamoComponentDeploymentSharedSpec{}, + }, + } + secretsRetriever := &mockSecretsRetriever{} + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + podSpec, err := GenerateBasePodSpec( + tt.component, + BackendFrameworkSGLang, + secretsRetriever, + "test-deployment", + "default", + RoleMain, + 1, + tt.controllerConfig, + commonconsts.MultinodeDeploymentTypeGrove, + "test-service", + ) + if !assert.NoError(t, err) { + return + } + if tt.wantEnvVar != "" { + assert.Contains(t, podSpec.Containers[0].Env, corev1.EnvVar{Name: commonconsts.DynamoDiscoveryBackendEnvVar, Value: tt.wantEnvVar}) + } else { + for _, env := range podSpec.Containers[0].Env { + if env.Name == commonconsts.DynamoDiscoveryBackendEnvVar { + t.Errorf("GenerateBasePodSpec() Discover backend env var should not be set, got %s", env.Value) + } + } + } + }) + } +} + func TestGenerateBasePodSpec_Worker(t *testing.T) { secretsRetriever := &mockSecretsRetriever{} controllerConfig := controller_common.Config{} @@ -4576,11 +4933,22 @@ func TestGenerateBasePodSpec_Worker(t *testing.T) { Env: []corev1.EnvVar{ {Name: "ANOTHER_COMPONENTENV", Value: "true"}, {Name: "ANOTHER_CONTAINER_ENV", Value: "true"}, - {Name: "DYN_NAMESPACE", Value: ""}, + {Name: commonconsts.DynamoComponentEnvVar, Value: "worker"}, + {Name: commonconsts.DynamoNamespaceEnvVar, Value: ""}, {Name: "DYN_PARENT_DGD_K8S_NAME", Value: "test-deployment"}, {Name: "DYN_PARENT_DGD_K8S_NAMESPACE", Value: "default"}, {Name: "DYN_SYSTEM_PORT", Value: "9090"}, {Name: "DYN_SYSTEM_USE_ENDPOINT_HEALTH_STATUS", Value: "[\"generate\"]"}, + {Name: "POD_NAME", ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.name", + }, + }}, + {Name: "POD_NAMESPACE", ValueFrom: &corev1.EnvVarSource{ + FieldRef: &corev1.ObjectFieldSelector{ + FieldPath: "metadata.namespace", + }, + }}, }, VolumeMounts: []corev1.VolumeMount{ {