From ac0f917601760611a3741e73dbab85264d3514b4 Mon Sep 17 00:00:00 2001 From: Dmitry Shmulevich <17212177+dmitsh@users.noreply.github.com> Date: Tue, 11 Aug 2026 15:56:34 -0700 Subject: [PATCH] refactor(providers): share Kubernetes accelerator discovery Signed-off-by: Dmitry Shmulevich <17212177+dmitsh@users.noreply.github.com> --- CHANGELOG.md | 2 + charts/topograph/templates/_validation.tpl | 8 +- charts/topograph/tests/validation_test.yaml | 22 ++++ charts/topograph/values.yaml | 12 +- demos/dra-slinky/values.dra-slinky.kwok.yaml | 4 + demos/oci-sim-slinky/demo.sh | 2 + docs/api.md | 6 +- docs/overview.md | 4 +- docs/providers/dra.md | 42 +++++-- docs/reference/node-labels.md | 2 +- internal/k8s/utils.go | 19 ++++ internal/k8s/utils_test.go | 38 +++++++ pkg/accelerator/accelerator.go | 58 ++++++++++ pkg/accelerator/accelerator_test.go | 64 +++++++++++ pkg/accelerator/kubernetes.go | 66 +++++++++++ pkg/providers/dra/provider.go | 105 +++++++----------- pkg/providers/dra/provider_test.go | 77 +++++++------ pkg/providers/infiniband/bm.go | 28 ----- pkg/providers/infiniband/bm_test.go | 30 ----- pkg/providers/infiniband/k8s.go | 5 +- pkg/providers/infiniband/provider_bm.go | 4 +- pkg/providers/infiniband/provider_k8s.go | 45 +------- pkg/providers/infiniband/provider_k8s_test.go | 46 ++------ pkg/topology/graph_test.go | 8 +- pkg/topology/topology.go | 1 - tests/ci/values.slinky-dra-block.yaml | 4 + 26 files changed, 436 insertions(+), 266 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 10edc243..fdbe534d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -8,6 +8,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/). ### Added +- The DRA provider now accepts the shared `provider.params.accelerator` Kubernetes-label configuration used by `infiniband-k8s`, allowing a custom Node label to supply accelerator domains while preserving `nvidia.com/gpu.clique` as the default when the section is omitted. - Pluggable accelerator-domain discovery for InfiniBand providers, independently selectable from fabric discovery with `nvidia-smi`, an explicitly configured Kubernetes Node label, or no accelerator source. Discovery is disabled when `accelerator` is omitted or empty; a non-empty section must set `source` explicitly. Helm defaults the `nvidia-smi` workload location to the `gpu-operator` namespace and `nvidia-device-plugin-daemonset` DaemonSet when those values are omitted. - Helm `kubeClient.qps` and `kubeClient.burst` values for tuning the DRA provider and the Kubernetes, NFD, and Slinky engine clients through deployment-level `KUBE_QPS` and `KUBE_BURST` settings. - The Kubernetes engine now publishes `accelerator.topograph.run/sub-domain` when a provider supplies `InstanceTopology.XclrSubDomainID`. @@ -70,6 +71,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.1.0/). ### Removed +- **BREAKING:** Removed the exported Go constant `topology.KeyNvidiaGPUClique`. The GPU Operator's `nvidia.com/gpu.clique` label is now a DRA provider default rather than part of the canonical topology API; configure label-backed accelerator discovery through `provider.params.accelerator.kubernetesLabel.key`. - Periodic node annotation refreshes, including the `--refresh-interval` broker flag and `nodeDataBroker.refreshInterval` Helm value. The broker now applies annotations once at startup. ### Security diff --git a/charts/topograph/templates/_validation.tpl b/charts/topograph/templates/_validation.tpl index b9bba24a..e50692b6 100644 --- a/charts/topograph/templates/_validation.tpl +++ b/charts/topograph/templates/_validation.tpl @@ -29,7 +29,7 @@ {{- fail "env.KUBE_BURST is managed by the chart; configure kubeClient.burst instead" }} {{- end }} -{{- if or (eq .Values.provider.name "infiniband-k8s") (eq .Values.provider.name "infiniband-bm") }} +{{- if or (eq .Values.provider.name "infiniband-k8s") (eq .Values.provider.name "infiniband-bm") (eq .Values.provider.name "dra") }} {{- $params := default dict .Values.provider.params }} {{- $acceleratorValue := get $params "accelerator" }} {{- $accelerator := default dict $acceleratorValue }} @@ -49,7 +49,7 @@ {{- fail (printf "unsupported provider.params.accelerator.source %q" $source) }} {{- end }} -{{- if and (eq .Values.provider.name "infiniband-k8s") (eq $source "kubernetes-label") }} +{{- if and (or (eq .Values.provider.name "infiniband-k8s") (eq .Values.provider.name "dra")) (eq $source "kubernetes-label") }} {{- $kubernetesLabel := default dict (get $accelerator "kubernetesLabel") }} {{- $key := trim (toString (get $kubernetesLabel "key")) }} {{- if eq $key "" }} @@ -61,6 +61,10 @@ {{- fail "provider.params.accelerator.source kubernetes-label is not supported by infiniband-bm" }} {{- end }} +{{- if and (eq .Values.provider.name "dra") (hasKey $params "accelerator") (ne $source "kubernetes-label") }} + {{- fail "provider.params.accelerator.source must be kubernetes-label for the dra provider" }} +{{- end }} + {{- end }} {{- if eq .Values.provider.name "gcp" }} diff --git a/charts/topograph/tests/validation_test.yaml b/charts/topograph/tests/validation_test.yaml index 74161fcb..daae83a3 100644 --- a/charts/topograph/tests/validation_test.yaml +++ b/charts/topograph/tests/validation_test.yaml @@ -121,3 +121,25 @@ tests: asserts: - failedTemplate: errorMessage: "provider.params.accelerator.kubernetesLabel.key must be set for source kubernetes-label" + + - it: rejects a DRA label accelerator source without a key + set: + provider: + name: dra + params: + accelerator: + source: kubernetes-label + asserts: + - failedTemplate: + errorMessage: "provider.params.accelerator.kubernetesLabel.key must be set for source kubernetes-label" + + - it: rejects a non-label accelerator source for DRA + set: + provider: + name: dra + params: + accelerator: + source: none + asserts: + - failedTemplate: + errorMessage: "provider.params.accelerator.source must be kubernetes-label for the dra provider" diff --git a/charts/topograph/values.yaml b/charts/topograph/values.yaml index 5e7f4c8c..fe3d819e 100644 --- a/charts/topograph/values.yaml +++ b/charts/topograph/values.yaml @@ -8,13 +8,15 @@ provider: # params: # accelerator: # # Accelerator-domain discovery is independent of fabric discovery. - # # Sources supported by infiniband-k8s: nvidia-smi, - # # kubernetes-label, and none. Omitting the accelerator section disables - # # accelerator discovery; an empty accelerator section is equivalent to - # # source: none. A non-empty section must set source explicitly. + # # Sources supported by infiniband-k8s: nvidia-smi, kubernetes-label, + # # and none. The dra provider supports kubernetes-label and defaults to + # # nvidia.com/gpu.clique when this section is omitted. For InfiniBand, + # # omitting accelerator disables discovery; an empty section is + # # equivalent to source: none. A non-empty section must set source. # source: kubernetes-label # kubernetesLabel: - # # Required for the kubernetes-label source; there is no default key. + # # Required for the kubernetes-label source. DRA defaults this key to + # # nvidia.com/gpu.clique only when the accelerator section is omitted. # key: nvidia.com/gpu.clique # nvidiaSmi: # # Optional; defaults used by the Helm-managed node-data-broker are diff --git a/demos/dra-slinky/values.dra-slinky.kwok.yaml b/demos/dra-slinky/values.dra-slinky.kwok.yaml index 48d1c72a..fa872bdf 100644 --- a/demos/dra-slinky/values.dra-slinky.kwok.yaml +++ b/demos/dra-slinky/values.dra-slinky.kwok.yaml @@ -1,6 +1,10 @@ provider: name: dra params: + accelerator: + source: kubernetes-label + kubernetesLabel: + key: nvidia.com/gpu.clique nodeSelector: kwok.x-k8s.io/node: "fake" engine: diff --git a/demos/oci-sim-slinky/demo.sh b/demos/oci-sim-slinky/demo.sh index cce510f4..181b771d 100755 --- a/demos/oci-sim-slinky/demo.sh +++ b/demos/oci-sim-slinky/demo.sh @@ -13,6 +13,8 @@ cd "$demo_dir/../.." source demos/utils.sh +step "make build TARGETS=kwok-nodes" + step "delete_cluster" step "kind create cluster --name \"${KUBE_CONTEXT#kind-}\" --wait 120s" diff --git a/docs/api.md b/docs/api.md index 80e5b44a..5bec19d2 100644 --- a/docs/api.md +++ b/docs/api.md @@ -72,9 +72,9 @@ Topograph exposes three endpoints for interacting with the service. Below are th - **name**: (optional) A string specifying the Service Provider, such as `aws`, `oci`, `gcp`, `nebius`, `nscale`, `netq`, `dra`, `infiniband-k8s`, `infiniband-bm` or `test`. This parameter will override the provider set in the topograph config. - **creds**: (optional) A key-value map with provider-specific parameters for authentication. - **params**: (optional) A key-value map with provider-specific parameters. The `test` provider uses these parameters for response simulation; for complete behavior and examples, see [Test Mode and Test Provider](./providers/test.md). - - **accelerator**: (optional) Used in: [`infiniband-k8s`, `infiniband-bm`]. Configures accelerator-domain discovery independently of network-fabric discovery. Omitting this section or setting it to an empty object disables accelerator-domain discovery. - - **source**: (required when `accelerator` is non-empty) `nvidia-smi`, `kubernetes-label` (`infiniband-k8s` only), or `none`. An empty object is equivalent to `source: none`. - - **kubernetesLabel.key**: (required for `kubernetes-label`) Kubernetes Node label read as the accelerator-domain ID. No default is assumed. + - **accelerator**: (optional) Used in: [`dra`, `infiniband-k8s`, `infiniband-bm`]. Configures accelerator-domain discovery independently of network-fabric discovery. For InfiniBand, omitting this section or setting it to an empty object disables accelerator-domain discovery. DRA supports only `kubernetes-label` and retains its legacy `nvidia.com/gpu.clique` default when the section is omitted. + - **source**: (required when `accelerator` is non-empty) `nvidia-smi`, `kubernetes-label` (`dra` and `infiniband-k8s`), or `none`. DRA accepts only `kubernetes-label`. + - **kubernetesLabel.key**: (required for `kubernetes-label`) Kubernetes Node label read as the accelerator-domain ID. No default is assumed for an explicit section. - For `infiniband-k8s`, a request with `source: nvidia-smi` reads accelerator-domain annotations previously collected by the node-data-broker. The request does not run `nvidia-smi` or reconfigure the broker. Deploy the broker with the same accelerator source before sending the request; see [Helm node-data-broker settings](./providers/infiniband.md#helm-node-data-broker-settings). - **engine**: (optional) Selects the topology output and provides any engine-specific parameters. - **name**: (optional) A string specifying the topology output, either `slurm`, `k8s`, `nfd`, `slinky`, or `graph`. This parameter will override the engine set in the topograph config. diff --git a/docs/overview.md b/docs/overview.md index 93e552f1..751ed06f 100644 --- a/docs/overview.md +++ b/docs/overview.md @@ -40,7 +40,7 @@ Currently supported providers: - [Nscale](./providers/nscale.md) - [Lambda](./providers/lambdai.md) - [NetQ](./providers/netq.md) -- [DRA](./providers/dra.md) — provides Slinky block topology from pre-existing `nvidia.com/gpu.clique` labels; it does not discover the backend switch fabric +- [DRA](./providers/dra.md) — provides Slinky block topology from a configured pre-existing Node label (default `nvidia.com/gpu.clique`); it does not discover the backend switch fabric - [InfiniBand (bare-metal)](./providers/infiniband.md#infiniband-bm-bare-metal) - [InfiniBand (Kubernetes)](./providers/infiniband.md#infiniband-k8s-kubernetes) - [Test](./providers/test.md) - simulates Topograph success, pending, and error responses for integration testing @@ -66,7 +66,7 @@ Currently supported engines: | InfiniBand fabric, no NetQ, Kubernetes | [InfiniBand (Kubernetes)](./providers/infiniband.md) | | Client integration and regression testing | [Test](./providers/test.md) | -The DRA provider is a narrow Slinky integration, not a general Kubernetes MNNVL topology provider. It groups nodes into Slurm `topology/block` domains from existing `nvidia.com/gpu.clique` labels. Those labels are not present in every GPU Operator deployment, and DRA does not discover the backend switch fabric between NVLink partitions. Consequently, if a workload cannot fit in one partition, selection of additional partitions is not informed by network proximity. Use NetQ or `infiniband-k8s` when cross-partition fabric locality is required. +The DRA provider is a narrow Slinky integration, not a general Kubernetes MNNVL topology provider. It groups nodes into Slurm `topology/block` domains from an existing Node label, defaulting to `nvidia.com/gpu.clique`. That label is not present in every GPU Operator deployment, and DRA does not discover the backend switch fabric between NVLink partitions. Consequently, if a workload cannot fit in one partition, selection of additional partitions is not informed by network proximity. Use NetQ or `infiniband-k8s` when cross-partition fabric locality is required. For non-MNNVL GPU clusters (such as DGX B200 or B300 SuperPODs), `nvidia.com/gpu.clique` is not set — Topograph with an InfiniBand provider is the only source of network topology for scheduling decisions on these systems. diff --git a/docs/providers/dra.md b/docs/providers/dra.md index 90288e97..9d595ec1 100644 --- a/docs/providers/dra.md +++ b/docs/providers/dra.md @@ -1,7 +1,8 @@ # DRA Topology Provider -The DRA provider reads existing `nvidia.com/gpu.clique` Kubernetes node labels -generated by the [NVIDIA GPU Operator](https://docs.nvidia.com/datacenter/cloud-native/gpu-operator/latest/index.html)'s +The DRA provider reads accelerator-domain IDs from an existing Kubernetes Node +label and defaults to `nvidia.com/gpu.clique`, generated by the +[NVIDIA GPU Operator](https://docs.nvidia.com/datacenter/cloud-native/gpu-operator/latest/index.html)'s [GPU Feature Discovery (GFD)](https://github.com/NVIDIA/k8s-device-plugin/blob/main/docs/gpu-feature-discovery/README.md), specifically its IMEX labeler, and groups nodes by NVLink partition. Its supported use is generating Slurm `topology/block` data with the Slinky engine. @@ -25,19 +26,19 @@ The DRA provider supplies the label-to-Slinky block-topology bridge for this eco Use the DRA provider only when all of the following are true: - You are using Slinky (Slurm-on-Kubernetes) with `topology/block` -- Every participating node already has a valid `nvidia.com/gpu.clique` label +- Every participating node already has a valid value in the configured accelerator-domain label - Workloads fit within one NVLink partition, or you accept that placement across multiple partitions will not account for backend-fabric locality The label is deployment- and state-dependent; it is not guaranteed to exist in every MNNVL GPU Operator installation. If it is absent, the provider cannot derive partition membership. If you need the switch hierarchy or topology-aware selection across partitions, use the [InfiniBand](./infiniband.md) or [NetQ](./netq.md) provider instead. ## How It Works -The DRA provider does not create `nvidia.com/gpu.clique`. Before selecting this provider, verify that the GPU Operator exposes the label on every participating node. +The DRA provider does not create its source label. Before selecting this provider, verify that the configured label exists on every participating node. When `provider.params.accelerator` is omitted, the source label defaults to `nvidia.com/gpu.clique` for backward compatibility. Topograph reads these labels from the Kubernetes API: 1. Lists all nodes (filtered by `nodeSelector` if provided) -2. For each node with a `nvidia.com/gpu.clique` label, reads the clique ID and groups nodes by domain +2. Uses the shared Kubernetes-label accelerator discoverer to read the configured Node label and group nodes by domain 3. Returns the NVLink domain map as block topology If no nodes with matching labels are found, Topograph returns a `502` error with a diagnostic message indicating which label and annotations to check. @@ -45,21 +46,40 @@ If no nodes with matching labels are found, Topograph returns a `502` error with ## Prerequisites - A Slinky (Slurm-on-Kubernetes) cluster configured to use `topology/block` -- A valid `nvidia.com/gpu.clique` label already present on every participating Kubernetes node +- A valid accelerator-domain label already present on every participating Kubernetes node ## Parameters | Parameter | Type | Required | Description | |---|---|---|---| | `nodeSelector` | `map[string]string` | No | Label selector to filter which nodes participate in topology discovery | +| `accelerator` | `object` | No | Shared accelerator discovery configuration. When omitted, DRA reads `nvidia.com/gpu.clique`. | +| `accelerator.source` | `string` | With `accelerator` | Must be `kubernetes-label`; DRA is an accelerator-only label provider. | +| `accelerator.kubernetesLabel.key` | `string` | With `accelerator` | Kubernetes Node label read as the accelerator-domain ID. | ## Configuration Deploy Topograph with the Helm chart and select the DRA provider and Slinky engine through the chart's `provider` and `engine` values. Set the optional -`nodeSelector` under `provider.params`. The chart manages the Topograph -configuration and topology request payload; they do not need to be supplied -separately. +source label and `nodeSelector` under `provider.params`. The chart manages the +Topograph configuration and topology request payload, so you do not need to +supply either one separately. + +```yaml +provider: + name: dra + params: + accelerator: + source: kubernetes-label + kubernetesLabel: + key: nvidia.com/gpu.clique + nodeSelector: + nvidia.com/gpu.present: "true" +``` + +The `accelerator` object has the same shape as the Kubernetes-label source for +`infiniband-k8s`. DRA supports only that source because it intentionally +produces accelerator block domains without discovering a network fabric. Configure Kubernetes client limits with the chart-wide settings: @@ -85,13 +105,13 @@ for the available Helm values. ## Verifying the Output -Before triggering topology generation, verify that clique labels exist on all participating nodes: +Before triggering topology generation, verify that the configured labels exist on all participating nodes. For the default label: ```bash kubectl get nodes -o json | jq '.items[] | {name: .metadata.name, clique: .metadata.labels["nvidia.com/gpu.clique"]}' ``` -If topology generation returns a `502` error, check that the expected nodes have the `nvidia.com/gpu.clique` label and the `topograph.nvidia.com/region` / `topograph.nvidia.com/instance` annotations (the latter two are set by Topograph itself during topology discovery): +If topology generation returns a `502` error, check that the expected nodes have the configured source label and the `topograph.nvidia.com/region` / `topograph.nvidia.com/instance` annotations (the latter two are set by Topograph itself during topology discovery). For the default label: ```bash kubectl get nodes -o json | jq '.items[] | {name: .metadata.name, clique: .metadata.labels["nvidia.com/gpu.clique"], region: .metadata.annotations["topograph.nvidia.com/region"], instance: .metadata.annotations["topograph.nvidia.com/instance"]}' diff --git a/docs/reference/node-labels.md b/docs/reference/node-labels.md index 5c2f9e66..a6032f54 100644 --- a/docs/reference/node-labels.md +++ b/docs/reference/node-labels.md @@ -62,7 +62,7 @@ provider discovery. Some GPU Operator deployments expose `nvidia.com/gpu.clique` on nodes with Multi-Node NVLink (MNNVL) GPUs; operators may select it explicitly as either an -engine source label or an `infiniband-k8s` provider discovery label. The `netq` +engine source label or a DRA/`infiniband-k8s` provider discovery label. The `netq` provider instead uses a `DomainUUID` from the NMX management API—a different identifier that refers to the same physical domain but cannot be compared as a string. diff --git a/internal/k8s/utils.go b/internal/k8s/utils.go index 030b156b..038267eb 100644 --- a/internal/k8s/utils.go +++ b/internal/k8s/utils.go @@ -21,9 +21,28 @@ import ( "k8s.io/client-go/tools/remotecommand" "k8s.io/klog/v2" + internalconfig "github.com/NVIDIA/topograph/internal/config" "github.com/NVIDIA/topograph/pkg/topology" ) +type nodeSelectorConfig struct { + NodeSelector map[string]string `mapstructure:"nodeSelector"` +} + +// NodeListOptions decodes a provider's optional nodeSelector into Kubernetes +// list options. Other provider parameters are intentionally ignored. +func NodeListOptions(params map[string]any) (*metav1.ListOptions, error) { + config := nodeSelectorConfig{} + if err := internalconfig.Decode(params, &config); err != nil { + return nil, err + } + if len(config.NodeSelector) == 0 { + return nil, nil + } + + return &metav1.ListOptions{LabelSelector: labels.Set(config.NodeSelector).String()}, nil +} + func GetNodes(ctx context.Context, client kubernetes.Interface, opt *metav1.ListOptions) (*corev1.NodeList, error) { if opt == nil { opt = &metav1.ListOptions{} diff --git a/internal/k8s/utils_test.go b/internal/k8s/utils_test.go index 14994fba..fa394706 100644 --- a/internal/k8s/utils_test.go +++ b/internal/k8s/utils_test.go @@ -10,6 +10,7 @@ import ( "github.com/stretchr/testify/require" corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" ) func TestIsPodReady(t *testing.T) { @@ -81,6 +82,43 @@ func TestIsPodReady(t *testing.T) { } } +func TestNodeListOptions(t *testing.T) { + tests := []struct { + name string + params map[string]any + want *metav1.ListOptions + err string + }{ + {name: "no parameters"}, + { + name: "other provider parameters are ignored", + params: map[string]any{"accelerator": map[string]any{"source": "none"}}, + }, + { + name: "node selector", + params: map[string]any{"nodeSelector": map[string]string{"key": "value"}}, + want: &metav1.ListOptions{LabelSelector: "key=value"}, + }, + { + name: "invalid node selector", + params: map[string]any{"nodeSelector": 0.1}, + err: "could not decode configuration", + }, + } + + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + got, err := NodeListOptions(test.params) + if test.err != "" { + require.ErrorContains(t, err, test.err) + return + } + require.NoError(t, err) + require.Equal(t, test.want, got) + }) + } +} + func TestValidateLabelKey(t *testing.T) { testCases := []struct { name string diff --git a/pkg/accelerator/accelerator.go b/pkg/accelerator/accelerator.go index e981edfe..93dac450 100644 --- a/pkg/accelerator/accelerator.go +++ b/pkg/accelerator/accelerator.go @@ -50,6 +50,25 @@ type Section struct { present bool } +// KubernetesLabelSection returns an accelerator section configured to read +// domain IDs from the supplied Kubernetes Node label. +func KubernetesLabelSection(key string) Section { + return Section{ + value: map[string]any{ + "source": SourceKubernetesLabel, + "kubernetesLabel": map[string]any{ + "key": key, + }, + }, + present: true, + } +} + +// Present reports whether the accelerator section was explicitly supplied. +func (s Section) Present() bool { + return s.present +} + // SectionFromProviderParams extracts the accelerator section without parsing // source-specific fields in the provider. func SectionFromProviderParams(providerParams map[string]any) Section { @@ -162,6 +181,37 @@ type Discoverer interface { Discover(context.Context, []Target) (Assignments, error) } +// TargetsFromComputeInstances converts the canonical request identity mapping +// into accelerator discovery targets. +func TargetsFromComputeInstances(instances []topology.ComputeInstances) []Target { + targets := make([]Target, 0) + for _, regionalInstances := range instances { + for instanceID, hostName := range regionalInstances.Instances { + targets = append(targets, Target{InstanceID: instanceID, HostName: hostName}) + } + } + return targets +} + +// DomainMapFromAssignments converts discovered accelerator assignments into +// the canonical topology domain map using the target identity mapping. +func DomainMapFromAssignments(assignments Assignments, targets []Target) topology.DomainMap { + domainMap := topology.NewDomainMap() + for _, target := range targets { + assignment, ok := assignments[target.InstanceID] + if !ok { + continue + } + domainMap.AddHostInfo(&topology.HostInfo{ + Domain: assignment.DomainID, + SubDomain: assignment.SubDomainID, + InstanceID: target.InstanceID, + HostName: target.HostName, + }) + } + return domainMap +} + type metadataDiscoverer struct { key string value func(Target, string) string @@ -174,7 +224,15 @@ func NewKubernetesDiscoverer(section Section) (Discoverer, error) { if err != nil { return nil, err } + return NewKubernetesDiscovererFromConfig(config) +} +// NewKubernetesDiscovererFromConfig returns a discoverer that resolves +// accelerator domains from Kubernetes Node metadata using validated config. +func NewKubernetesDiscovererFromConfig(config Config) (Discoverer, error) { + if err := config.Validate(); err != nil { + return nil, err + } switch config.Source { case SourceNvidiaSMI: return &metadataDiscoverer{ diff --git a/pkg/accelerator/accelerator_test.go b/pkg/accelerator/accelerator_test.go index 568aa9c4..4456d554 100644 --- a/pkg/accelerator/accelerator_test.go +++ b/pkg/accelerator/accelerator_test.go @@ -11,6 +11,8 @@ import ( "testing" "github.com/stretchr/testify/require" + corev1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" kubernetesfake "k8s.io/client-go/kubernetes/fake" "k8s.io/client-go/rest" @@ -208,6 +210,68 @@ func TestKubernetesDiscoverer(t *testing.T) { } } +func TestDiscoverKubernetesDomainsUsesCanonicalIdentity(t *testing.T) { + nodes := &corev1.NodeList{Items: []corev1.Node{ + {ObjectMeta: metav1.ObjectMeta{ + Name: "kubernetes-node-1", + Labels: map[string]string{testKubernetesLabel: "domain-1"}, + Annotations: map[string]string{ + topology.KeyNodeInstance: "instance-1", + topology.KeyNodeRegion: "region-1", + }, + }}, + {ObjectMeta: metav1.ObjectMeta{Name: "missing-annotations"}}, + {ObjectMeta: metav1.ObjectMeta{ + Name: "unrequested-node", + Labels: map[string]string{testKubernetesLabel: "domain-2"}, + Annotations: map[string]string{ + topology.KeyNodeInstance: "instance-2", + topology.KeyNodeRegion: "region-2", + }, + }}, + }} + instances := []topology.ComputeInstances{{ + Region: "region-1", + Instances: map[string]string{"instance-1": "scheduler-node-1"}, + }} + discoverer, err := NewKubernetesDiscoverer(KubernetesLabelSection(testKubernetesLabel)) + require.NoError(t, err) + + domains, err := DiscoverKubernetesDomains(context.Background(), discoverer, nodes, instances) + require.NoError(t, err) + expected := topology.NewDomainMap() + expected.AddHost("domain-1", "instance-1", "scheduler-node-1") + require.Equal(t, expected, domains) +} + +func TestTargetsAndDomainMapFromComputeInstances(t *testing.T) { + targets := TargetsFromComputeInstances([]topology.ComputeInstances{{ + Region: "region-1", + Instances: map[string]string{ + "instance-1": "node-1", + "instance-2": "node-2", + }, + }}) + require.ElementsMatch(t, []Target{ + {InstanceID: "instance-1", HostName: "node-1"}, + {InstanceID: "instance-2", HostName: "node-2"}, + }, targets) + + domains := DomainMapFromAssignments(Assignments{ + "instance-1": {DomainID: "domain-1", SubDomainID: "sub-domain-1"}, + }, targets) + require.Equal(t, topology.DomainMap{ + "domain-1": { + "node-1": { + Domain: "domain-1", + SubDomain: "sub-domain-1", + InstanceID: "instance-1", + HostName: "node-1", + }, + }, + }, domains) +} + type fakeCommandRunner struct { outputs map[string]string err error diff --git a/pkg/accelerator/kubernetes.go b/pkg/accelerator/kubernetes.go index 955c466b..a07f0441 100644 --- a/pkg/accelerator/kubernetes.go +++ b/pkg/accelerator/kubernetes.go @@ -10,13 +10,79 @@ import ( "fmt" "strings" + corev1 "k8s.io/api/core/v1" "k8s.io/client-go/kubernetes" "k8s.io/client-go/rest" "k8s.io/klog/v2" internalK8s "github.com/NVIDIA/topograph/internal/k8s" + "github.com/NVIDIA/topograph/pkg/topology" ) +// BaseKubernetesNodeAnnotations returns the identity annotations shared by +// Kubernetes providers whose instance ID and region are node-local. +func BaseKubernetesNodeAnnotations(hostName string) map[string]string { + return map[string]string{ + topology.KeyNodeInstance: hostName, + topology.KeyNodeRegion: "local", + } +} + +// TargetsFromKubernetesNodes resolves selected Nodes through the canonical +// region/instance mapping and attaches their metadata for discovery. +func TargetsFromKubernetesNodes(nodes *corev1.NodeList, instances []topology.ComputeInstances) []Target { + instancesByRegion := make(map[string]map[string]string, len(instances)) + for _, regionalInstances := range instances { + instancesByRegion[regionalInstances.Region] = regionalInstances.Instances + } + + targets := make([]Target, 0, len(nodes.Items)) + for _, node := range nodes.Items { + instanceID := strings.TrimSpace(node.Annotations[topology.KeyNodeInstance]) + if instanceID == "" { + klog.Warningf("missing or empty %q annotation in node %s", topology.KeyNodeInstance, node.Name) + continue + } + region := strings.TrimSpace(node.Annotations[topology.KeyNodeRegion]) + if region == "" { + klog.Warningf("missing or empty %q annotation in node %s", topology.KeyNodeRegion, node.Name) + continue + } + regionalInstances, ok := instancesByRegion[region] + if !ok { + continue + } + hostName, ok := regionalInstances[instanceID] + if !ok { + continue + } + + targets = append(targets, Target{ + InstanceID: instanceID, + HostName: hostName, + Labels: node.Labels, + Annotations: node.Annotations, + }) + } + return targets +} + +// DiscoverKubernetesDomains runs a configured metadata discoverer against +// selected Nodes and converts its result into the canonical domain map. +func DiscoverKubernetesDomains( + ctx context.Context, + discoverer Discoverer, + nodes *corev1.NodeList, + instances []topology.ComputeInstances, +) (topology.DomainMap, error) { + targets := TargetsFromKubernetesNodes(nodes, instances) + assignments, err := discoverer.Discover(ctx, targets) + if err != nil { + return nil, err + } + return DomainMapFromAssignments(assignments, targets), nil +} + // NewKubernetesNodeDiscoverer returns the node-local discoverer used by the // node-data-broker. Sources that read existing Kubernetes metadata require no // node-local collection and therefore return an empty discoverer. diff --git a/pkg/providers/dra/provider.go b/pkg/providers/dra/provider.go index 21799de2..7f673e15 100644 --- a/pkg/providers/dra/provider.go +++ b/pkg/providers/dra/provider.go @@ -11,14 +11,12 @@ import ( "net/http" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/labels" "k8s.io/client-go/kubernetes" "k8s.io/client-go/rest" - "k8s.io/klog/v2" - "github.com/NVIDIA/topograph/internal/config" "github.com/NVIDIA/topograph/internal/httperr" "github.com/NVIDIA/topograph/internal/k8s" + "github.com/NVIDIA/topograph/pkg/accelerator" "github.com/NVIDIA/topograph/pkg/providers" "github.com/NVIDIA/topograph/pkg/topology" ) @@ -26,21 +24,15 @@ import ( const ( NAME = "dra" - DomainLabel = topology.KeyNvidiaGPUClique + defaultDomainLabel = "nvidia.com/gpu.clique" ) type Provider struct { - config *rest.Config - client kubernetes.Interface - params *Params -} - -type Params struct { - // NodeSelector (optional) specifies nodes participating in the topology - NodeSelector map[string]string `mapstructure:"nodeSelector"` - - // derived fields - nodeListOpt *metav1.ListOptions + config *rest.Config + client kubernetes.Interface + nodeListOpt *metav1.ListOptions + accelerator accelerator.Discoverer + acceleratorLabel string } func NamedLoader() (string, providers.Loader) { @@ -48,7 +40,11 @@ func NamedLoader() (string, providers.Loader) { } func Loader(ctx context.Context, config providers.Config) (providers.Provider, *httperr.Error) { - p, err := getParameters(config.Params) + nodeListOpt, err := k8s.NodeListOptions(config.Params) + if err != nil { + return nil, httperr.NewError(http.StatusBadRequest, err.Error()) + } + acceleratorDiscoverer, acceleratorLabel, err := newAcceleratorDiscoverer(config.Params) if err != nil { return nil, httperr.NewError(http.StatusBadRequest, err.Error()) } @@ -68,66 +64,48 @@ func Loader(ctx context.Context, config providers.Config) (providers.Provider, * } return &Provider{ - config: cfg, - client: client, - params: p, + config: cfg, + client: client, + nodeListOpt: nodeListOpt, + accelerator: acceleratorDiscoverer, + acceleratorLabel: acceleratorLabel, }, nil } -func getParameters(params map[string]any) (*Params, error) { - p := &Params{} - if err := config.Decode(params, p); err != nil { - return nil, err +func newAcceleratorDiscoverer(params map[string]any) (accelerator.Discoverer, string, error) { + section := accelerator.SectionFromProviderParams(params) + if !section.Present() { + section = accelerator.KubernetesLabelSection(defaultDomainLabel) } - if len(p.NodeSelector) != 0 { - p.nodeListOpt = &metav1.ListOptions{ - LabelSelector: labels.Set(p.NodeSelector).String(), - } + acceleratorConfig, err := accelerator.ParseConfig(section) + if err != nil { + return nil, "", err } - - return p, nil + if acceleratorConfig.Source != accelerator.SourceKubernetesLabel { + return nil, "", fmt.Errorf("dra provider supports only accelerator source %q", accelerator.SourceKubernetesLabel) + } + acceleratorDiscoverer, err := accelerator.NewKubernetesDiscovererFromConfig(acceleratorConfig) + if err != nil { + return nil, "", err + } + return acceleratorDiscoverer, acceleratorConfig.KubernetesLabel.Key, nil } func (p *Provider) GenerateTopologyConfig(ctx context.Context, _ *int, instances []topology.ComputeInstances) (*topology.Graph, *httperr.Error) { - regIndices := make(map[string]int) // map[region : index] - for i, ci := range instances { - regIndices[ci.Region] = i - } - - nodes, err := k8s.GetNodes(ctx, p.client, p.params.nodeListOpt) + nodes, err := k8s.GetNodes(ctx, p.client, p.nodeListOpt) if err != nil { return nil, httperr.NewError(http.StatusBadGateway, err.Error()) } - domainMap := topology.NewDomainMap() - for _, node := range nodes.Items { - clusterID, ok := node.Labels[DomainLabel] - if !ok { - continue - } - - region := node.Annotations[topology.KeyNodeRegion] - indx, ok := regIndices[region] - if !ok { - continue - } - - i2n := instances[indx].Instances - instanceID, ok := node.Annotations[topology.KeyNodeInstance] - if !ok || instanceID == "" { - klog.Warningf("missing or empty %q annotation in node %s", topology.KeyNodeInstance, node.Name) - continue - } - - if host, ok := i2n[instanceID]; ok { - domainMap.AddHost(clusterID, instanceID, host) - } + domainMap, err := accelerator.DiscoverKubernetesDomains(ctx, p.accelerator, nodes, instances) + if err != nil { + return nil, httperr.NewError(http.StatusBadGateway, fmt.Sprintf("failed to discover accelerator domains: %v", err)) } if len(domainMap) == 0 { return nil, httperr.NewError(http.StatusBadGateway, fmt.Sprintf("no matching nodes found; check label %q and annotations %q and %q", - DomainLabel, topology.KeyNodeRegion, topology.KeyNodeInstance)) + p.acceleratorLabel, topology.KeyNodeRegion, topology.KeyNodeInstance)) } return &topology.Graph{ @@ -135,11 +113,6 @@ func (p *Provider) GenerateTopologyConfig(ctx context.Context, _ *int, instances }, nil } -func GetNodeAnnotations(ctx context.Context, hostname string) (map[string]string, error) { - annotations := map[string]string{ - topology.KeyNodeInstance: hostname, - topology.KeyNodeRegion: "local", - } - - return annotations, nil +func GetNodeAnnotations(_ context.Context, hostname string) (map[string]string, error) { + return accelerator.BaseKubernetesNodeAnnotations(hostname), nil } diff --git a/pkg/providers/dra/provider_test.go b/pkg/providers/dra/provider_test.go index 2ecaa5d2..3814aec1 100644 --- a/pkg/providers/dra/provider_test.go +++ b/pkg/providers/dra/provider_test.go @@ -9,6 +9,7 @@ import ( "context" "testing" + "github.com/NVIDIA/topograph/pkg/accelerator" "github.com/NVIDIA/topograph/pkg/topology" "github.com/stretchr/testify/require" corev1 "k8s.io/api/core/v1" @@ -16,44 +17,46 @@ import ( "k8s.io/client-go/kubernetes/fake" ) -func TestGetParameters(t *testing.T) { - testCases := []struct { - name string - params map[string]any - ret *Params - err string +func TestNewAcceleratorDiscoverer(t *testing.T) { + tests := []struct { + name string + params map[string]any + labelKey string + err string }{ + {name: "legacy default", labelKey: defaultDomainLabel}, { - name: "Case 1: no params", - params: nil, - ret: &Params{}, + name: "configured label", + params: map[string]any{"accelerator": map[string]any{ + "source": accelerator.SourceKubernetesLabel, + "kubernetesLabel": map[string]any{"key": "example.com/domain"}, + }}, + labelKey: "example.com/domain", }, { - name: "Case 2: bad params", - params: map[string]any{"nodeSelector": .1}, - err: "could not decode configuration: 1 error(s) decoding:\n\n* 'nodeSelector' expected a map, got 'float64'", - }, - { - name: "Case 3: valid input", - params: map[string]any{"nodeSelector": map[string]string{"key": "val"}}, - ret: &Params{ - NodeSelector: map[string]string{"key": "val"}, - nodeListOpt: &metav1.ListOptions{ - LabelSelector: "key=val", - }, - }, + name: "unsupported source", + params: map[string]any{"accelerator": map[string]any{"source": accelerator.SourceNone}}, + err: `dra provider supports only accelerator source "kubernetes-label"`, }, } - for _, tc := range testCases { - t.Run(tc.name, func(t *testing.T) { - p, err := getParameters(tc.params) - if len(tc.err) != 0 { - require.ErrorContains(t, err, tc.err) - } else { - require.NoError(t, err) - require.Equal(t, tc.ret, p) + for _, test := range tests { + t.Run(test.name, func(t *testing.T) { + discoverer, labelKey, err := newAcceleratorDiscoverer(test.params) + if test.err != "" { + require.EqualError(t, err, test.err) + return } + require.NoError(t, err) + require.Equal(t, test.labelKey, labelKey) + assignments, err := discoverer.Discover(context.Background(), []accelerator.Target{{ + InstanceID: "instance-1", + Labels: map[string]string{test.labelKey: "domain-1"}, + }}) + require.NoError(t, err) + require.Equal(t, accelerator.Assignments{ + "instance-1": {DomainID: "domain-1"}, + }, assignments) }) } } @@ -62,7 +65,7 @@ func TestGenerateTopologyConfigUsesAnnotatedInstanceID(t *testing.T) { node := &corev1.Node{ObjectMeta: metav1.ObjectMeta{ Name: "k8s-node-1", Labels: map[string]string{ - DomainLabel: "clique-1", + defaultDomainLabel: "clique-1", }, Annotations: map[string]string{ topology.KeyNodeInstance: "instance-123", @@ -70,8 +73,9 @@ func TestGenerateTopologyConfigUsesAnnotatedInstanceID(t *testing.T) { }, }} provider := &Provider{ - client: fake.NewSimpleClientset(node), - params: &Params{}, + client: fake.NewSimpleClientset(node), + accelerator: mustAcceleratorDiscoverer(t, defaultDomainLabel), + acceleratorLabel: defaultDomainLabel, } instances := []topology.ComputeInstances{{ Region: "local", @@ -87,3 +91,10 @@ func TestGenerateTopologyConfigUsesAnnotatedInstanceID(t *testing.T) { expectedDomains.AddHost("clique-1", "instance-123", "scheduler-node-1") require.Equal(t, expectedDomains, graph.Domains) } + +func mustAcceleratorDiscoverer(t *testing.T, labelKey string) accelerator.Discoverer { + t.Helper() + discoverer, err := accelerator.NewKubernetesDiscoverer(accelerator.KubernetesLabelSection(labelKey)) + require.NoError(t, err) + return discoverer +} diff --git a/pkg/providers/infiniband/bm.go b/pkg/providers/infiniband/bm.go index 19dfb24d..35655b19 100644 --- a/pkg/providers/infiniband/bm.go +++ b/pkg/providers/infiniband/bm.go @@ -13,7 +13,6 @@ import ( "github.com/NVIDIA/topograph/internal/exec" "github.com/NVIDIA/topograph/pkg/accelerator" - "github.com/NVIDIA/topograph/pkg/topology" ) type IBNetDiscoverBM struct{} @@ -55,30 +54,3 @@ func parsePdshNvidiaSMIOutput(stdout *bytes.Buffer) (map[string]string, error) { return outputs, nil } - -func acceleratorTargets(cis []topology.ComputeInstances) []accelerator.Target { - targets := make([]accelerator.Target, 0) - for _, ci := range cis { - for instanceID, hostName := range ci.Instances { - targets = append(targets, accelerator.Target{InstanceID: instanceID, HostName: hostName}) - } - } - return targets -} - -func domainMapFromAssignments(assignments accelerator.Assignments, targets []accelerator.Target) topology.DomainMap { - domainMap := topology.NewDomainMap() - for _, target := range targets { - assignment, ok := assignments[target.InstanceID] - if !ok { - continue - } - domainMap.AddHostInfo(&topology.HostInfo{ - Domain: assignment.DomainID, - SubDomain: assignment.SubDomainID, - InstanceID: target.InstanceID, - HostName: target.HostName, - }) - } - return domainMap -} diff --git a/pkg/providers/infiniband/bm_test.go b/pkg/providers/infiniband/bm_test.go index eafc8402..89a3e851 100644 --- a/pkg/providers/infiniband/bm_test.go +++ b/pkg/providers/infiniband/bm_test.go @@ -10,9 +10,6 @@ import ( "testing" "github.com/stretchr/testify/require" - - "github.com/NVIDIA/topograph/pkg/accelerator" - "github.com/NVIDIA/topograph/pkg/topology" ) func TestParsePdshNvidiaSMIOutput(t *testing.T) { @@ -29,30 +26,3 @@ node-2: uuid-2, 8 "node-2": "uuid-2, 8\n", }, outputs) } - -func TestAcceleratorTargetsAndDomainMap(t *testing.T) { - targets := acceleratorTargets([]topology.ComputeInstances{{ - Instances: map[string]string{ - "instance-1": "node-1", - "instance-2": "node-2", - }, - }}) - require.ElementsMatch(t, []accelerator.Target{ - {InstanceID: "instance-1", HostName: "node-1"}, - {InstanceID: "instance-2", HostName: "node-2"}, - }, targets) - - domainMap := domainMapFromAssignments(accelerator.Assignments{ - "instance-1": {DomainID: "domain-1", SubDomainID: "partition-1"}, - }, targets) - require.Equal(t, topology.DomainMap{ - "domain-1": { - "node-1": { - Domain: "domain-1", - SubDomain: "partition-1", - InstanceID: "instance-1", - HostName: "node-1", - }, - }, - }, domainMap) -} diff --git a/pkg/providers/infiniband/k8s.go b/pkg/providers/infiniband/k8s.go index 12cd6d9e..8d06a3d8 100644 --- a/pkg/providers/infiniband/k8s.go +++ b/pkg/providers/infiniband/k8s.go @@ -48,10 +48,7 @@ func (h *IBNetDiscoverK8S) Run(ctx context.Context, node string) (*bytes.Buffer, } func GetNodeAnnotations(ctx context.Context, client kubernetes.Interface, config *rest.Config, hostname string, section accelerator.Section) (map[string]string, error) { - annotations := map[string]string{ - topology.KeyNodeInstance: hostname, - topology.KeyNodeRegion: "local", - } + annotations := accelerator.BaseKubernetesNodeAnnotations(hostname) discoverer, err := accelerator.NewKubernetesNodeDiscoverer(section, client, config) if err != nil { diff --git a/pkg/providers/infiniband/provider_bm.go b/pkg/providers/infiniband/provider_bm.go index dd0b1087..eeb2a400 100644 --- a/pkg/providers/infiniband/provider_bm.go +++ b/pkg/providers/infiniband/provider_bm.go @@ -43,12 +43,12 @@ func (p *ProviderBM) GenerateTopologyConfig(ctx context.Context, _ *int, cis []t return nil, httperr.NewError(http.StatusBadRequest, "on-prem does not support multi-region topology requests") } - targets := acceleratorTargets(cis) + targets := accelerator.TargetsFromComputeInstances(cis) assignments, err := p.accelerator.Discover(ctx, targets) if err != nil { return nil, httperr.NewError(http.StatusInternalServerError, fmt.Sprintf("failed to discover accelerator domains: %v", err)) } - domainMap := domainMapFromAssignments(assignments, targets) + domainMap := accelerator.DomainMapFromAssignments(assignments, targets) treeRoot, err := getIbTree(ctx, cis, &IBNetDiscoverBM{}) if err != nil { diff --git a/pkg/providers/infiniband/provider_k8s.go b/pkg/providers/infiniband/provider_k8s.go index 4ed3844a..66f04be6 100644 --- a/pkg/providers/infiniband/provider_k8s.go +++ b/pkg/providers/infiniband/provider_k8s.go @@ -11,11 +11,9 @@ import ( "net/http" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/labels" "k8s.io/client-go/kubernetes" "k8s.io/client-go/rest" - "github.com/NVIDIA/topograph/internal/config" "github.com/NVIDIA/topograph/internal/httperr" "github.com/NVIDIA/topograph/internal/k8s" "github.com/NVIDIA/topograph/pkg/accelerator" @@ -28,16 +26,8 @@ const NAME_K8S = "infiniband-k8s" type ProviderK8S struct { config *rest.Config client *kubernetes.Clientset - params *Params - accelerator accelerator.Discoverer -} - -type Params struct { - // NodeSelector (optional) specifies nodes participating in the topology - NodeSelector map[string]string `mapstructure:"nodeSelector"` - - // derived fields nodeListOpt *metav1.ListOptions + accelerator accelerator.Discoverer } func NamedLoaderK8S() (string, providers.Loader) { @@ -45,7 +35,7 @@ func NamedLoaderK8S() (string, providers.Loader) { } func LoaderK8S(ctx context.Context, config providers.Config) (providers.Provider, *httperr.Error) { - p, err := getParameters(config.Params) + nodeListOpt, err := k8s.NodeListOptions(config.Params) if err != nil { return nil, httperr.NewError(http.StatusBadRequest, err.Error()) } @@ -69,50 +59,25 @@ func LoaderK8S(ctx context.Context, config providers.Config) (providers.Provider return &ProviderK8S{ config: cfg, client: client, - params: p, + nodeListOpt: nodeListOpt, accelerator: acceleratorDiscoverer, }, nil } -func getParameters(params map[string]any) (*Params, error) { - p := &Params{} - if err := config.Decode(params, p); err != nil { - return nil, err - } - - if len(p.NodeSelector) != 0 { - p.nodeListOpt = &metav1.ListOptions{ - LabelSelector: labels.Set(p.NodeSelector).String(), - } - } - - return p, nil -} - func (p *ProviderK8S) GenerateTopologyConfig(ctx context.Context, _ *int, cis []topology.ComputeInstances) (*topology.Graph, *httperr.Error) { if len(cis) > 1 { return nil, httperr.NewError(http.StatusBadRequest, "on-prem does not support multi-region topology requests") } - nodes, err := k8s.GetNodes(ctx, p.client, p.params.nodeListOpt) + nodes, err := k8s.GetNodes(ctx, p.client, p.nodeListOpt) if err != nil { return nil, httperr.NewError(http.StatusBadGateway, err.Error()) } - targets := make([]accelerator.Target, 0, len(nodes.Items)) - for _, node := range nodes.Items { - targets = append(targets, accelerator.Target{ - InstanceID: node.Name, - HostName: node.Name, - Labels: node.Labels, - Annotations: node.Annotations, - }) - } - assignments, err := p.accelerator.Discover(ctx, targets) + domainMap, err := accelerator.DiscoverKubernetesDomains(ctx, p.accelerator, nodes, cis) if err != nil { return nil, httperr.NewError(http.StatusBadGateway, fmt.Sprintf("failed to discover accelerator domains: %v", err)) } - domainMap := domainMapFromAssignments(assignments, targets) ibnetdiscover := NewIBNetDiscoverK8S(p.config, p.client) treeRoot, err := getIbTree(ctx, cis, ibnetdiscover) diff --git a/pkg/providers/infiniband/provider_k8s_test.go b/pkg/providers/infiniband/provider_k8s_test.go index d49c67a8..e8a19837 100644 --- a/pkg/providers/infiniband/provider_k8s_test.go +++ b/pkg/providers/infiniband/provider_k8s_test.go @@ -6,45 +6,21 @@ package infiniband import ( + "context" "testing" "github.com/stretchr/testify/require" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + + "github.com/NVIDIA/topograph/pkg/topology" ) -func TestGetParameters(t *testing.T) { - tests := []struct { - name string - params map[string]any - labelSelector string - err string - }{ - {name: "no parameters"}, - { - name: "bad node selector", - params: map[string]any{"nodeSelector": .1}, - err: "could not decode configuration", - }, - { - name: "node selector", - params: map[string]any{"nodeSelector": map[string]string{"key": "val"}}, - labelSelector: "key=val", - }, - } +func TestProviderK8SRejectsMultiRegionRequest(t *testing.T) { + provider := &ProviderK8S{} + graph, httpErr := provider.GenerateTopologyConfig(context.Background(), nil, []topology.ComputeInstances{ + {Region: "region-1"}, + {Region: "region-2"}, + }) - for _, test := range tests { - t.Run(test.name, func(t *testing.T) { - params, err := getParameters(test.params) - if test.err != "" { - require.ErrorContains(t, err, test.err) - return - } - require.NoError(t, err) - if test.labelSelector == "" { - require.Nil(t, params.nodeListOpt) - } else { - require.Equal(t, &metav1.ListOptions{LabelSelector: test.labelSelector}, params.nodeListOpt) - } - }) - } + require.Nil(t, graph) + require.EqualError(t, httpErr, "on-prem does not support multi-region topology requests") } diff --git a/pkg/topology/graph_test.go b/pkg/topology/graph_test.go index 3d1ee48a..b2ecb33d 100644 --- a/pkg/topology/graph_test.go +++ b/pkg/topology/graph_test.go @@ -234,21 +234,23 @@ func TestToInstanceOmitsXclrSubDomainWithoutDomainLabel(t *testing.T) { require.NotContains(t, instance.Labels, KeyTopologyXclrSubDomain) } -func TestToInstanceIgnoresUnrelatedGPUCliqueLabel(t *testing.T) { +func TestToInstancePreservesUnrelatedAcceleratorDomainLabel(t *testing.T) { + const acceleratorDomainLabel = "example.cpm/accelerator-domain" + inst := &InstanceTopology{ InstanceID: "i-001", XclrDomainID: "discovered-domain", XclrSubDomainID: "discovered-sub-domain", Instance: &Instance{ Labels: map[string]string{ - KeyNvidiaGPUClique: "gpu-clique", + acceleratorDomainLabel: "external-domain", }, }, } instance := inst.toInstance(0) - require.Equal(t, "gpu-clique", instance.Labels[KeyNvidiaGPUClique]) + require.Equal(t, "external-domain", instance.Labels[acceleratorDomainLabel]) require.Equal(t, "discovered-domain", instance.Labels[KeyTopologyXclrDomain]) require.Equal(t, "discovered-sub-domain", instance.Labels[KeyTopologyXclrSubDomain]) } diff --git a/pkg/topology/topology.go b/pkg/topology/topology.go index 3b81defe..ab5db3a3 100644 --- a/pkg/topology/topology.go +++ b/pkg/topology/topology.go @@ -37,7 +37,6 @@ const ( KeyGpuClusterID = "topograph.nvidia.com/cluster-id" // NVIDIA GPU Operator node labels - KeyNvidiaGPUClique = "nvidia.com/gpu.clique" KeyNvidiaGPUProduct = "nvidia.com/gpu.product" // Topograph default node labels. Fabric tier zero is closest to the compute diff --git a/tests/ci/values.slinky-dra-block.yaml b/tests/ci/values.slinky-dra-block.yaml index e49593eb..b4127fd8 100644 --- a/tests/ci/values.slinky-dra-block.yaml +++ b/tests/ci/values.slinky-dra-block.yaml @@ -3,6 +3,10 @@ provider: name: dra params: + accelerator: + source: kubernetes-label + kubernetesLabel: + key: nvidia.com/gpu.clique nodeSelector: kwok.x-k8s.io/node: "fake"