Просмотр исходного кода

Update device and include in bingen (#4014)

Signed-off-by: Sean Holcomb <seanholcomb@gmail.com>
Sean Holcomb 2 дней назад
Родитель
Сommit
ba0e947ad0
42 измененных файлов с 1295 добавлено и 847 удалено
  1. 41 0
      core/pkg/compute/kubemodel/container.go
  2. 70 0
      core/pkg/compute/kubemodel/container_test.go
  3. 14 0
      core/pkg/compute/kubemodel/daemonset.go
  4. 31 1
      core/pkg/compute/kubemodel/daemonset_test.go
  5. 0 89
      core/pkg/compute/kubemodel/dcgmdevice.go
  6. 0 172
      core/pkg/compute/kubemodel/dcgmdevice_test.go
  7. 54 0
      core/pkg/compute/kubemodel/device.go
  8. 112 0
      core/pkg/compute/kubemodel/device_test.go
  9. 1 1
      core/pkg/compute/kubemodel/kubemodel.go
  10. 1 1
      core/pkg/compute/kubemodel/kubemodel_test.go
  11. 1 1
      core/pkg/model/kubemodel/bingen.go
  12. 28 14
      core/pkg/model/kubemodel/container.go
  13. 8 8
      core/pkg/model/kubemodel/daemonset.go
  14. 0 66
      core/pkg/model/kubemodel/dcgm.go
  15. 52 0
      core/pkg/model/kubemodel/device.go
  16. 20 20
      core/pkg/model/kubemodel/device_test.go
  17. 3 3
      core/pkg/model/kubemodel/kubemodel.go
  18. 358 447
      core/pkg/model/kubemodel/kubemodel_codecs.go
  19. 1 2
      core/pkg/model/kubemodel/kubemodel_helpers_test.go
  20. 10 11
      core/pkg/model/kubemodel/mock.go
  21. 2 0
      core/pkg/source/datasource.go
  22. 25 0
      core/pkg/source/decoders.go
  23. 6 0
      core/pkg/source/mock.go
  24. 4 0
      core/pkg/source/noop.go
  25. 5 0
      core/pkg/source/record.go
  26. 2 0
      core/pkg/source/record_test.go
  27. 25 0
      core/pkg/util/args.go
  28. 14 0
      modules/collector-source/pkg/collector/collector.go
  29. 4 0
      modules/collector-source/pkg/collector/metricsquerier.go
  30. 1 0
      modules/collector-source/pkg/metric/collector.go
  31. 1 0
      modules/collector-source/pkg/metric/metrics.go
  32. 21 0
      modules/collector-source/pkg/scrape/clustercache.go
  33. 90 0
      modules/collector-source/pkg/scrape/clustercache_test.go
  34. 27 3
      modules/collector-source/pkg/scrape/dcgm.go
  35. 150 0
      modules/collector-source/pkg/scrape/dcgm_test.go
  36. 15 0
      modules/collector-source/pkg/scrape/index.go
  37. 2 1
      modules/collector-source/pkg/scrape/network.go
  38. 2 1
      modules/collector-source/pkg/scrape/opencost.go
  39. 15 3
      modules/collector-source/pkg/scrape/targetscraper.go
  40. 45 3
      modules/collector-source/pkg/scrape/targetscraper_test.go
  41. 18 0
      modules/prometheus-source/pkg/prom/metricsquerier.go
  42. 16 0
      pkg/metrics/kubemodel.go

+ 41 - 0
core/pkg/compute/kubemodel/container.go

@@ -38,6 +38,9 @@ func (km *KubeModel) computeContainers(kms *kubemodel.KubeModelSet, start, end t
 	ramUsageAvgFuture := source.WithGroup(grp, metrics.QueryRAMUsageAvg(start, end))
 	ramUsageMaxFuture := source.WithGroup(grp, metrics.QueryRAMUsageMax(start, end))
 
+	deviceUsageAvgFuture := source.WithGroup(grp, metrics.QueryDCGMContainerUsageAvg(start, end))
+	deviceUsageMaxFuture := source.WithGroup(grp, metrics.QueryDCGMContainerUsageMax(start, end))
+
 	type containerKey struct {
 		podUID string
 		name   string
@@ -161,6 +164,44 @@ func (km *KubeModel) computeContainers(kms *kubemodel.KubeModelSet, start, end t
 		}
 	}
 
+	deviceUsageAvgResult, _ := deviceUsageAvgFuture.Await()
+	for _, res := range deviceUsageAvgResult {
+		if res.PodUID == "" || res.Container == "" {
+			continue
+		}
+		key := containerKey{podUID: res.PodUID, name: res.Container}
+		container, ok := containerMap[key]
+		if !ok {
+			log.Warnf("container %s/%s has not been initialized to add device usage avg", res.PodUID, res.Container)
+			continue
+		}
+		if container.DeviceUsages == nil {
+			container.DeviceUsages = make(map[string]kubemodel.DeviceUsage)
+		}
+		usage := container.DeviceUsages[res.UUID]
+		usage.UsageAvg = res.Value
+		container.DeviceUsages[res.UUID] = usage
+	}
+
+	deviceUsageMaxResult, _ := deviceUsageMaxFuture.Await()
+	for _, res := range deviceUsageMaxResult {
+		if res.PodUID == "" || res.Container == "" {
+			continue
+		}
+		key := containerKey{podUID: res.PodUID, name: res.Container}
+		container, ok := containerMap[key]
+		if !ok {
+			log.Warnf("container %s/%s has not been initialized to add device usage max", res.PodUID, res.Container)
+			continue
+		}
+		if container.DeviceUsages == nil {
+			container.DeviceUsages = make(map[string]kubemodel.DeviceUsage)
+		}
+		usage := container.DeviceUsages[res.UUID]
+		usage.UsageMax = res.Value
+		container.DeviceUsages[res.UUID] = usage
+	}
+
 	for _, container := range containerMap {
 		err := kms.RegisterContainer(container)
 		if err != nil {

+ 70 - 0
core/pkg/compute/kubemodel/container_test.go

@@ -161,6 +161,76 @@ func TestComputeContainers(t *testing.T) {
 				},
 			},
 		},
+		{
+			name: "device usage avg and max are populated",
+			overrides: map[string]any{
+				source.QueryContainerUptime: []*source.ContainerUptimeResult{
+					{UptimeResult: source.UptimeResult{UID: "pod-1", First: start, Last: end}, Container: "training"},
+				},
+				source.QueryDCGMContainerUsageAvg: []*source.DCGMDeviceContainerUsageResult{
+					{UUID: "GPU-abc123", PodUID: "pod-1", Container: "training", Value: 0.75},
+				},
+				source.QueryDCGMContainerUsageMax: []*source.DCGMDeviceContainerUsageResult{
+					{UUID: "GPU-abc123", PodUID: "pod-1", Container: "training", Value: 0.95},
+				},
+			},
+			want: map[string]*kubemodel.Container{
+				"pod-1/training": {
+					PodUID:           "pod-1",
+					Name:             "training",
+					Start:            start,
+					End:              end,
+					ResourceRequests: kubemodel.ResourceQuantities{},
+					ResourceLimits:   kubemodel.ResourceQuantities{},
+					DeviceUsages: map[string]kubemodel.DeviceUsage{
+						"GPU-abc123": {UsageAvg: 0.75, UsageMax: 0.95},
+					},
+				},
+			},
+		},
+		{
+			name: "device usage with empty pod uid or container is ignored",
+			overrides: map[string]any{
+				source.QueryContainerUptime: []*source.ContainerUptimeResult{
+					{UptimeResult: source.UptimeResult{UID: "pod-1", First: start, Last: end}, Container: "training"},
+				},
+				source.QueryDCGMContainerUsageAvg: []*source.DCGMDeviceContainerUsageResult{
+					{UUID: "GPU-abc123", PodUID: "", Container: "training", Value: 0.5},
+					{UUID: "GPU-abc123", PodUID: "pod-1", Container: "", Value: 0.5},
+				},
+			},
+			want: map[string]*kubemodel.Container{
+				"pod-1/training": {
+					PodUID:           "pod-1",
+					Name:             "training",
+					Start:            start,
+					End:              end,
+					ResourceRequests: kubemodel.ResourceQuantities{},
+					ResourceLimits:   kubemodel.ResourceQuantities{},
+				},
+			},
+		},
+		{
+			name: "device usage for unknown container is ignored",
+			overrides: map[string]any{
+				source.QueryContainerUptime: []*source.ContainerUptimeResult{
+					{UptimeResult: source.UptimeResult{UID: "pod-1", First: start, Last: end}, Container: "main"},
+				},
+				source.QueryDCGMContainerUsageAvg: []*source.DCGMDeviceContainerUsageResult{
+					{UUID: "GPU-abc123", PodUID: "pod-1", Container: "training", Value: 0.75},
+				},
+			},
+			want: map[string]*kubemodel.Container{
+				"pod-1/main": {
+					PodUID:           "pod-1",
+					Name:             "main",
+					Start:            start,
+					End:              end,
+					ResourceRequests: kubemodel.ResourceQuantities{},
+					ResourceLimits:   kubemodel.ResourceQuantities{},
+				},
+			},
+		},
 		{
 			name: "resource requests for unknown container are ignored",
 			overrides: map[string]any{

+ 14 - 0
core/pkg/compute/kubemodel/daemonset.go

@@ -16,6 +16,7 @@ func (km *KubeModel) computeDaemonSets(kms *kubemodel.KubeModelSet, start, end t
 	daemonSetUptimeResultFuture := source.WithGroup(grp, metrics.QueryDaemonSetUptime(start, end))
 	daemonSetLabelsResultFuture := source.WithGroup(grp, metrics.QueryDaemonSetLabels(start, end))
 	daemonSetAnnotationsResultFuture := source.WithGroup(grp, metrics.QueryDaemonSetAnnotations(start, end))
+	daemonSetArgumentsResultFuture := source.WithGroup(grp, metrics.QueryDaemonSetArguments(start, end))
 
 	daemonSetMap := make(map[string]*kubemodel.DaemonSet)
 
@@ -60,6 +61,19 @@ func (km *KubeModel) computeDaemonSets(kms *kubemodel.KubeModelSet, start, end t
 		daemonSet.Annotations = res.Annotations
 	}
 
+	daemonSetArgumentsResult, _ := daemonSetArgumentsResultFuture.Await()
+	for _, res := range daemonSetArgumentsResult {
+		daemonSet, ok := daemonSetMap[res.UID]
+		if !ok {
+			log.Warnf("daemonset with UID '%s' has not been initialized to add arguments", res.UID)
+			continue
+		}
+		if daemonSet.Arguments == nil {
+			daemonSet.Arguments = make(map[string]string)
+		}
+		daemonSet.Arguments[res.Arg] = res.Value
+	}
+
 	for _, daemonSet := range daemonSetMap {
 		err := kms.RegisterDaemonSet(daemonSet)
 		if err != nil {

+ 31 - 1
core/pkg/compute/kubemodel/daemonset_test.go

@@ -67,7 +67,7 @@ func TestComputeDaemonSets(t *testing.T) {
 			want: map[string]*kubemodel.DaemonSet{},
 		},
 		{
-			name: "daemonset labels and annotations are attached",
+			name: "daemonset labels, annotations, and arguments are attached",
 			overrides: map[string]any{
 				source.QueryDaemonSetInfo: []*source.DaemonSetInfoResult{
 					{UID: "ds-1", DaemonSet: "fluentd", NamespaceUID: "ns-1"},
@@ -81,6 +81,10 @@ func TestComputeDaemonSets(t *testing.T) {
 				source.QueryDaemonSetAnnotations: []*source.AnnotationsResult{
 					{UID: "ds-1", Annotations: map[string]string{"managed-by": "helm"}},
 				},
+				source.QueryDaemonSetArguments: []*source.DaemonSetArgumentResult{
+					{UID: "ds-1", Arg: "vgpu", Value: "2"},
+					{UID: "ds-1", Arg: "log-level", Value: "debug"},
+				},
 			},
 			want: map[string]*kubemodel.DaemonSet{
 				"ds-1": {
@@ -91,6 +95,32 @@ func TestComputeDaemonSets(t *testing.T) {
 					End:          end,
 					Labels:       map[string]string{"component": "logging"},
 					Annotations:  map[string]string{"managed-by": "helm"},
+					Arguments:    map[string]string{"vgpu": "2", "log-level": "debug"},
+				},
+			},
+		},
+		{
+			name: "arguments for unknown daemonset are ignored",
+			overrides: map[string]any{
+				source.QueryDaemonSetInfo: []*source.DaemonSetInfoResult{
+					{UID: "ds-1", DaemonSet: "fluentd", NamespaceUID: "ns-1"},
+				},
+				source.QueryDaemonSetUptime: []*source.UptimeResult{
+					{UID: "ds-1", First: start, Last: end},
+				},
+				source.QueryDaemonSetArguments: []*source.DaemonSetArgumentResult{
+					{UID: "ds-1", Arg: "vgpu", Value: "2"},
+					{UID: "unknown-ds", Arg: "vgpu", Value: "4"},
+				},
+			},
+			want: map[string]*kubemodel.DaemonSet{
+				"ds-1": {
+					UID:          "ds-1",
+					Name:         "fluentd",
+					NamespaceUID: "ns-1",
+					Start:        start,
+					End:          end,
+					Arguments:    map[string]string{"vgpu": "2"},
 				},
 			},
 		},

+ 0 - 89
core/pkg/compute/kubemodel/dcgmdevice.go

@@ -1,89 +0,0 @@
-package kubemodel
-
-import (
-	"time"
-
-	"github.com/opencost/opencost/core/pkg/log"
-	"github.com/opencost/opencost/core/pkg/model/kubemodel"
-	"github.com/opencost/opencost/core/pkg/source"
-)
-
-func (km *KubeModel) computeDCGMDevices(kms *kubemodel.KubeModelSet, start, end time.Time) error {
-	grp := source.NewQueryGroup()
-	metrics := km.ds.Metrics()
-
-	dcgmInfoFuture := source.WithGroup(grp, metrics.QueryDCGMDeviceInfo(start, end))
-	dcgmUptimeFuture := source.WithGroup(grp, metrics.QueryDCGMDeviceUptime(start, end))
-	dcgmUsageAvgFuture := source.WithGroup(grp, metrics.QueryDCGMContainerUsageAvg(start, end))
-	dcgmUsageMaxFuture := source.WithGroup(grp, metrics.QueryDCGMContainerUsageMax(start, end))
-
-	deviceMap := make(map[string]*kubemodel.DCGMDevice)
-
-	dcgmInfoResult, _ := dcgmInfoFuture.Await()
-	for _, res := range dcgmInfoResult {
-		if res.UUID == "" {
-			continue
-		}
-		if _, ok := deviceMap[res.UUID]; ok {
-			continue
-		}
-		deviceMap[res.UUID] = &kubemodel.DCGMDevice{
-			UUID:      res.UUID,
-			Device:    res.Device,
-			ModelName: res.ModelName,
-			PodUsages: make(map[string]kubemodel.DCGMPod),
-		}
-	}
-
-	dcgmUptimeResult, _ := dcgmUptimeFuture.Await()
-	for _, res := range dcgmUptimeResult {
-		d, ok := deviceMap[res.UUID]
-		if !ok {
-			log.Warnf("DCGM uptime result for unknown device UUID '%s'", res.UUID)
-			continue
-		}
-		s, e := res.GetStartEnd(start, end, km.ds.Resolution())
-		d.Start = s
-		d.End = e
-	}
-
-	dcgmUsageAvgResult, _ := dcgmUsageAvgFuture.Await()
-	for _, res := range dcgmUsageAvgResult {
-		device, ok := deviceMap[res.UUID]
-		if !ok || res.PodUID == "" || res.Container == "" {
-			continue
-		}
-		pod, ok := device.PodUsages[res.PodUID]
-		if !ok {
-			pod = kubemodel.DCGMPod{ContainerUsages: make(map[string]kubemodel.DCGMContainer)}
-		}
-		c := pod.ContainerUsages[res.Container]
-		c.UsageAvg = res.Value
-		pod.ContainerUsages[res.Container] = c
-		device.PodUsages[res.PodUID] = pod
-	}
-
-	dcgmUsageMaxResult, _ := dcgmUsageMaxFuture.Await()
-	for _, res := range dcgmUsageMaxResult {
-		device, ok := deviceMap[res.UUID]
-		if !ok || res.PodUID == "" || res.Container == "" {
-			continue
-		}
-		pod, ok := device.PodUsages[res.PodUID]
-		if !ok {
-			pod = kubemodel.DCGMPod{ContainerUsages: make(map[string]kubemodel.DCGMContainer)}
-		}
-		c := pod.ContainerUsages[res.Container]
-		c.UsageMax = res.Value
-		pod.ContainerUsages[res.Container] = c
-		device.PodUsages[res.PodUID] = pod
-	}
-
-	for _, device := range deviceMap {
-		if err := kms.RegisterDCGMDevice(device); err != nil {
-			log.Warnf("Failed to register DCGM device: %s", err.Error())
-		}
-	}
-
-	return nil
-}

+ 0 - 172
core/pkg/compute/kubemodel/dcgmdevice_test.go

@@ -1,172 +0,0 @@
-package kubemodel
-
-import (
-	"testing"
-	"time"
-
-	"github.com/stretchr/testify/assert"
-	"github.com/stretchr/testify/require"
-
-	"github.com/opencost/opencost/core/pkg/model/kubemodel"
-	"github.com/opencost/opencost/core/pkg/source"
-)
-
-func TestComputeDCGMDevices(t *testing.T) {
-	start := time.Date(2024, 1, 1, 0, 0, 0, 0, time.UTC)
-	end := start.Add(time.Hour)
-
-	tests := []struct {
-		name      string
-		overrides map[string]any
-		want      map[string]*kubemodel.DCGMDevice
-	}{
-		{
-			name:      "no data returns empty dcgm device map",
-			overrides: map[string]any{},
-			want:      map[string]*kubemodel.DCGMDevice{},
-		},
-		{
-			name: "basic dcgm device info and uptime",
-			overrides: map[string]any{
-				source.QueryDCGMDeviceInfo: []*source.DCGMDeviceInfoResult{
-					{UUID: "GPU-abc123", Device: "nvidia0", ModelName: "A100"},
-				},
-				source.QueryDCGMDeviceUptime: []*source.DCGMDeviceUptimeResult{
-					{UUID: "GPU-abc123", First: start, Last: end},
-				},
-			},
-			want: map[string]*kubemodel.DCGMDevice{
-				"GPU-abc123": {
-					UUID:      "GPU-abc123",
-					Device:    "nvidia0",
-					ModelName: "A100",
-					Start:     start,
-					End:       end,
-					PodUsages: map[string]kubemodel.DCGMPod{},
-				},
-			},
-		},
-		{
-			name: "dcgm device without uptime is not registered",
-			overrides: map[string]any{
-				source.QueryDCGMDeviceInfo: []*source.DCGMDeviceInfoResult{
-					{UUID: "GPU-abc123", Device: "nvidia0", ModelName: "A100"},
-				},
-			},
-			want: map[string]*kubemodel.DCGMDevice{},
-		},
-		{
-			name: "dcgm device with empty uuid is skipped",
-			overrides: map[string]any{
-				source.QueryDCGMDeviceInfo: []*source.DCGMDeviceInfoResult{
-					{UUID: "", Device: "nvidia0", ModelName: "A100"},
-				},
-				source.QueryDCGMDeviceUptime: []*source.DCGMDeviceUptimeResult{
-					{UUID: "GPU-abc123", First: start, Last: end},
-				},
-			},
-			want: map[string]*kubemodel.DCGMDevice{},
-		},
-		{
-			name: "dcgm container usage avg and max are populated",
-			overrides: map[string]any{
-				source.QueryDCGMDeviceInfo: []*source.DCGMDeviceInfoResult{
-					{UUID: "GPU-abc123", Device: "nvidia0", ModelName: "A100"},
-				},
-				source.QueryDCGMDeviceUptime: []*source.DCGMDeviceUptimeResult{
-					{UUID: "GPU-abc123", First: start, Last: end},
-				},
-				source.QueryDCGMContainerUsageAvg: []*source.DCGMDeviceContainerUsageResult{
-					{UUID: "GPU-abc123", PodUID: "pod-1", Container: "training", Value: 0.75},
-				},
-				source.QueryDCGMContainerUsageMax: []*source.DCGMDeviceContainerUsageResult{
-					{UUID: "GPU-abc123", PodUID: "pod-1", Container: "training", Value: 0.95},
-				},
-			},
-			want: map[string]*kubemodel.DCGMDevice{
-				"GPU-abc123": {
-					UUID:      "GPU-abc123",
-					Device:    "nvidia0",
-					ModelName: "A100",
-					Start:     start,
-					End:       end,
-					PodUsages: map[string]kubemodel.DCGMPod{
-						"pod-1": {
-							ContainerUsages: map[string]kubemodel.DCGMContainer{
-								"training": {UsageAvg: 0.75, UsageMax: 0.95},
-							},
-						},
-					},
-				},
-			},
-		},
-		{
-			name: "usage with empty pod uid or container is ignored",
-			overrides: map[string]any{
-				source.QueryDCGMDeviceInfo: []*source.DCGMDeviceInfoResult{
-					{UUID: "GPU-abc123", Device: "nvidia0", ModelName: "A100"},
-				},
-				source.QueryDCGMDeviceUptime: []*source.DCGMDeviceUptimeResult{
-					{UUID: "GPU-abc123", First: start, Last: end},
-				},
-				source.QueryDCGMContainerUsageAvg: []*source.DCGMDeviceContainerUsageResult{
-					{UUID: "GPU-abc123", PodUID: "", Container: "training", Value: 0.5},
-					{UUID: "GPU-abc123", PodUID: "pod-1", Container: "", Value: 0.5},
-				},
-			},
-			want: map[string]*kubemodel.DCGMDevice{
-				"GPU-abc123": {
-					UUID:      "GPU-abc123",
-					Device:    "nvidia0",
-					ModelName: "A100",
-					Start:     start,
-					End:       end,
-					PodUsages: map[string]kubemodel.DCGMPod{},
-				},
-			},
-		},
-		{
-			name: "duplicate device info entries use first occurrence",
-			overrides: map[string]any{
-				source.QueryDCGMDeviceInfo: []*source.DCGMDeviceInfoResult{
-					{UUID: "GPU-abc123", Device: "nvidia0", ModelName: "A100"},
-					{UUID: "GPU-abc123", Device: "nvidia0-dup", ModelName: "A100-dup"},
-				},
-				source.QueryDCGMDeviceUptime: []*source.DCGMDeviceUptimeResult{
-					{UUID: "GPU-abc123", First: start, Last: end},
-				},
-			},
-			want: map[string]*kubemodel.DCGMDevice{
-				"GPU-abc123": {
-					UUID:      "GPU-abc123",
-					Device:    "nvidia0",
-					ModelName: "A100",
-					Start:     start,
-					End:       end,
-					PodUsages: map[string]kubemodel.DCGMPod{},
-				},
-			},
-		},
-	}
-
-	for _, tt := range tests {
-		t.Run(tt.name, func(t *testing.T) {
-			ds := source.NewMockOpenCostDataSource()
-			ds.ResolutionValue = 5 * time.Minute
-			seedCluster(ds, start, end)
-			for method, result := range tt.overrides {
-				ds.Querier.SetOverride(method, result)
-			}
-
-			km, err := NewKubeModel(testClusterUID, false, ds)
-			require.NoError(t, err)
-
-			kms := kubemodel.NewKubeModelSet(start, end)
-
-			err = km.computeDCGMDevices(kms, start, end)
-			require.NoError(t, err)
-
-			assert.Equal(t, tt.want, kms.DCGMDevices)
-		})
-	}
-}

+ 54 - 0
core/pkg/compute/kubemodel/device.go

@@ -0,0 +1,54 @@
+package kubemodel
+
+import (
+	"time"
+
+	"github.com/opencost/opencost/core/pkg/log"
+	"github.com/opencost/opencost/core/pkg/model/kubemodel"
+	"github.com/opencost/opencost/core/pkg/source"
+)
+
+func (km *KubeModel) computeDevices(kms *kubemodel.KubeModelSet, start, end time.Time) error {
+	grp := source.NewQueryGroup()
+	metrics := km.ds.Metrics()
+
+	infoFuture := source.WithGroup(grp, metrics.QueryDCGMDeviceInfo(start, end))
+	uptimeFuture := source.WithGroup(grp, metrics.QueryDCGMDeviceUptime(start, end))
+
+	deviceMap := make(map[string]*kubemodel.Device)
+
+	infoResult, _ := infoFuture.Await()
+	for _, res := range infoResult {
+		if res.UUID == "" {
+			continue
+		}
+		if _, ok := deviceMap[res.UUID]; ok {
+			continue
+		}
+		deviceMap[res.UUID] = &kubemodel.Device{
+			UUID:      res.UUID,
+			Device:    res.Device,
+			ModelName: res.ModelName,
+		}
+	}
+
+	uptimeResult, _ := uptimeFuture.Await()
+	for _, res := range uptimeResult {
+		d, ok := deviceMap[res.UUID]
+		if !ok {
+			log.Warnf("DCGM uptime result for unknown device UUID '%s'", res.UUID)
+			continue
+		}
+		s, e := res.GetStartEnd(start, end, km.ds.Resolution())
+		d.Start = s
+		d.End = e
+	}
+
+	for _, device := range deviceMap {
+		if err := kms.RegisterDevice(device); err != nil {
+			log.Warnf("Failed to register device: %s", err.Error())
+		}
+	}
+
+	return nil
+}

+ 112 - 0
core/pkg/compute/kubemodel/device_test.go

@@ -0,0 +1,112 @@
+package kubemodel
+
+import (
+	"testing"
+	"time"
+
+	"github.com/stretchr/testify/assert"
+	"github.com/stretchr/testify/require"
+
+	"github.com/opencost/opencost/core/pkg/model/kubemodel"
+	"github.com/opencost/opencost/core/pkg/source"
+)
+
+func TestComputeDevices(t *testing.T) {
+	start := time.Date(2024, 1, 1, 0, 0, 0, 0, time.UTC)
+	end := start.Add(time.Hour)
+
+	tests := []struct {
+		name      string
+		overrides map[string]any
+		want      map[string]*kubemodel.Device
+	}{
+		{
+			name:      "no data returns empty device map",
+			overrides: map[string]any{},
+			want:      map[string]*kubemodel.Device{},
+		},
+		{
+			name: "basic device info and uptime",
+			overrides: map[string]any{
+				source.QueryDCGMDeviceInfo: []*source.DCGMDeviceInfoResult{
+					{UUID: "GPU-abc123", Device: "nvidia0", ModelName: "A100"},
+				},
+				source.QueryDCGMDeviceUptime: []*source.DCGMDeviceUptimeResult{
+					{UUID: "GPU-abc123", First: start, Last: end},
+				},
+			},
+			want: map[string]*kubemodel.Device{
+				"GPU-abc123": {
+					UUID:      "GPU-abc123",
+					Device:    "nvidia0",
+					ModelName: "A100",
+					Start:     start,
+					End:       end,
+				},
+			},
+		},
+		{
+			name: "device without uptime is not registered",
+			overrides: map[string]any{
+				source.QueryDCGMDeviceInfo: []*source.DCGMDeviceInfoResult{
+					{UUID: "GPU-abc123", Device: "nvidia0", ModelName: "A100"},
+				},
+			},
+			want: map[string]*kubemodel.Device{},
+		},
+		{
+			name: "device with empty uuid is skipped",
+			overrides: map[string]any{
+				source.QueryDCGMDeviceInfo: []*source.DCGMDeviceInfoResult{
+					{UUID: "", Device: "nvidia0", ModelName: "A100"},
+				},
+				source.QueryDCGMDeviceUptime: []*source.DCGMDeviceUptimeResult{
+					{UUID: "GPU-abc123", First: start, Last: end},
+				},
+			},
+			want: map[string]*kubemodel.Device{},
+		},
+		{
+			name: "duplicate device info entries use first occurrence",
+			overrides: map[string]any{
+				source.QueryDCGMDeviceInfo: []*source.DCGMDeviceInfoResult{
+					{UUID: "GPU-abc123", Device: "nvidia0", ModelName: "A100"},
+					{UUID: "GPU-abc123", Device: "nvidia0-dup", ModelName: "A100-dup"},
+				},
+				source.QueryDCGMDeviceUptime: []*source.DCGMDeviceUptimeResult{
+					{UUID: "GPU-abc123", First: start, Last: end},
+				},
+			},
+			want: map[string]*kubemodel.Device{
+				"GPU-abc123": {
+					UUID:      "GPU-abc123",
+					Device:    "nvidia0",
+					ModelName: "A100",
+					Start:     start,
+					End:       end,
+				},
+			},
+		},
+	}
+
+	for _, tt := range tests {
+		t.Run(tt.name, func(t *testing.T) {
+			ds := source.NewMockOpenCostDataSource()
+			ds.ResolutionValue = 5 * time.Minute
+			seedCluster(ds, start, end)
+			for method, result := range tt.overrides {
+				ds.Querier.SetOverride(method, result)
+			}
+
+			km, err := NewKubeModel(testClusterUID, false, ds)
+			require.NoError(t, err)
+
+			kms := kubemodel.NewKubeModelSet(start, end)
+
+			err = km.computeDevices(kms, start, end)
+			require.NoError(t, err)
+
+			assert.Equal(t, tt.want, kms.Devices)
+		})
+	}
+}

+ 1 - 1
core/pkg/compute/kubemodel/kubemodel.go

@@ -104,6 +104,6 @@ func (km *KubeModel) computeFuncs(start, end time.Time) []computeFunc {
 		km.computePersistentVolumeClaims,
 		km.computePods,
 		km.computeContainers,
-		//km.computeDCGMDevices,
+		km.computeDevices,
 	}
 }

+ 1 - 1
core/pkg/compute/kubemodel/kubemodel_test.go

@@ -185,7 +185,7 @@ func TestComputeKubeModelSet(t *testing.T) {
 				assert.NotEmpty(t, kms.Services)
 				assert.NotEmpty(t, kms.PersistentVolumes)
 				assert.NotEmpty(t, kms.PersistentVolumeClaims)
-				//assert.NotEmpty(t, kms.DCGMDevices)
+				assert.NotEmpty(t, kms.Devices)
 			},
 		},
 	}

+ 1 - 1
core/pkg/model/kubemodel/bingen.go

@@ -22,4 +22,4 @@ package kubemodel
 
 // @bingen:define[string]:github.com/opencost/opencost/core/pkg/cloud.Provider
 
-//go:generate bingen -package=kubemodel -version=2
+//go:generate bingen -package=kubemodel -version=3

+ 28 - 14
core/pkg/model/kubemodel/container.go

@@ -7,22 +7,28 @@ import (
 
 // @bingen:generate:Container
 type Container struct {
-	PodUID                string             `json:"podUid"`
-	Name                  string             `json:"name"`
-	ResourceRequests      ResourceQuantities `json:"resourceRequests"`
-	ResourceLimits        ResourceQuantities `json:"resourceLimits"`
-	CPUCoreAllocationAvg  float64            `json:"cpuCoreAllocationAvg"`
-	CPUCoreUsageAvg       float64            `json:"cpuCoreUsageAvg"`
-	CPUCoreUsageMax       float64            `json:"cpuCoreUsageMax"`
-	RAMBytesAllocationAvg float64            `json:"ramBytesAllocationAvg"`
-	RAMBytesUsageAvg      float64            `json:"ramBytesUsageAvg"`
-	RAMBytesUsageMax      float64            `json:"ramBytesUsageMax"`
-	Start                 time.Time          `json:"start"`
-	End                   time.Time          `json:"end"`
+	PodUID                string                 `json:"podUid"`
+	Name                  string                 `json:"name"`
+	ResourceRequests      ResourceQuantities     `json:"resourceRequests"`
+	ResourceLimits        ResourceQuantities     `json:"resourceLimits"`
+	CPUCoreAllocationAvg  float64                `json:"cpuCoreAllocationAvg"`
+	CPUCoreUsageAvg       float64                `json:"cpuCoreUsageAvg"`
+	CPUCoreUsageMax       float64                `json:"cpuCoreUsageMax"`
+	RAMBytesAllocationAvg float64                `json:"ramBytesAllocationAvg"`
+	RAMBytesUsageAvg      float64                `json:"ramBytesUsageAvg"`
+	RAMBytesUsageMax      float64                `json:"ramBytesUsageMax"`
+	DeviceUsages          map[string]DeviceUsage `json:"deviceUsages"` // @bingen:field[version=3]
+	Start                 time.Time              `json:"start"`
+	End                   time.Time              `json:"end"`
 }
 
-func (c *Container) GetKey() string {
-	return fmt.Sprintf("%s/%s", c.PodUID, c.Name)
+// DeviceUsage holds usage metrics for a single container/device pairing. The shape is
+// vendor-agnostic, but the only populating source currently implemented is the DCGM exporter.
+// It is keyed by Device.UUID under Container.DeviceUsages.
+// @bingen:generate:DeviceUsage
+type DeviceUsage struct {
+	UsageAvg float64 `json:"usageAvg"`
+	UsageMax float64 `json:"usageMax"`
 }
 
 func (c *Container) ValidateContainer(window Window) error {
@@ -56,3 +62,11 @@ func (kms *KubeModelSet) RegisterContainer(container *Container) error {
 
 	return nil
 }
+
+func (c *Container) GetKey() string {
+	return ContainerKey(c.PodUID, c.Name)
+}
+
+func ContainerKey(podUID, containerName string) string {
+	return fmt.Sprintf("%s/%s", podUID, containerName)
+}

+ 8 - 8
core/pkg/model/kubemodel/daemonset.go

@@ -8,14 +8,14 @@ import (
 // @bingen:generate:DaemonSet
 // DaemonSet represents a Kubernetes DaemonSet resource
 type DaemonSet struct {
-	UID              string            `json:"uid"`
-	NamespaceUID     string            `json:"namespaceUid"`
-	Name             string            `json:"name"`
-	Labels           map[string]string `json:"labels,omitempty"`
-	Annotations      map[string]string `json:"annotations,omitempty"`
-	DevicePluginInfo map[string]string `json:"devicePluginInfo"` // bingen:field[ignore]
-	Start            time.Time         `json:"start,omitempty"`
-	End              time.Time         `json:"end,omitempty"`
+	UID          string            `json:"uid"`
+	NamespaceUID string            `json:"namespaceUid"`
+	Name         string            `json:"name"`
+	Labels       map[string]string `json:"labels,omitempty"`
+	Annotations  map[string]string `json:"annotations,omitempty"`
+	Arguments    map[string]string `json:"arguments,omitempty"`
+	Start        time.Time         `json:"start,omitempty"`
+	End          time.Time         `json:"end,omitempty"`
 }
 
 func (d *DaemonSet) ValidateDaemonSet(window Window) error {

+ 0 - 66
core/pkg/model/kubemodel/dcgm.go

@@ -1,66 +0,0 @@
-package kubemodel
-
-import (
-	"fmt"
-	"time"
-)
-
-// DCGMDevice holds recording from the DCGM exporter which provides identification and usage metrics for
-// Nvidia gpu. These Nvidia devices can be incorporated into the cluster via k8s Device Plugin API or DRAs.
-// While the DCGM exporter does provide unique identifiers for the containers that it is reporting metrics on,
-// It is split out here to provide some isolation from the rest of the KubeModel which represent universal structures
-// from the k8s API. It is left to the end user to interpret the relationships to the rest of the cluster based on
-// container unique identifiers
-// @bingen:generate:DCGMDevice
-type DCGMDevice struct {
-	UUID      string             `json:"uuid"`
-	Start     time.Time          `json:"start"`
-	End       time.Time          `json:"end"`
-	Device    string             `json:"device"`
-	ModelName string             `json:"modelName"`
-	PodUsages map[string]DCGMPod `json:"podUsages"`
-}
-
-// @bingen:generate:DCGMPod
-type DCGMPod struct {
-	ContainerUsages map[string]DCGMContainer `json:"container-usages"`
-}
-
-// @bingen:generate:DCGMContainer
-type DCGMContainer struct {
-	UsageAvg float64 `json:"usageAvg"`
-	UsageMax float64 `json:"usageMax"`
-}
-
-func (d *DCGMDevice) ValidateDCGMDevice(window Window) error {
-	if d.UUID == "" {
-		return fmt.Errorf("UUID is missing for DCGMDevice with device '%s'", d.Device)
-	}
-
-	if err := checkWindow(window, d.Start, d.End); err != nil {
-		return err
-	}
-
-	return nil
-}
-
-// RegisterDCGMDevice validates and adds a DCGMDevice to the set, keyed by UUID.
-func (kms *KubeModelSet) RegisterDCGMDevice(device *DCGMDevice) error {
-	if err := device.ValidateDCGMDevice(kms.Window); err != nil {
-		err = fmt.Errorf("RegisterDCGMDevice: invalid dcgm device: %w", err)
-		kms.Error(err)
-		return err
-	}
-
-	if _, ok := kms.DCGMDevices[device.UUID]; !ok {
-		if kms.Cluster == nil {
-			kms.Warnf("RegisterDCGMDevice: Cluster is nil")
-		}
-
-		kms.DCGMDevices[device.UUID] = device
-
-		kms.Metadata.ObjectCount++
-	}
-
-	return nil
-}

+ 52 - 0
core/pkg/model/kubemodel/device.go

@@ -0,0 +1,52 @@
+package kubemodel
+
+import (
+	"fmt"
+	"time"
+)
+
+// Device holds identification for an accelerator device (e.g. an Nvidia GPU) attached to the
+// cluster via the k8s Device Plugin API or DRAs. The shape is vendor-agnostic, but the only
+// populating source currently implemented is the DCGM exporter. Usage of a Device by a specific
+// container is recorded on Container.DeviceUsages, keyed by Device.UUID.
+// @bingen:generate:Device
+type Device struct {
+	UUID      string    `json:"uuid"`
+	Start     time.Time `json:"start"`
+	End       time.Time `json:"end"`
+	Device    string    `json:"device"`
+	ModelName string    `json:"modelName"`
+}
+
+func (d *Device) ValidateDevice(window Window) error {
+	if d.UUID == "" {
+		return fmt.Errorf("UUID is missing for Device with device '%s'", d.Device)
+	}
+
+	if err := checkWindow(window, d.Start, d.End); err != nil {
+		return err
+	}
+
+	return nil
+}
+
+// RegisterDevice validates and adds a Device to the set, keyed by UUID.
+func (kms *KubeModelSet) RegisterDevice(device *Device) error {
+	if err := device.ValidateDevice(kms.Window); err != nil {
+		err = fmt.Errorf("RegisterDevice: invalid device: %w", err)
+		kms.Error(err)
+		return err
+	}
+
+	if _, ok := kms.Devices[device.UUID]; !ok {
+		if kms.Cluster == nil {
+			kms.Warnf("RegisterDevice: Cluster is nil")
+		}
+
+		kms.Devices[device.UUID] = device
+
+		kms.Metadata.ObjectCount++
+	}
+
+	return nil
+}

+ 20 - 20
core/pkg/model/kubemodel/dcgm_test.go → core/pkg/model/kubemodel/device_test.go

@@ -7,35 +7,35 @@ import (
 	"github.com/stretchr/testify/require"
 )
 
-func TestValidateDCGMDevice(t *testing.T) {
+func TestValidateDevice(t *testing.T) {
 	start := time.Now().UTC().Truncate(time.Hour)
 	end := start.Add(time.Hour)
 	window := Window{Start: start, End: end}
 
 	tests := []struct {
 		name    string
-		device  *DCGMDevice
+		device  *Device
 		wantErr string
 	}{
 		{
 			name:    "empty UUID",
-			device:  &DCGMDevice{Device: "GPU-0", Start: start, End: end},
-			wantErr: "UUID is missing for DCGMDevice with device 'GPU-0'",
+			device:  &Device{Device: "GPU-0", Start: start, End: end},
+			wantErr: "UUID is missing for Device with device 'GPU-0'",
 		},
 		{
 			name:    "outside window",
-			device:  &DCGMDevice{UUID: "gpu-uuid", Device: "GPU-0", Start: start.Add(-time.Hour), End: end},
+			device:  &Device{UUID: "gpu-uuid", Device: "GPU-0", Start: start.Add(-time.Hour), End: end},
 			wantErr: checkWindow(window, start.Add(-time.Hour), end).Error(),
 		},
 		{
 			name:   "valid",
-			device: &DCGMDevice{UUID: "gpu-uuid", Device: "GPU-0", Start: start, End: end},
+			device: &Device{UUID: "gpu-uuid", Device: "GPU-0", Start: start, End: end},
 		},
 	}
 
 	for _, tt := range tests {
 		t.Run(tt.name, func(t *testing.T) {
-			err := tt.device.ValidateDCGMDevice(window)
+			err := tt.device.ValidateDevice(window)
 			if tt.wantErr != "" {
 				require.EqualError(t, err, tt.wantErr)
 			} else {
@@ -45,12 +45,12 @@ func TestValidateDCGMDevice(t *testing.T) {
 	}
 }
 
-func TestRegisterDCGMDevice(t *testing.T) {
+func TestRegisterDevice(t *testing.T) {
 	start := time.Now().UTC().Truncate(time.Hour)
 	end := start.Add(time.Hour)
 
-	newDevice := func(uuid, device string) *DCGMDevice {
-		return &DCGMDevice{UUID: uuid, Device: device, Start: start, End: end}
+	newDevice := func(uuid, device string) *Device {
+		return &Device{UUID: uuid, Device: device, Start: start, End: end}
 	}
 	withCluster := func(kms *KubeModelSet) {
 		kms.RegisterCluster(&Cluster{UID: "cluster-uid", Start: start, End: end})
@@ -59,18 +59,18 @@ func TestRegisterDCGMDevice(t *testing.T) {
 	tests := []struct {
 		name    string
 		setup   func(*KubeModelSet)
-		device  *DCGMDevice
+		device  *Device
 		wantErr string
 		want    *KubeModelSet
 	}{
 		{
 			name:    "validation failure",
-			device:  &DCGMDevice{UUID: "", Device: "GPU-0", Start: start, End: end},
-			wantErr: "RegisterDCGMDevice: invalid dcgm device: UUID is missing for DCGMDevice with device 'GPU-0'",
+			device:  &Device{UUID: "", Device: "GPU-0", Start: start, End: end},
+			wantErr: "RegisterDevice: invalid device: UUID is missing for Device with device 'GPU-0'",
 			want: func() *KubeModelSet {
 				kms := NewKubeModelSet(start, end)
 				kms.Metadata.Diagnostics = []Diagnostic{
-					{Level: DiagnosticLevelError, Message: "RegisterDCGMDevice: invalid dcgm device: UUID is missing for DCGMDevice with device 'GPU-0'"},
+					{Level: DiagnosticLevelError, Message: "RegisterDevice: invalid device: UUID is missing for Device with device 'GPU-0'"},
 				}
 				return kms
 			}(),
@@ -80,10 +80,10 @@ func TestRegisterDCGMDevice(t *testing.T) {
 			device: newDevice("gpu-uuid", "GPU-0"),
 			want: func() *KubeModelSet {
 				kms := NewKubeModelSet(start, end)
-				kms.DCGMDevices["gpu-uuid"] = newDevice("gpu-uuid", "GPU-0")
+				kms.Devices["gpu-uuid"] = newDevice("gpu-uuid", "GPU-0")
 				kms.Metadata.ObjectCount = 1
 				kms.Metadata.Diagnostics = []Diagnostic{
-					{Level: DiagnosticLevelWarning, Message: "RegisterDCGMDevice: Cluster is nil"},
+					{Level: DiagnosticLevelWarning, Message: "RegisterDevice: Cluster is nil"},
 				}
 				return kms
 			}(),
@@ -95,7 +95,7 @@ func TestRegisterDCGMDevice(t *testing.T) {
 			want: func() *KubeModelSet {
 				kms := NewKubeModelSet(start, end)
 				withCluster(kms)
-				kms.DCGMDevices["gpu-uuid"] = newDevice("gpu-uuid", "GPU-0")
+				kms.Devices["gpu-uuid"] = newDevice("gpu-uuid", "GPU-0")
 				kms.Metadata.ObjectCount = 1
 				return kms
 			}(),
@@ -104,13 +104,13 @@ func TestRegisterDCGMDevice(t *testing.T) {
 			name: "duplicate registration is a no-op",
 			setup: func(kms *KubeModelSet) {
 				withCluster(kms)
-				kms.RegisterDCGMDevice(newDevice("gpu-uuid", "GPU-0"))
+				kms.RegisterDevice(newDevice("gpu-uuid", "GPU-0"))
 			},
 			device: newDevice("gpu-uuid", "GPU-1"),
 			want: func() *KubeModelSet {
 				kms := NewKubeModelSet(start, end)
 				withCluster(kms)
-				kms.DCGMDevices["gpu-uuid"] = newDevice("gpu-uuid", "GPU-0")
+				kms.Devices["gpu-uuid"] = newDevice("gpu-uuid", "GPU-0")
 				kms.Metadata.ObjectCount = 1
 				return kms
 			}(),
@@ -124,7 +124,7 @@ func TestRegisterDCGMDevice(t *testing.T) {
 				tt.setup(kms)
 			}
 
-			err := kms.RegisterDCGMDevice(tt.device)
+			err := kms.RegisterDevice(tt.device)
 
 			if tt.wantErr != "" {
 				require.EqualError(t, err, tt.wantErr)

+ 3 - 3
core/pkg/model/kubemodel/kubemodel.go

@@ -24,7 +24,7 @@ type KubeModelSet struct {
 	PersistentVolumeClaims map[string]*PersistentVolumeClaim `json:"pvcs"`              // @bingen:field[version=2]
 	Pods                   map[string]*Pod                   `json:"pods"`              // @bingen:field[version=2]
 	Containers             map[string]*Container             `json:"containers"`        // @bingen:field[version=2]
-	DCGMDevices            map[string]*DCGMDevice            `json:"dcgmDevices"`       // @bingen:field[ignore]
+	Devices                map[string]*Device                `json:"devices"`           // @bingen:field[version=3]
 }
 
 func NewKubeModelSet(start time.Time, end time.Time) *KubeModelSet {
@@ -48,7 +48,7 @@ func NewKubeModelSet(start time.Time, end time.Time) *KubeModelSet {
 		ReplicaSets:            map[string]*ReplicaSet{},
 		Namespaces:             map[string]*Namespace{},
 		Nodes:                  map[string]*Node{},
-		DCGMDevices:            map[string]*DCGMDevice{},
+		Devices:                map[string]*Device{},
 		Pods:                   map[string]*Pod{},
 		PersistentVolumeClaims: map[string]*PersistentVolumeClaim{},
 		ResourceQuotas:         map[string]*ResourceQuota{},
@@ -74,7 +74,7 @@ func (kms *KubeModelSet) IsEmpty() bool {
 		len(kms.ReplicaSets) == 0 &&
 		len(kms.Namespaces) == 0 &&
 		len(kms.Nodes) == 0 &&
-		len(kms.DCGMDevices) == 0 &&
+		len(kms.Devices) == 0 &&
 		len(kms.Pods) == 0 &&
 		len(kms.PersistentVolumeClaims) == 0 &&
 		len(kms.ResourceQuotas) == 0 &&

Разница между файлами не показана из-за своего большого размера
+ 358 - 447
core/pkg/model/kubemodel/kubemodel_codecs.go


+ 1 - 2
core/pkg/model/kubemodel/kubemodel_helpers_test.go

@@ -59,6 +59,5 @@ func KubeModelSetEquals(t *testing.T, this, that *KubeModelSet) {
 	require.Equal(t, this.PersistentVolumeClaims, that.PersistentVolumeClaims)
 	require.Equal(t, this.Services, that.Services)
 	require.Equal(t, this.PersistentVolumes, that.PersistentVolumes)
-	// DCGM is ignored by bingen
-	// require.Equal(t, this.DCGMDevices, that.DCGMDevices)
+	require.Equal(t, this.Devices, that.Devices)
 }

+ 10 - 11
core/pkg/model/kubemodel/mock.go

@@ -162,6 +162,7 @@ func NewMockKubeModelSet(start, end time.Time) *KubeModelSet {
 		Name:         "my-daemonset",
 		Labels:       map[string]string{"app": "my-daemonset"},
 		Annotations:  map[string]string{"note": "test"},
+		Arguments:    map[string]string{"vgpu": "2"},
 		Start:        start,
 		End:          end,
 	})
@@ -236,21 +237,19 @@ func NewMockKubeModelSet(start, end time.Time) *KubeModelSet {
 		End:             end,
 	})
 
-	// --- DCGMDevice ---
-	kms.RegisterDCGMDevice(&DCGMDevice{
+	// --- Device ---
+	kms.RegisterDevice(&Device{
 		UUID:      "GPU-abc123def-456-789",
 		Device:    "0",
 		ModelName: "Tesla T4",
-		PodUsages: map[string]DCGMPod{
-			"pod-uid": {
-				ContainerUsages: map[string]DCGMContainer{
-					"app": {UsageAvg: 0.65, UsageMax: 0.92},
-				},
-			},
-		},
-		Start: start,
-		End:   end,
+		Start:     start,
+		End:       end,
 	})
+	if c, ok := kms.Containers["pod-uid/app"]; ok {
+		c.DeviceUsages = map[string]DeviceUsage{
+			"GPU-abc123def-456-789": {UsageAvg: 0.65, UsageMax: 0.92},
+		}
+	}
 
 	// --- Diagnostics ---
 	kms.Error(errMock("mock error"))

+ 2 - 0
core/pkg/source/datasource.go

@@ -130,6 +130,7 @@ const (
 	QueryDaemonSetUptime      = "QueryDaemonSetUptime"
 	QueryDaemonSetLabels      = "QueryDaemonSetLabels"
 	QueryDaemonSetAnnotations = "QueryDaemonSetAnnotations"
+	QueryDaemonSetArguments   = "QueryDaemonSetArguments"
 
 	// Job
 	QueryJobInfo        = "QueryJobInfo"
@@ -352,6 +353,7 @@ type MetricsQuerier interface {
 	QueryDaemonSetUptime(start, end time.Time) *Future[UptimeResult]
 	QueryDaemonSetLabels(start, end time.Time) *Future[LabelsResult]
 	QueryDaemonSetAnnotations(start, end time.Time) *Future[AnnotationsResult]
+	QueryDaemonSetArguments(start, end time.Time) *Future[DaemonSetArgumentResult]
 
 	// Job
 	QueryJobInfo(start, end time.Time) *Future[JobInfoResult]

+ 25 - 0
core/pkg/source/decoders.go

@@ -63,6 +63,8 @@ const (
 	SameRegionLabel      = "same_region"
 	NatGatewayLabel      = "nat_gateway"
 	KubeModelVersion     = "kubemodel_version"
+	ArgLabel             = "arg"
+	ValueLabel           = "value"
 )
 
 const (
@@ -1804,6 +1806,29 @@ func DecodeDaemonSetInfoResult(result *QueryResult) *DaemonSetInfoResult {
 	}
 }
 
+// DaemonSetArgumentResult represents a single "--key=value" container argument parsed off a
+// DaemonSet. One row is emitted per argument, so a DaemonSet with N arguments produces N results.
+type DaemonSetArgumentResult struct {
+	UID     string
+	Cluster string
+	Arg     string
+	Value   string
+}
+
+func DecodeDaemonSetArgumentResult(result *QueryResult) *DaemonSetArgumentResult {
+	uid, _ := result.GetString(UIDLabel)
+	cluster, _ := result.GetCluster()
+	arg, _ := result.GetString(ArgLabel)
+	value, _ := result.GetString(ValueLabel)
+
+	return &DaemonSetArgumentResult{
+		UID:     uid,
+		Cluster: cluster,
+		Arg:     arg,
+		Value:   value,
+	}
+}
+
 type JobInfoResult struct {
 	UID          string
 	Cluster      string

+ 6 - 0
core/pkg/source/mock.go

@@ -626,6 +626,12 @@ func (m *MockMetricsQuerier) QueryDaemonSetAnnotations(start, end time.Time) *Fu
 	})
 }
 
+func (m *MockMetricsQuerier) QueryDaemonSetArguments(start, end time.Time) *Future[DaemonSetArgumentResult] {
+	return getFutureFromOverride(m.overrides, QueryDaemonSetArguments, func() *Future[DaemonSetArgumentResult] {
+		return m.noop.QueryDaemonSetArguments(start, end)
+	})
+}
+
 // Job
 
 func (m *MockMetricsQuerier) QueryJobInfo(start, end time.Time) *Future[JobInfoResult] {

+ 4 - 0
core/pkg/source/noop.go

@@ -411,6 +411,10 @@ func (m *NoOpMetricsQuerier) QueryDaemonSetAnnotations(start, end time.Time) *Fu
 	return newEmptyResult(DecodeAnnotationsResult)
 }
 
+func (m *NoOpMetricsQuerier) QueryDaemonSetArguments(start, end time.Time) *Future[DaemonSetArgumentResult] {
+	return newEmptyResult(DecodeDaemonSetArgumentResult)
+}
+
 // Job
 
 func (m *NoOpMetricsQuerier) QueryJobInfo(start, end time.Time) *Future[JobInfoResult] {

+ 5 - 0
core/pkg/source/record.go

@@ -511,6 +511,11 @@ func (m *RecordMetricsQuerier) QueryDaemonSetAnnotations(start, end time.Time) *
 	return m.Querier.QueryDaemonSetAnnotations(start, end)
 }
 
+func (m *RecordMetricsQuerier) QueryDaemonSetArguments(start, end time.Time) *Future[DaemonSetArgumentResult] {
+	m.recordCall(QueryDaemonSetArguments)
+	return m.Querier.QueryDaemonSetArguments(start, end)
+}
+
 // Job
 
 func (m *RecordMetricsQuerier) QueryJobInfo(start, end time.Time) *Future[JobInfoResult] {

+ 2 - 0
core/pkg/source/record_test.go

@@ -103,6 +103,7 @@ func TestRecordMetricsQuerier_Counts(t *testing.T) {
 	r.QueryDaemonSetUptime(start, end)
 	r.QueryDaemonSetLabels(start, end)
 	r.QueryDaemonSetAnnotations(start, end)
+	r.QueryDaemonSetArguments(start, end)
 	r.QueryJobInfo(start, end)
 	r.QueryJobUptime(start, end)
 	r.QueryJobLabels(start, end)
@@ -259,6 +260,7 @@ func TestRecordMetricsQuerier_Counts(t *testing.T) {
 		"QueryDaemonSetUptime":                          1,
 		"QueryDaemonSetLabels":                          1,
 		"QueryDaemonSetAnnotations":                     1,
+		"QueryDaemonSetArguments":                       1,
 		"QueryJobInfo":                                  1,
 		"QueryJobUptime":                                1,
 		"QueryJobLabels":                                1,

+ 25 - 0
core/pkg/util/args.go

@@ -0,0 +1,25 @@
+package util
+
+import (
+	"strings"
+
+	v1 "k8s.io/api/core/v1"
+)
+
+// ParseContainerArgs extracts "--key=value" style arguments from the given containers into a
+// flat map. Bare flags without a "=" (e.g. "--verbose") are recorded with an empty value. Args
+// are merged across all containers; on key collision, the last container/arg seen wins.
+func ParseContainerArgs(containers []v1.Container) map[string]string {
+	args := make(map[string]string)
+	for _, container := range containers {
+		for _, arg := range container.Args {
+			trimmed := strings.TrimLeft(arg, "-")
+			if trimmed == "" {
+				continue
+			}
+			key, value, _ := strings.Cut(trimmed, "=")
+			args[key] = value
+		}
+	}
+	return args
+}

+ 14 - 0
modules/collector-source/pkg/collector/collector.go

@@ -121,6 +121,7 @@ func NewOpenCostMetricStore() metric.MetricStore {
 	memStore.Register(NewDaemonSetUptimeMetricCollector())
 	memStore.Register(NewDaemonSetLabelsMetricCollector())
 	memStore.Register(NewDaemonSetAnnotationsMetricCollector())
+	memStore.Register(NewDaemonSetArgumentsMetricCollector())
 	memStore.Register(NewJobInfoMetricCollector())
 	memStore.Register(NewJobUptimeMetricCollector())
 	memStore.Register(NewJobLabelsMetricCollector())
@@ -2528,6 +2529,19 @@ func NewDaemonSetAnnotationsMetricCollector() *metric.MetricCollector {
 	)
 }
 
+func NewDaemonSetArgumentsMetricCollector() *metric.MetricCollector {
+	return metric.NewMetricCollector(
+		metric.DaemonSetArgumentsID,
+		metric.DaemonSetArguments,
+		[]string{
+			source.UIDLabel,
+			source.ArgLabel,
+		},
+		aggregator.Info,
+		nil,
+	)
+}
+
 func NewJobInfoMetricCollector() *metric.MetricCollector {
 	return metric.NewMetricCollector(
 		metric.JobInfoID,

+ 4 - 0
modules/collector-source/pkg/collector/metricsquerier.go

@@ -601,6 +601,10 @@ func (c *collectorMetricsQuerier) QueryDaemonSetAnnotations(start, end time.Time
 	return queryCollector(c, start, end, metric.DaemonSetAnnotationsID, source.DecodeAnnotationsResult)
 }
 
+func (c *collectorMetricsQuerier) QueryDaemonSetArguments(start, end time.Time) *source.Future[source.DaemonSetArgumentResult] {
+	return queryCollector(c, start, end, metric.DaemonSetArgumentsID, source.DecodeDaemonSetArgumentResult)
+}
+
 func (c *collectorMetricsQuerier) QueryJobInfo(start, end time.Time) *source.Future[source.JobInfoResult] {
 	return queryCollector(c, start, end, metric.JobInfoID, source.DecodeJobInfoResult)
 }

+ 1 - 0
modules/collector-source/pkg/metric/collector.go

@@ -131,6 +131,7 @@ const (
 	DaemonSetUptimeID                          MetricCollectorID = "DaemonSetUptime"
 	DaemonSetLabelsID                          MetricCollectorID = "DaemonSetLabels"
 	DaemonSetAnnotationsID                     MetricCollectorID = "DaemonSetAnnotations"
+	DaemonSetArgumentsID                       MetricCollectorID = "DaemonSetArguments"
 	JobInfoID                                  MetricCollectorID = "JobInfo"
 	JobUptimeID                                MetricCollectorID = "JobUptime"
 	JobLabelsID                                MetricCollectorID = "JobLabels"

+ 1 - 0
modules/collector-source/pkg/metric/metrics.go

@@ -34,6 +34,7 @@ const (
 	DaemonSetInfo                                         = "daemonset_info"
 	DaemonSetLabels                                       = "daemonset_labels"
 	DaemonSetAnnotations                                  = "daemonset_annotations"
+	DaemonSetArguments                                    = "daemonset_arguments"
 	JobInfo                                               = "job_info"
 	JobLabels                                             = "job_labels"
 	JobAnnotations                                        = "job_annotations"

+ 21 - 0
modules/collector-source/pkg/scrape/clustercache.go

@@ -910,6 +910,27 @@ func (ccs *ClusterCacheScraper) scrapeDaemonSets(daemonSets []*clustercache.Daem
 			Value:          0,
 			AdditionalInfo: daemonSetAnnotations,
 		})
+
+		// daemonSet arguments
+		daemonSetArguments := coreutil.ParseContainerArgs(daemonSet.SpecContainers)
+		argKeys := maps.Keys(daemonSetArguments)
+		slices.Sort(argKeys)
+		for _, arg := range argKeys {
+			value := daemonSetArguments[arg]
+			argLabels := map[string]string{
+				source.UIDLabel:          string(daemonSet.UID),
+				source.NamespaceUIDLabel: string(nsUID),
+				source.DaemonSetLabel:    daemonSet.Name,
+				source.ArgLabel:          arg,
+				source.ValueLabel:        value,
+			}
+			scrapeResults = append(scrapeResults, metric.Update{
+				Name:           metric.DaemonSetArguments,
+				Labels:         argLabels,
+				Value:          0,
+				AdditionalInfo: argLabels,
+			})
+		}
 	}
 
 	events.Dispatch(event.ScrapeEvent{

+ 90 - 0
modules/collector-source/pkg/scrape/clustercache_test.go

@@ -2393,6 +2393,96 @@ func Test_kubernetesScraper_scrapeDaemonSets(t *testing.T) {
 				},
 			},
 		},
+		{
+			name: "with container arguments",
+			scrapes: []scrape{
+				{
+					DaemonSets: []*clustercache.DaemonSet{
+						{
+							Name:      "daemonSet1",
+							Namespace: "namespace1",
+							UID:       "uuid1",
+							SpecContainers: []v1.Container{
+								{Args: []string{"--vgpu=2", "--bare-flag"}},
+							},
+						},
+					},
+					Timestamp: start1,
+				},
+			},
+			expected: []metric.Update{
+				{
+					Name: metric.DaemonSetInfo,
+					Labels: map[string]string{
+						source.UIDLabel:          "uuid1",
+						source.NamespaceUIDLabel: "",
+						source.DaemonSetLabel:    "daemonSet1",
+					},
+					Value: 0,
+					AdditionalInfo: map[string]string{
+						source.UIDLabel:          "uuid1",
+						source.NamespaceUIDLabel: "",
+						source.DaemonSetLabel:    "daemonSet1",
+					},
+				},
+				{
+					Name: metric.DaemonSetLabels,
+					Labels: map[string]string{
+						source.UIDLabel:          "uuid1",
+						source.NamespaceUIDLabel: "",
+						source.DaemonSetLabel:    "daemonSet1",
+					},
+					Value:          0,
+					AdditionalInfo: map[string]string{},
+				},
+				{
+					Name: metric.DaemonSetAnnotations,
+					Labels: map[string]string{
+						source.UIDLabel:          "uuid1",
+						source.NamespaceUIDLabel: "",
+						source.DaemonSetLabel:    "daemonSet1",
+					},
+					Value:          0,
+					AdditionalInfo: map[string]string{},
+				},
+				{
+					Name: metric.DaemonSetArguments,
+					Labels: map[string]string{
+						source.UIDLabel:          "uuid1",
+						source.NamespaceUIDLabel: "",
+						source.DaemonSetLabel:    "daemonSet1",
+						source.ArgLabel:          "bare-flag",
+						source.ValueLabel:        "",
+					},
+					Value: 0,
+					AdditionalInfo: map[string]string{
+						source.UIDLabel:          "uuid1",
+						source.NamespaceUIDLabel: "",
+						source.DaemonSetLabel:    "daemonSet1",
+						source.ArgLabel:          "bare-flag",
+						source.ValueLabel:        "",
+					},
+				},
+				{
+					Name: metric.DaemonSetArguments,
+					Labels: map[string]string{
+						source.UIDLabel:          "uuid1",
+						source.NamespaceUIDLabel: "",
+						source.DaemonSetLabel:    "daemonSet1",
+						source.ArgLabel:          "vgpu",
+						source.ValueLabel:        "2",
+					},
+					Value: 0,
+					AdditionalInfo: map[string]string{
+						source.UIDLabel:          "uuid1",
+						source.NamespaceUIDLabel: "",
+						source.DaemonSetLabel:    "daemonSet1",
+						source.ArgLabel:          "vgpu",
+						source.ValueLabel:        "2",
+					},
+				},
+			},
+		},
 	}
 
 	for _, tt := range tests {

+ 27 - 3
modules/collector-source/pkg/scrape/dcgm.go

@@ -6,6 +6,7 @@ import (
 
 	"github.com/opencost/opencost/core/pkg/clustercache"
 	"github.com/opencost/opencost/core/pkg/log"
+	"github.com/opencost/opencost/core/pkg/source"
 	"github.com/opencost/opencost/modules/collector-source/pkg/event"
 	"github.com/opencost/opencost/modules/collector-source/pkg/metric"
 	"github.com/opencost/opencost/modules/collector-source/pkg/scrape/target"
@@ -16,10 +17,10 @@ var dcgmRegex = regexp.MustCompile("(?i)(.*dcgm-exporter.*)")
 
 func newDCGMScrapper(clusterCache clustercache.ClusterCache) Scraper {
 	tp := newDCGMTargetProvider(clusterCache)
-	return newDCGMTargetScraper(tp)
+	return newDCGMTargetScraper(tp, podUIDEnricher(clusterCache))
 }
 
-func newDCGMTargetScraper(provider target.TargetProvider) *TargetScraper {
+func newDCGMTargetScraper(provider target.TargetProvider, enrich UpdateEnricher) *TargetScraper {
 	return newTargetScrapper(
 		event.DCGMScraperName,
 		provider,
@@ -27,7 +28,30 @@ func newDCGMTargetScraper(provider target.TargetProvider) *TargetScraper {
 			metric.DCGMFIPROFGRENGINEACTIVE,
 			metric.DCGMFIDEVDECUTIL,
 		},
-		true)
+		true,
+		enrich)
+}
+
+// podUIDEnricher backfills pod_uid on a DCGM update using its own namespace/pod name labels,
+// resolved against a freshly built index of the cluster's current pods. Left unset if pod_uid
+// is already present, or if namespace/pod can't be resolved to a known pod.
+func podUIDEnricher(clusterCache clustercache.ClusterCache) UpdateEnricher {
+	return func(updates []metric.Update) {
+		index := buildPodIndex(clusterCache.GetAllPods())
+		for _, update := range updates {
+			if update.Labels[source.PodUIDLabel] != "" {
+				continue
+			}
+			namespace, pod := update.Labels[source.NamespaceLabel], update.Labels[source.PodLabel]
+			if namespace == "" || pod == "" {
+				continue
+			}
+
+			if uid, ok := index[podKey{namespace: namespace, name: pod}]; ok {
+				update.Labels[source.PodUIDLabel] = string(uid)
+			}
+		}
+	}
 }
 
 type DCGMTargetProvider struct {

+ 150 - 0
modules/collector-source/pkg/scrape/dcgm_test.go

@@ -2,6 +2,10 @@ package scrape
 
 import (
 	"testing"
+
+	"github.com/opencost/opencost/core/pkg/clustercache"
+	"github.com/opencost/opencost/core/pkg/source"
+	"github.com/opencost/opencost/modules/collector-source/pkg/metric"
 )
 
 func Test_isDCGM(t *testing.T) {
@@ -62,3 +66,149 @@ func Test_isDCGM(t *testing.T) {
 		})
 	}
 }
+
+func Test_podUIDEnricher(t *testing.T) {
+	cache := &clustercache.MockClusterCache{
+		Pods: []*clustercache.Pod{
+			{UID: "pod-uid-1", Name: "pod1", Namespace: "namespace1"},
+			{UID: "pod-uid-2", Name: "pod2", Namespace: "namespace2"},
+		},
+	}
+
+	tests := map[string]struct {
+		labels map[string]string
+		want   map[string]string
+	}{
+		"resolves pod_uid from namespace and pod": {
+			labels: map[string]string{
+				source.NamespaceLabel: "namespace1",
+				source.PodLabel:       "pod1",
+			},
+			want: map[string]string{
+				source.NamespaceLabel: "namespace1",
+				source.PodLabel:       "pod1",
+				source.PodUIDLabel:    "pod-uid-1",
+			},
+		},
+		"unknown pod is left unset": {
+			labels: map[string]string{
+				source.NamespaceLabel: "namespace1",
+				source.PodLabel:       "unknown-pod",
+			},
+			want: map[string]string{
+				source.NamespaceLabel: "namespace1",
+				source.PodLabel:       "unknown-pod",
+			},
+		},
+		"pod name from wrong namespace is left unset": {
+			labels: map[string]string{
+				source.NamespaceLabel: "namespace2",
+				source.PodLabel:       "pod1",
+			},
+			want: map[string]string{
+				source.NamespaceLabel: "namespace2",
+				source.PodLabel:       "pod1",
+			},
+		},
+		"missing namespace label is left unset": {
+			labels: map[string]string{
+				source.PodLabel: "pod1",
+			},
+			want: map[string]string{
+				source.PodLabel: "pod1",
+			},
+		},
+		"missing pod label is left unset": {
+			labels: map[string]string{
+				source.NamespaceLabel: "namespace1",
+			},
+			want: map[string]string{
+				source.NamespaceLabel: "namespace1",
+			},
+		},
+		"nil labels is left untouched": {
+			labels: nil,
+			want:   nil,
+		},
+		"existing non-empty pod_uid is left untouched": {
+			labels: map[string]string{
+				source.NamespaceLabel: "namespace1",
+				source.PodLabel:       "pod1",
+				source.PodUIDLabel:    "already-set",
+			},
+			want: map[string]string{
+				source.NamespaceLabel: "namespace1",
+				source.PodLabel:       "pod1",
+				source.PodUIDLabel:    "already-set",
+			},
+		},
+	}
+
+	for name, tt := range tests {
+		t.Run(name, func(t *testing.T) {
+			enrich := podUIDEnricher(cache)
+			updates := []metric.Update{{Labels: cloneLabels(tt.labels)}}
+
+			enrich(updates)
+
+			assertLabels(t, updates[0].Labels, tt.want)
+		})
+	}
+
+	// A single call must correctly enrich every update in the batch, not just
+	// the first, since updates are processed together rather than one at a time.
+	t.Run("enriches every update in a multi-update batch", func(t *testing.T) {
+		enrich := podUIDEnricher(cache)
+		updates := []metric.Update{
+			{Labels: map[string]string{source.NamespaceLabel: "namespace1", source.PodLabel: "pod1"}},
+			{Labels: map[string]string{source.NamespaceLabel: "namespace2", source.PodLabel: "pod2"}},
+			{Labels: map[string]string{source.NamespaceLabel: "namespace1", source.PodLabel: "unknown-pod"}},
+		}
+
+		enrich(updates)
+
+		assertLabels(t, updates[0].Labels, map[string]string{
+			source.NamespaceLabel: "namespace1",
+			source.PodLabel:       "pod1",
+			source.PodUIDLabel:    "pod-uid-1",
+		})
+		assertLabels(t, updates[1].Labels, map[string]string{
+			source.NamespaceLabel: "namespace2",
+			source.PodLabel:       "pod2",
+			source.PodUIDLabel:    "pod-uid-2",
+		})
+		assertLabels(t, updates[2].Labels, map[string]string{
+			source.NamespaceLabel: "namespace1",
+			source.PodLabel:       "unknown-pod",
+		})
+	})
+
+	t.Run("empty batch does not panic", func(t *testing.T) {
+		enrich := podUIDEnricher(cache)
+		enrich(nil)
+		enrich([]metric.Update{})
+	})
+}
+
+func cloneLabels(labels map[string]string) map[string]string {
+	if labels == nil {
+		return nil
+	}
+	clone := make(map[string]string, len(labels))
+	for k, v := range labels {
+		clone[k] = v
+	}
+	return clone
+}
+
+func assertLabels(t *testing.T, got, want map[string]string) {
+	t.Helper()
+	if len(got) != len(want) {
+		t.Fatalf("got labels %v, want %v", got, want)
+	}
+	for k, v := range want {
+		if got[k] != v {
+			t.Errorf("label %q = %q, want %q", k, got[k], v)
+		}
+	}
+}

+ 15 - 0
modules/collector-source/pkg/scrape/index.go

@@ -46,3 +46,18 @@ func buildPVIndex(pvs []*clustercache.PersistentVolume) map[string]types.UID {
 	}
 	return m
 }
+
+// podKey is a composite key for a Pod (namespace + name).
+type podKey struct {
+	namespace string
+	name      string
+}
+
+// buildPodIndex returns a map from (namespace, name) to Pod UID.
+func buildPodIndex(pods []*clustercache.Pod) map[podKey]types.UID {
+	m := make(map[podKey]types.UID, len(pods))
+	for _, pod := range pods {
+		m[podKey{namespace: pod.Namespace, name: pod.Name}] = pod.UID
+	}
+	return m
+}

+ 2 - 1
modules/collector-source/pkg/scrape/network.go

@@ -32,7 +32,8 @@ func newNetworkTargetScraper(provider target.TargetProvider) *TargetScraper {
 			metric.KubecostPodNetworkEgressBytesTotal,
 			metric.KubecostPodNetworkIngressBytesTotal,
 		},
-		true)
+		true,
+		nil)
 }
 
 type NetworkTargetProvider struct {

+ 2 - 1
modules/collector-source/pkg/scrape/opencost.go

@@ -36,5 +36,6 @@ func newOpencostTargetScraper(provider target.TargetProvider) *TargetScraper {
 			metric.NodeGPUCount,
 			metric.KubecostNodeIsSpot,
 		},
-		true)
+		true,
+		nil)
 }

+ 15 - 3
modules/collector-source/pkg/scrape/targetscraper.go

@@ -12,14 +12,19 @@ import (
 	"github.com/opencost/opencost/modules/collector-source/pkg/scrape/target"
 )
 
+// UpdateEnricher optionally batch transforms entire set of updates before they returned from
+// Scrape().
+type UpdateEnricher func(update []metric.Update)
+
 type TargetScraper struct {
 	name           string // identifier for the scraper
 	targetProvider target.TargetProvider
 	metricNames    map[string]struct{} // filter for which metrics will be processed
 	includeMetrics bool                // toggle to make metrics an include or exclude list
+	enrich         UpdateEnricher      // optional per-update enrichment, nil means no-op
 }
 
-func newTargetScrapper(name string, provider target.TargetProvider, metricNames []string, includeMetrics bool) *TargetScraper {
+func newTargetScrapper(name string, provider target.TargetProvider, metricNames []string, includeMetrics bool, enrich UpdateEnricher) *TargetScraper {
 	metricSet := make(map[string]struct{})
 	for _, metricName := range metricNames {
 		metricSet[metricName] = struct{}{}
@@ -29,6 +34,7 @@ func newTargetScrapper(name string, provider target.TargetProvider, metricNames
 		targetProvider: provider,
 		metricNames:    metricSet,
 		includeMetrics: includeMetrics,
+		enrich:         enrich,
 	}
 }
 
@@ -70,11 +76,13 @@ func (s *TargetScraper) Scrape() []metric.Update {
 				if _, ok := s.metricNames[result.Name]; ok != s.includeMetrics {
 					continue
 				}
-				scrapeResults = append(scrapeResults, metric.Update{
+				update := metric.Update{
 					Name:   result.Name,
 					Labels: result.Labels,
 					Value:  result.Value,
-				})
+				}
+
+				scrapeResults = append(scrapeResults, update)
 			}
 			return scrapeResults
 		}
@@ -83,6 +91,10 @@ func (s *TargetScraper) Scrape() []metric.Update {
 
 	updates := concurrentScrape(scrapeFuncs...)
 
+	if s.enrich != nil {
+		s.enrich(updates)
+	}
+
 	// dispatch a scrape event for this specific scrape
 	events.Dispatch(event.ScrapeEvent{
 		ScraperName: s.name,

+ 45 - 3
modules/collector-source/pkg/scrape/targetscraper_test.go

@@ -95,6 +95,13 @@ pod_pvc_allocation{namespace="namespace1",persistentvolume="pvc-1",persistentvol
 pod_pvc_allocation{namespace="namespace1",persistentvolume="pvc-2",persistentvolumeclaim="pvc2",pod="pod2"} 3.4359738368e+10
 `
 
+const enrichScrape = `
+# HELP test_metric test metric
+# TYPE test_metric gauge
+test_metric{name="a"} 1
+test_metric{name="b",extra="already-set"} 2
+`
+
 const dcgmScrape = `
 # HELP DCGM_FI_PROF_GR_ENGINE_ACTIVE Ratio of time the graphics engine is active.
 # TYPE DCGM_FI_PROF_GR_ENGINE_ACTIVE gauge
@@ -411,9 +418,11 @@ func TestTargetScraper_Scrape(t *testing.T) {
 			},
 		},
 		{
-			name:                 "GPU Metric",
-			scrapeText:           dcgmScrape,
-			targetScraperFactory: newDCGMTargetScraper,
+			name:       "GPU Metric",
+			scrapeText: dcgmScrape,
+			targetScraperFactory: func(provider target.TargetProvider) *TargetScraper {
+				return newDCGMTargetScraper(provider, nil)
+			},
 			expected: []metric.Update{
 				{
 					Name: metric.DCGMFIPROFGRENGINEACTIVE,
@@ -441,6 +450,39 @@ func TestTargetScraper_Scrape(t *testing.T) {
 				},
 			},
 		},
+		{
+			name:       "Enrichment",
+			scrapeText: enrichScrape,
+			targetScraperFactory: func(provider target.TargetProvider) *TargetScraper {
+				enrich := func(updates []metric.Update) {
+					for i := range updates {
+						if updates[i].Labels["extra"] != "" {
+							continue
+						}
+						updates[i].Labels["extra"] = "enriched"
+					}
+				}
+				return newTargetScrapper("test-enrich", provider, nil, false, enrich)
+			},
+			expected: []metric.Update{
+				{
+					Name: "test_metric",
+					Labels: map[string]string{
+						"name":  "a",
+						"extra": "enriched",
+					},
+					Value: 1,
+				},
+				{
+					Name: "test_metric",
+					Labels: map[string]string{
+						"name":  "b",
+						"extra": "already-set",
+					},
+					Value: 2,
+				},
+			},
+		},
 	}
 	for _, tt := range tests {
 		t.Run(tt.name, func(t *testing.T) {

+ 18 - 0
modules/prometheus-source/pkg/prom/metricsquerier.go

@@ -2244,6 +2244,24 @@ func (pds *PrometheusMetricsQuerier) QueryDaemonSetAnnotations(start, end time.T
 	return source.NewFuture(source.DecodeAnnotationsResult, ctx.QueryAtTime(queryDaemonSetAnnotations, end))
 }
 
+func (pds *PrometheusMetricsQuerier) QueryDaemonSetArguments(start, end time.Time) *source.Future[source.DaemonSetArgumentResult] {
+	const queryName = "QueryDaemonSetArguments"
+	const queryFmtDaemonSetArguments = `avg(avg_over_time(daemonset_arguments{%s}[%s])) by (%s, uid, arg, value)`
+
+	cfg := pds.promConfig
+
+	durStr := timeutil.DurationString(end.Sub(start))
+	if durStr == "" {
+		panic(fmt.Sprintf("failed to parse duration string passed to %s", queryName))
+	}
+
+	queryDaemonSetArguments := fmt.Sprintf(queryFmtDaemonSetArguments, cfg.ClusterFilter, durStr, cfg.ClusterLabel)
+	log.Debugf(PrometheusMetricsQueryLogFormat, queryName, end.Unix(), queryDaemonSetArguments)
+
+	ctx := pds.promContexts.NewNamedContext(KubeModelContextName)
+	return source.NewFuture(source.DecodeDaemonSetArgumentResult, ctx.QueryAtTime(queryDaemonSetArguments, end))
+}
+
 func (pds *PrometheusMetricsQuerier) QueryJobInfo(start, end time.Time) *source.Future[source.JobInfoResult] {
 	const queryName = "QueryJobInfo"
 	const queryFmtJobInfo = `avg(avg_over_time(job_info{%s}[%s])) by (%s, uid, namespace_uid, job)`

+ 16 - 0
pkg/metrics/kubemodel.go

@@ -2,6 +2,8 @@ package metrics
 
 import (
 	"fmt"
+	"maps"
+	"slices"
 
 	"github.com/opencost/opencost/core/pkg/clustercache"
 	"github.com/opencost/opencost/core/pkg/clusters"
@@ -38,6 +40,7 @@ var kubeModelMetricNames = []string{
 	"daemonset_info",
 	"daemonset_labels",
 	"daemonset_annotations",
+	"daemonset_arguments",
 	"job_info",
 	"job_labels",
 	"job_annotations",
@@ -358,6 +361,7 @@ func (c KubeModelCollector) scrapeDaemonSets(
 	emitInfo := !isDisabled(disabled, "daemonset_info")
 	emitLabels := !isDisabled(disabled, "daemonset_labels")
 	emitAnno := !isDisabled(disabled, "daemonset_annotations")
+	emitArgs := !isDisabled(disabled, "daemonset_arguments")
 
 	for _, ds := range sets {
 		nsUID, ok := nsIndex[ds.Namespace]
@@ -378,6 +382,18 @@ func (c KubeModelCollector) scrapeDaemonSets(
 		if emitAnno {
 			out = append(out, kubeAnnotationsMetric("daemonset_annotations", string(ds.UID), ds.Annotations))
 		}
+		if emitArgs {
+			daemonSetArguments := coreutil.ParseContainerArgs(ds.SpecContainers)
+			for _, arg := range slices.Sorted(maps.Keys(daemonSetArguments)) {
+				out = append(out, newInfoMetric("daemonset_arguments", map[string]string{
+					"uid":           string(ds.UID),
+					"namespace_uid": string(nsUID),
+					"daemonset":     ds.Name,
+					"arg":           arg,
+					"value":         daemonSetArguments[arg],
+				}))
+			}
+		}
 	}
 	return out
 }

Некоторые файлы не были показаны из-за большого количества измененных файлов