diff --git a/internal/pkg/collector/gpu_collector.go b/internal/pkg/collector/gpu_collector.go index c1fea6a4..f3b3fdff 100644 --- a/internal/pkg/collector/gpu_collector.go +++ b/internal/pkg/collector/gpu_collector.go @@ -260,6 +260,17 @@ func (c *DCGMCollector) latestValues(mi devicemonitoring.Info) ([]dcgm.FieldValu ) } + // Watch registration is scoped per GPU model (see the devicewatcher + // package), but the field list above isn't - it's every counter + // configured for this entity type, regardless of which model actually + // had it watched. For an ordinary field that's fine: DCGM reports one + // it doesn't support as a per-entity blank value. For a DCP field it + // isn't: asking for a value that was never watched returns + // DCGM_ST_NOT_WATCHED, and the stale-watch repair path then retries + // forever trying to fix a watch that was deliberately never created. + // Filter to what this entity's model actually has watched. + fields = filterDCPFieldsForModel(fields, mi.DeviceInfo) + return dcgmprovider.Client().EntityGetLatestValues( mi.Entity.EntityGroupId, mi.Entity.EntityId, @@ -267,6 +278,22 @@ func (c *DCGMCollector) latestValues(mi devicemonitoring.Info) ([]dcgm.FieldValu ) } +// filterDCPFieldsForModel drops DCP fields the entity's GPU model doesn't +// support, leaving non-DCP fields untouched. See NVIDIA/dcgm-exporter#736. +func filterDCPFieldsForModel(fields []dcgm.Short, device dcgm.Device) []dcgm.Short { + if device.Identifiers.Model == "" { + return fields + } + + filtered := make([]dcgm.Short, 0, len(fields)) + for _, fieldID := range fields { + if devicewatcher.ModelSupportsDCPField(device.Identifiers.Model, device.GPU, fieldID) { + filtered = append(filtered, fieldID) + } + } + return filtered +} + // addMetrics renders values with the labels and identity fields for their entity type. func (c *DCGMCollector) addMetrics( metrics MetricsByCounter, diff --git a/internal/pkg/counters/const.go b/internal/pkg/counters/const.go index 226ab17f..56b20a54 100644 --- a/internal/pkg/counters/const.go +++ b/internal/pkg/counters/const.go @@ -27,3 +27,12 @@ const ( DCGMExpXIDErrorsTotal = "DCGM_EXP_XID_ERRORS_TOTAL" DCGMExpClockEventsTotal = "DCGM_EXP_CLOCK_EVENTS_TOTAL" ) + +// IsDCPField reports whether fieldID falls in the DCP/profiling field ID range +// (DCGM_FI_PROF_*). These fields are validated for the whole DCGM watch group +// at registration time, unlike ordinary fields which degrade to a per-entity +// NOT_SUPPORTED value at scrape time - see devicewatcher's model-partitioned +// watch path for why that distinction matters. +func IsDCPField(fieldID uint) bool { + return fieldID >= dcpFieldsStart && fieldID < cpuFieldsStart +} diff --git a/internal/pkg/counters/counter_config.go b/internal/pkg/counters/counter_config.go index 3dd328e3..616d31d0 100644 --- a/internal/pkg/counters/counter_config.go +++ b/internal/pkg/counters/counter_config.go @@ -191,7 +191,7 @@ func ExtractCounters(records [][]string, c *appconfig.Config) (*CounterSet, erro } func fieldIsSupported(fieldID uint, c *appconfig.Config) bool { - if fieldID < dcpFieldsStart || fieldID >= cpuFieldsStart { + if !IsDCPField(fieldID) { return true } diff --git a/internal/pkg/devicewatcher/dcp_model_partition_test.go b/internal/pkg/devicewatcher/dcp_model_partition_test.go new file mode 100644 index 00000000..1a3bb569 --- /dev/null +++ b/internal/pkg/devicewatcher/dcp_model_partition_test.go @@ -0,0 +1,358 @@ +/* + * Copyright (c) 2024, NVIDIA CORPORATION. All rights reserved. + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package devicewatcher + +// Regression coverage for https://github.com/NVIDIA/dcgm-exporter/issues/657: +// a node mixing GPU models (e.g. A30 + RTX2000) failed to start entirely +// because one shared DCGM watch group covered every GPU, and DCP/profiling +// fields fail the whole group's registration if any member GPU doesn't +// support them. These tests cover the per-model partition that replaces +// that shared group for exactly this situation. + +import ( + "errors" + "testing" + + "github.com/NVIDIA/go-dcgm/pkg/dcgm" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "go.uber.org/mock/gomock" + + mockdcgm "github.com/NVIDIA/dcgm-exporter/internal/mocks/pkg/dcgmprovider" + mockdeviceinfo "github.com/NVIDIA/dcgm-exporter/internal/mocks/pkg/deviceinfo" + "github.com/NVIDIA/dcgm-exporter/internal/pkg/appconfig" + "github.com/NVIDIA/dcgm-exporter/internal/pkg/dcgmprovider" + "github.com/NVIDIA/dcgm-exporter/internal/pkg/deviceinfo" + "github.com/NVIDIA/dcgm-exporter/internal/pkg/devicemonitoring" +) + +const ( + // fbUsed is an ordinary field: DCGM reports it as a per-entity + // NOT_SUPPORTED value rather than failing the watch, so it should never + // be filtered out. + fbUsed = dcgm.DCGM_FI_DEV_FB_USED // 252, < dcpFieldsStart + + // fp64Active is the DCP field from the original bug report. It's the one + // that needs per-model filtering. + fp64Active = dcgm.DCGM_FI_PROF_FP64_UTIL_RATIO // 1006, DCP range +) + +func TestAnyDCPField(t *testing.T) { + tests := []struct { + name string + groups []FieldWatchGroup + want bool + }{ + {"no groups", nil, false}, + {"only ordinary fields", []FieldWatchGroup{{Fields: []dcgm.Short{fbUsed}}}, false}, + {"one DCP field among ordinary ones", []FieldWatchGroup{{Fields: []dcgm.Short{fbUsed, fp64Active}}}, true}, + {"DCP field in a later group", []FieldWatchGroup{{Fields: []dcgm.Short{fbUsed}}, {Fields: []dcgm.Short{fp64Active}}}, true}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + assert.Equal(t, tt.want, anyDCPField(tt.groups)) + }) + } +} + +func TestPartitionMonitoredEntitiesByModel(t *testing.T) { + a30 := devicemonitoring.Info{ + Entity: dcgm.GroupEntityPair{EntityGroupId: dcgm.FE_GPU, EntityId: 0}, + DeviceInfo: dcgm.Device{GPU: 0, Identifiers: dcgm.DeviceIdentifiers{Model: "NVIDIA A30"}}, + } + rtx1 := devicemonitoring.Info{ + Entity: dcgm.GroupEntityPair{EntityGroupId: dcgm.FE_GPU, EntityId: 1}, + DeviceInfo: dcgm.Device{GPU: 1, Identifiers: dcgm.DeviceIdentifiers{Model: "NVIDIA RTX2000"}}, + } + rtx2 := devicemonitoring.Info{ + Entity: dcgm.GroupEntityPair{EntityGroupId: dcgm.FE_GPU, EntityId: 2}, + DeviceInfo: dcgm.Device{GPU: 2, Identifiers: dcgm.DeviceIdentifiers{Model: "NVIDIA RTX2000"}}, + } + + models, byModel := partitionMonitoredEntitiesByModel([]devicemonitoring.Info{a30, rtx1, rtx2}) + + require.Equal(t, []string{"NVIDIA A30", "NVIDIA RTX2000"}, models, "models should be sorted for deterministic group creation order") + assert.Equal(t, []devicemonitoring.Info{a30}, byModel["NVIDIA A30"]) + assert.Equal(t, []devicemonitoring.Info{rtx1, rtx2}, byModel["NVIDIA RTX2000"]) +} + +func TestFilterFieldsForModel(t *testing.T) { + tests := []struct { + name string + fields []dcgm.Short + supported map[dcgm.Short]bool + want []dcgm.Short + }{ + { + name: "ordinary field always kept even with no DCP support", + fields: []dcgm.Short{fbUsed}, + supported: map[dcgm.Short]bool{}, + want: []dcgm.Short{fbUsed}, + }, + { + name: "DCP field dropped when model doesn't support it", + fields: []dcgm.Short{fbUsed, fp64Active}, + supported: map[dcgm.Short]bool{}, + want: []dcgm.Short{fbUsed}, + }, + { + name: "DCP field kept when model supports it", + fields: []dcgm.Short{fbUsed, fp64Active}, + supported: map[dcgm.Short]bool{fp64Active: true}, + want: []dcgm.Short{fbUsed, fp64Active}, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + assert.Equal(t, tt.want, filterFieldsForModel(tt.fields, tt.supported)) + }) + } +} + +func TestSupportedDCPFields(t *testing.T) { + realDCGM := dcgmprovider.Client() + defer dcgmprovider.SetClient(realDCGM) + + ctrl := gomock.NewController(t) + mockDCGM := mockdcgm.NewMockDCGM(ctrl) + dcgmprovider.SetClient(mockDCGM) + + mockDCGM.EXPECT().GetSupportedMetricGroups(uint(0)).Return([]dcgm.MetricGroup{ + {FieldIds: []uint{uint(fp64Active)}}, + }, nil) + + supported, err := supportedDCPFields(0) + require.NoError(t, err) + assert.Equal(t, map[dcgm.Short]bool{fp64Active: true}, supported) +} + +// TestModelSupportsDCPField_NonDCPFieldNeverQueriesDCGM covers the cheap +// path: a non-DCP field is always reported supported without touching DCGM +// or the cache at all. +func TestModelSupportsDCPField_NonDCPFieldNeverQueriesDCGM(t *testing.T) { + realDCGM := dcgmprovider.Client() + defer dcgmprovider.SetClient(realDCGM) + defer ResetDCPCapabilityCache() + + ctrl := gomock.NewController(t) + mockDCGM := mockdcgm.NewMockDCGM(ctrl) + dcgmprovider.SetClient(mockDCGM) + // No GetSupportedMetricGroups expectation set at all: any call would + // fail the test via gomock's unexpected-call panic. + + assert.True(t, ModelSupportsDCPField("NVIDIA RTX2000", 1, fbUsed)) +} + +// TestModelSupportsDCPField_QueriesOncePerModelThenCaches is the regression +// test for the actual bug report on NVIDIA/dcgm-exporter#736: a scrape used +// to ask DCGM for a DCP field's value on every entity regardless of model, +// which DCGM answers with DCGM_ST_NOT_WATCHED for a model that was never +// watched for it, triggering an unrecoverable "repairing stale watch" loop. +// This proves the model-aware answer, and that repeated scrapes (the normal +// case - this gets called every scrape interval) don't requery DCGM each time. +func TestModelSupportsDCPField_QueriesOncePerModelThenCaches(t *testing.T) { + realDCGM := dcgmprovider.Client() + defer dcgmprovider.SetClient(realDCGM) + defer ResetDCPCapabilityCache() + + ctrl := gomock.NewController(t) + mockDCGM := mockdcgm.NewMockDCGM(ctrl) + dcgmprovider.SetClient(mockDCGM) + + // Exactly one query expected for the A30 despite three lookups below - + // the second and third must come from the cache. + mockDCGM.EXPECT().GetSupportedMetricGroups(uint(0)).Return([]dcgm.MetricGroup{ + {FieldIds: []uint{uint(fp64Active)}}, + }, nil).Times(1) + // The RTX2000 is never queried with GPU ID 0: it must be looked up + // under its own representative GPU ID, and it supports nothing. + mockDCGM.EXPECT().GetSupportedMetricGroups(uint(1)).Return([]dcgm.MetricGroup{}, nil).Times(1) + + assert.True(t, ModelSupportsDCPField("NVIDIA A30", 0, fp64Active), "A30 supports fp64Active") + assert.True(t, ModelSupportsDCPField("NVIDIA A30", 0, fp64Active), "second call for the same model must hit the cache, not DCGM again") + assert.False(t, ModelSupportsDCPField("NVIDIA RTX2000", 1, fp64Active), "RTX2000 does not support fp64Active") +} + +// TestModelSupportsDCPField_ResetClearsCache proves ResetDCPCapabilityCache +// actually forces a fresh query instead of serving a stale cached answer - +// this is what queryDCPMetrics calls on every DCGM reinit. +func TestModelSupportsDCPField_ResetClearsCache(t *testing.T) { + realDCGM := dcgmprovider.Client() + defer dcgmprovider.SetClient(realDCGM) + defer ResetDCPCapabilityCache() + + ctrl := gomock.NewController(t) + mockDCGM := mockdcgm.NewMockDCGM(ctrl) + dcgmprovider.SetClient(mockDCGM) + + gomock.InOrder( + mockDCGM.EXPECT().GetSupportedMetricGroups(uint(0)).Return([]dcgm.MetricGroup{}, nil), + mockDCGM.EXPECT().GetSupportedMetricGroups(uint(0)).Return([]dcgm.MetricGroup{ + {FieldIds: []uint{uint(fp64Active)}}, + }, nil), + ) + + assert.False(t, ModelSupportsDCPField("NVIDIA A30", 0, fp64Active), "first query reports unsupported") + ResetDCPCapabilityCache() + assert.True(t, ModelSupportsDCPField("NVIDIA A30", 0, fp64Active), "after reset, the second query's answer must be used, not the cached first one") +} + +// TestModelSupportsDCPField_QueryErrorFailsOpen covers the case where DCGM +// can't answer the capability query at scrape time: filtering must not +// silently drop a field that might genuinely be fine, since the +// watch-registration path (not this one) is what actually decides whether a +// field gets watched at all. +func TestModelSupportsDCPField_QueryErrorFailsOpen(t *testing.T) { + realDCGM := dcgmprovider.Client() + defer dcgmprovider.SetClient(realDCGM) + defer ResetDCPCapabilityCache() + + ctrl := gomock.NewController(t) + mockDCGM := mockdcgm.NewMockDCGM(ctrl) + dcgmprovider.SetClient(mockDCGM) + + mockDCGM.EXPECT().GetSupportedMetricGroups(uint(0)).Return(nil, errors.New("dcgm connection error")) + + assert.True(t, ModelSupportsDCPField("NVIDIA A30", 0, fp64Active)) +} + +// TestWatchDeviceFieldGroups_PartitionsMixedGPUModels reproduces issue #657 +// end to end: a node with an A30 (supports fp64Active) and an RTX2000 +// (doesn't) used to register one shared group with fp64Active in it, and +// DCGM rejected the whole watch. This asserts the fix instead: one group +// per model, with fp64Active only watched against the A30's group. +func TestWatchDeviceFieldGroups_PartitionsMixedGPUModels(t *testing.T) { + realDCGM := dcgmprovider.Client() + defer dcgmprovider.SetClient(realDCGM) + + ctrl := gomock.NewController(t) + mockDCGM := mockdcgm.NewMockDCGM(ctrl) + dcgmprovider.SetClient(mockDCGM) + + mockDeviceInfo := mockdeviceinfo.NewMockProvider(ctrl) + mockDeviceInfo.EXPECT().InfoType().Return(dcgm.FE_GPU).AnyTimes() + mockDeviceInfo.EXPECT().GOpts().Return(appconfig.DeviceOptions{Flex: true}).AnyTimes() + mockDeviceInfo.EXPECT().GPUCount().Return(uint(2)).AnyTimes() + + a30 := deviceinfo.GPUInfo{DeviceInfo: dcgm.Device{GPU: 0, Identifiers: dcgm.DeviceIdentifiers{Model: "NVIDIA A30"}}} + rtx2000 := deviceinfo.GPUInfo{DeviceInfo: dcgm.Device{GPU: 1, Identifiers: dcgm.DeviceIdentifiers{Model: "NVIDIA RTX2000"}}} + mockDeviceInfo.EXPECT().GPU(uint(0)).Return(a30).AnyTimes() + mockDeviceInfo.EXPECT().GPU(uint(1)).Return(rtx2000).AnyTimes() + mockDeviceInfo.EXPECT().GPUs().Return([]deviceinfo.GPUInfo{a30, rtx2000}).AnyTimes() + + a30Group := dcgm.GroupHandle{} + a30Group.SetHandle(1) + rtxGroup := dcgm.GroupHandle{} + rtxGroup.SetHandle(2) + mockDCGM.EXPECT().CreateGroup(gomock.Any()).Return(a30Group, nil) + mockDCGM.EXPECT().CreateGroup(gomock.Any()).Return(rtxGroup, nil) + mockDCGM.EXPECT().AddEntityToGroup(a30Group, dcgm.FE_GPU, uint(0)).Return(nil) + mockDCGM.EXPECT().AddEntityToGroup(rtxGroup, dcgm.FE_GPU, uint(1)).Return(nil) + + // A30 supports the DCP field; RTX2000 doesn't report it as supported. + mockDCGM.EXPECT().GetSupportedMetricGroups(uint(0)).Return([]dcgm.MetricGroup{ + {FieldIds: []uint{uint(fp64Active)}}, + }, nil) + mockDCGM.EXPECT().GetSupportedMetricGroups(uint(1)).Return([]dcgm.MetricGroup{}, nil) + + a30FieldGroup := dcgm.FieldHandle{} + a30FieldGroup.SetHandle(10) + rtxFieldGroup := dcgm.FieldHandle{} + rtxFieldGroup.SetHandle(20) + + // The A30 group watches both fields; the RTX2000 group watches only the + // ordinary one. If this were still one shared group, the RTX2000's lack + // of fp64Active support would fail the single WatchFieldsWithGroupEx + // call for both GPUs - that's the bug. + mockDCGM.EXPECT().FieldGroupCreate(gomock.Any(), []dcgm.Short{fbUsed, fp64Active}).Return(a30FieldGroup, nil) + mockDCGM.EXPECT().WatchFieldsWithGroupEx(a30FieldGroup, a30Group, gomock.Any(), gomock.Any(), gomock.Any()).Return(nil) + + mockDCGM.EXPECT().FieldGroupCreate(gomock.Any(), []dcgm.Short{fbUsed}).Return(rtxFieldGroup, nil) + mockDCGM.EXPECT().WatchFieldsWithGroupEx(rtxFieldGroup, rtxGroup, gomock.Any(), gomock.Any(), gomock.Any()).Return(nil) + + watcher := NewDeviceWatcher() + groups, fieldGroups, cleanups, err := watcher.WatchDeviceFieldGroups( + []FieldWatchGroup{{Name: "default", Fields: []dcgm.Short{fbUsed, fp64Active}, IntervalMSec: 1000}}, + mockDeviceInfo, + ) + + require.NoError(t, err) + assert.ElementsMatch(t, []dcgm.GroupHandle{a30Group, rtxGroup}, groups) + assert.ElementsMatch(t, []dcgm.FieldHandle{a30FieldGroup, rtxFieldGroup}, fieldGroups) + require.Len(t, cleanups, 1) + + // Cleanup should tear down both groups without error. + mockDCGM.EXPECT().UnwatchFields(gomock.Any(), gomock.Any()).Return(nil).AnyTimes() + mockDCGM.EXPECT().FieldGroupDestroy(gomock.Any()).Return(nil).AnyTimes() + mockDCGM.EXPECT().DestroyGroup(gomock.Any()).Return(nil).AnyTimes() + cleanups[0]() +} + +// TestWatchDeviceFieldGroups_SingleModelUnchanged makes sure the common case +// - one GPU model on the node - still goes through the original single +// shared-group path, not the partitioned one, so there's no behavior change +// for the overwhelming majority of real deployments. +func TestWatchDeviceFieldGroups_SingleModelUnchanged(t *testing.T) { + realDCGM := dcgmprovider.Client() + defer dcgmprovider.SetClient(realDCGM) + + ctrl := gomock.NewController(t) + mockDCGM := mockdcgm.NewMockDCGM(ctrl) + dcgmprovider.SetClient(mockDCGM) + + mockDeviceInfo := mockdeviceinfo.NewMockProvider(ctrl) + mockDeviceInfo.EXPECT().InfoType().Return(dcgm.FE_GPU).AnyTimes() + mockDeviceInfo.EXPECT().GOpts().Return(appconfig.DeviceOptions{Flex: true}).AnyTimes() + mockDeviceInfo.EXPECT().GPUCount().Return(uint(2)).AnyTimes() + + gpu0 := deviceinfo.GPUInfo{DeviceInfo: dcgm.Device{GPU: 0, Identifiers: dcgm.DeviceIdentifiers{Model: "NVIDIA A30"}}} + gpu1 := deviceinfo.GPUInfo{DeviceInfo: dcgm.Device{GPU: 1, Identifiers: dcgm.DeviceIdentifiers{Model: "NVIDIA A30"}}} + mockDeviceInfo.EXPECT().GPU(uint(0)).Return(gpu0).AnyTimes() + mockDeviceInfo.EXPECT().GPU(uint(1)).Return(gpu1).AnyTimes() + mockDeviceInfo.EXPECT().GPUs().Return([]deviceinfo.GPUInfo{gpu0, gpu1}).AnyTimes() + + group := dcgm.GroupHandle{} + group.SetHandle(1) + mockDCGM.EXPECT().CreateGroup(gomock.Any()).Return(group, nil) + mockDCGM.EXPECT().AddEntityToGroup(group, dcgm.FE_GPU, uint(0)).Return(nil) + mockDCGM.EXPECT().AddEntityToGroup(group, dcgm.FE_GPU, uint(1)).Return(nil) + + // No GetSupportedMetricGroups call expected: a single-model node never + // enters the partitioned path at all. + fieldGroup := dcgm.FieldHandle{} + fieldGroup.SetHandle(10) + mockDCGM.EXPECT().FieldGroupCreate(gomock.Any(), []dcgm.Short{fbUsed, fp64Active}).Return(fieldGroup, nil) + mockDCGM.EXPECT().WatchFieldsWithGroupEx(fieldGroup, group, gomock.Any(), gomock.Any(), gomock.Any()).Return(nil) + + watcher := NewDeviceWatcher() + groups, fieldGroups, cleanups, err := watcher.WatchDeviceFieldGroups( + []FieldWatchGroup{{Name: "default", Fields: []dcgm.Short{fbUsed, fp64Active}, IntervalMSec: 1000}}, + mockDeviceInfo, + ) + + require.NoError(t, err) + assert.Equal(t, []dcgm.GroupHandle{group}, groups) + assert.Equal(t, []dcgm.FieldHandle{fieldGroup}, fieldGroups) + require.Len(t, cleanups, 1) + + mockDCGM.EXPECT().UnwatchFields(gomock.Any(), gomock.Any()).Return(nil).AnyTimes() + mockDCGM.EXPECT().FieldGroupDestroy(gomock.Any()).Return(nil).AnyTimes() + mockDCGM.EXPECT().DestroyGroup(gomock.Any()).Return(nil).AnyTimes() + cleanups[0]() +} diff --git a/internal/pkg/devicewatcher/device_watcher.go b/internal/pkg/devicewatcher/device_watcher.go index 3106cbbe..67b3784b 100644 --- a/internal/pkg/devicewatcher/device_watcher.go +++ b/internal/pkg/devicewatcher/device_watcher.go @@ -20,7 +20,9 @@ import ( "context" "fmt" "log/slog" + "sort" "strings" + "sync" "github.com/NVIDIA/go-dcgm/pkg/dcgm" @@ -166,6 +168,24 @@ func (d *DeviceWatcher) WatchDeviceFields( func (d *DeviceWatcher) WatchDeviceFieldGroups( fieldWatchGroups []FieldWatchGroup, deviceInfo deviceinfo.Provider, ) ([]dcgm.GroupHandle, []dcgm.FieldHandle, []func(), error) { + // DCP/profiling fields (DCGM_FI_PROF_*) are validated for an entire DCGM + // watch group at registration time: if any GPU model in the group lacks + // a requested profiling field, the whole watch call fails. Ordinary + // fields don't have this problem - unsupported ones just come back as a + // per-entity NOT_SUPPORTED value at scrape time (see isBlankValue in the + // collector package). So the single-shared-group path below is fine for + // almost every field on almost every node; it only breaks on a node + // mixing GPU models where at least one profiling field isn't supported + // everywhere. Route that specific case through a per-model partition + // instead of touching the common path. + if deviceInfo.InfoType() == dcgm.FE_GPU && anyDCPField(fieldWatchGroups) { + if handled, groups, fieldGroups, cleanups, err := d.watchFieldGroupsPartitionedByModel( + fieldWatchGroups, deviceInfo, + ); handled { + return groups, fieldGroups, cleanups, err + } + } + resources := &WatchResources{} // Create groups based on device type @@ -216,6 +236,204 @@ func (d *DeviceWatcher) WatchDeviceFieldGroups( return resources.groups, resources.fieldGroups, []func(){cleanup}, nil } +// anyDCPField reports whether any field across the given watch groups is a +// DCP/profiling field (DCGM_FI_PROF_*). +func anyDCPField(fieldWatchGroups []FieldWatchGroup) bool { + for _, fieldWatchGroup := range fieldWatchGroups { + for _, fieldID := range fieldWatchGroup.Fields { + if counters.IsDCPField(uint(fieldID)) { + return true + } + } + } + return false +} + +// watchFieldGroupsPartitionedByModel watches GPU fields one DCGM group per +// GPU model instead of one shared group for the whole node. handled is false +// when the node only has one GPU model present, so the caller should fall +// back to its normal single-group path unchanged. +// +// This exists because DCP/profiling fields fail the entire group's watch +// registration if any member GPU doesn't support them (unlike ordinary +// fields, which just report a per-entity NOT_SUPPORTED value). On a node +// mixing GPU models, that turns "one model doesn't support this profiling +// field" into "the exporter won't start at all". Partitioning by model and +// asking DCGM what each model actually supports keeps a field enabled +// wherever it works instead of dropping it - or crashing - node-wide. +func (d *DeviceWatcher) watchFieldGroupsPartitionedByModel( + fieldWatchGroups []FieldWatchGroup, deviceInfo deviceinfo.Provider, +) (handled bool, groups []dcgm.GroupHandle, fieldGroups []dcgm.FieldHandle, cleanups []func(), err error) { + monitoringInfo := devicemonitoring.GetMonitoredEntities(deviceInfo) + models, byModel := partitionMonitoredEntitiesByModel(monitoringInfo) + if len(models) < 2 { + return false, nil, nil, nil, nil + } + + resources := &WatchResources{} + for _, model := range models { + entities := byModel[model] + + groupID, _, gerr := createGroupFromEntities(entities) + if gerr != nil { + resources.Cleanup() + return true, nil, nil, nil, gerr + } + resources.groups = append(resources.groups, *groupID) + + supported, serr := supportedDCPFields(entities[0].DeviceInfo.GPU) + if serr != nil { + // Can't determine this model's profiling capability. Fall back + // to watching only what we know is safe for every model rather + // than risk the same all-or-nothing failure this path exists + // to avoid. + slog.Warn("Could not query DCP metric group support for GPU model; "+ + "watching non-profiling fields only for it", + slog.String("gpu_model", model), slog.String(ErrorKey, serr.Error())) + supported = map[dcgm.Short]bool{} + } + + for _, fieldWatchGroup := range fieldWatchGroups { + fields := filterFieldsForModel(dedupeFields(fieldWatchGroup.Fields), supported) + if len(fields) == 0 { + continue + } + + fieldGroup, ferr := newFieldGroupSimple(fields) + if ferr != nil { + resources.Cleanup() + return true, nil, nil, nil, ferr + } + resources.fieldGroups = append(resources.fieldGroups, fieldGroup) + + logWatchFieldsCall(deviceInfo, *groupID, fieldGroup, fields) + if werr := watchFieldGroupSimple(*groupID, fieldGroup, fieldWatchGroup.IntervalMSec*1000); werr != nil { + logWatchFieldsFailure(deviceInfo, *groupID, fieldGroup, fields, werr) + resources.Cleanup() + return true, nil, nil, nil, werr + } + resources.hasWatch = true + } + } + + cleanup := func() { resources.Cleanup() } + return true, resources.groups, resources.fieldGroups, []func(){cleanup}, nil +} + +// partitionMonitoredEntitiesByModel groups monitored entities by their +// parent GPU's model name. models is sorted for deterministic group +// creation order (matters for tests, harmless in production). +func partitionMonitoredEntitiesByModel( + monitoringInfo []devicemonitoring.Info, +) ([]string, map[string][]devicemonitoring.Info) { + byModel := make(map[string][]devicemonitoring.Info) + for _, mi := range monitoringInfo { + model := mi.DeviceInfo.Identifiers.Model + byModel[model] = append(byModel[model], mi) + } + + models := make([]string, 0, len(byModel)) + for model := range byModel { + models = append(models, model) + } + sort.Strings(models) + + return models, byModel +} + +// supportedDCPFields returns the DCP/profiling field IDs DCGM reports as +// supported for the GPU model represented by gpuID. Every GPU of the same +// model supports the same profiling fields, so one representative GPU per +// model is enough - no need to query every GPU individually. +func supportedDCPFields(gpuID uint) (map[dcgm.Short]bool, error) { + metricGroups, err := dcgmprovider.Client().GetSupportedMetricGroups(gpuID) + if err != nil { + return nil, err + } + + supported := make(map[dcgm.Short]bool) + for _, metricGroup := range metricGroups { + for _, fieldID := range metricGroup.FieldIds { + supported[dcgm.Short(fieldID)] = true + } + } + + return supported, nil +} + +var ( + dcpCapabilityCacheMu sync.RWMutex + dcpCapabilityCache = map[string]map[dcgm.Short]bool{} // GPU model -> supported DCP fields +) + +// ModelSupportsDCPField reports whether the GPU model identified by +// representativeGPUID supports the given field, so callers can avoid asking +// DCGM for a value that was never watched. Watch registration is already +// scoped per GPU model (see watchFieldGroupsPartitionedByModel), but a +// scrape reads every entity with the same configured field list regardless +// of model. For an ordinary field that's harmless - DCGM reports an +// unsupported field as a per-entity blank value - but for a DCP field it +// isn't: asking for a DCP field's value on an entity whose model was never +// watched for it returns DCGM_ST_NOT_WATCHED, and the exporter's stale-watch +// repair logic then retries forever trying to fix a watch that was +// deliberately never created. Filtering the scrape-time field list with this +// avoids ever asking. The answer is cached per model, since it can't change +// for a fixed physical GPU model within one DCGM session; ResetDCPCapabilityCache +// clears it across a reinit. Non-DCP fields always report supported without +// touching the cache or DCGM. See NVIDIA/dcgm-exporter#736. +func ModelSupportsDCPField(model string, representativeGPUID uint, fieldID dcgm.Short) bool { + if !counters.IsDCPField(uint(fieldID)) { + return true + } + + dcpCapabilityCacheMu.RLock() + supported, cached := dcpCapabilityCache[model] + dcpCapabilityCacheMu.RUnlock() + if cached { + return supported[fieldID] + } + + supported, err := supportedDCPFields(representativeGPUID) + if err != nil { + // Can't determine support one way or the other. Fail open (treat it + // as supported) rather than silently dropping a field that might be + // fine - the watch-registration path already handles the real + // unsupported case; this is scrape-time filtering only. + slog.Warn("Could not query DCP metric group support while filtering scrape fields; leaving field in scrape list", + slog.String("gpu_model", model), slog.String(ErrorKey, err.Error())) + return true + } + + dcpCapabilityCacheMu.Lock() + dcpCapabilityCache[model] = supported + dcpCapabilityCacheMu.Unlock() + + return supported[fieldID] +} + +// ResetDCPCapabilityCache clears the cached per-model DCP support answers. +// Called on DCGM reinit, since profiling capability is queried through the +// DCGM connection and could read differently on a fresh session. +func ResetDCPCapabilityCache() { + dcpCapabilityCacheMu.Lock() + defer dcpCapabilityCacheMu.Unlock() + dcpCapabilityCache = map[string]map[dcgm.Short]bool{} +} + +// filterFieldsForModel drops DCP fields a model's profiling query didn't +// report as supported. Non-DCP fields always pass through unfiltered: DCGM +// already reports those as a per-entity NOT_SUPPORTED value at scrape time +// instead of failing the watch, so there's nothing to filter for them. +func filterFieldsForModel(fields []dcgm.Short, supported map[dcgm.Short]bool) []dcgm.Short { + filtered := make([]dcgm.Short, 0, len(fields)) + for _, fieldID := range fields { + if !counters.IsDCPField(uint(fieldID)) || supported[fieldID] { + filtered = append(filtered, fieldID) + } + } + return filtered +} + func (d *DeviceWatcher) createGenericGroup(deviceInfo deviceinfo.Provider) (*dcgm.GroupHandle, func(), error, ) { @@ -224,6 +442,17 @@ func (d *DeviceWatcher) createGenericGroup(deviceInfo deviceinfo.Provider) (*dcg return nil, doNothing, nil } + return createGroupFromEntities(monitoringInfo) +} + +// createGroupFromEntities creates one DCGM group containing exactly the given +// entities. Used both for the single-group path (all monitored entities) and +// for the per-model partitioned path (one model's entities at a time). +func createGroupFromEntities(monitoringInfo []devicemonitoring.Info) (*dcgm.GroupHandle, func(), error) { + if len(monitoringInfo) == 0 { + return nil, doNothing, nil + } + groupID, cleanup, err := createGroup() if err != nil { return nil, cleanup, err diff --git a/pkg/cmd/app.go b/pkg/cmd/app.go index de0397df..2170b2fc 100644 --- a/pkg/cmd/app.go +++ b/pkg/cmd/app.go @@ -1251,6 +1251,12 @@ func getCounters(ctx context.Context, config *appconfig.Config) (*counters.Count func (r *reloadCoordinator) queryDCPMetrics(cfg *appconfig.Config, reloadID uint64) { slog.Debug("Querying DCGM profiling metric groups", slog.Uint64("reload_id", reloadID)) + // This runs on startup and after every DCGM reinit, the same cadence the + // per-model DCP capability cache needs to stay valid on: capability is + // queried through the DCGM connection, so a fresh session could in + // principle answer differently than the last one. + devicewatcher.ResetDCPCapabilityCache() + defer func() { if p := recover(); p != nil { slog.Warn("Profiling API panic - DCP metrics disabled",