Skip to content
1 change: 1 addition & 0 deletions pkg/emitter/emitter.go
Original file line number Diff line number Diff line change
Expand Up @@ -50,6 +50,7 @@ type MetricsSummary struct {
// MetricsSnapshot contains the metrics results from opencost data source queries.
type MetricsSnapshot struct {
Window opencost.Window
FailedQueries []string // names of queries that returned errors
PVActiveMinutes []*source.PVActiveMinutesResult
PVUsedAverage []*source.PVUsedAvgResult
PVUsedMax []*source.PVUsedMaxResult
Expand Down
227 changes: 128 additions & 99 deletions pkg/emitter/snapshot.go
Original file line number Diff line number Diff line change
Expand Up @@ -13,8 +13,33 @@ import (
"github.com/opencost/opencost/core/pkg/log"
"github.com/opencost/opencost/core/pkg/opencost"
"github.com/opencost/opencost/core/pkg/source"
"github.com/prometheus/client_golang/prometheus"
"github.com/prometheus/client_golang/prometheus/promauto"
)

var (
// metricQueryFailures tracks the number of failed metric queries by query name
metricQueryFailures = promauto.NewCounterVec(
prometheus.CounterOpts{
Name: "finops_agent_metric_query_failures_total",
Help: "Total number of failed metric queries by query name",
},
[]string{"query_name"},
)
)

// awaitWithLog awaits a QueryGroupFuture and logs any errors, incrementing the failure counter
// Returns the result slice
func awaitWithLog[T any](name string, future *source.QueryGroupFuture[T], failedQueries *[]string) []*T {
result, err := future.Await()
if err != nil {
log.Warnf("metric query failed for %s: %v", name, err)
metricQueryFailures.WithLabelValues(name).Inc()
*failedQueries = append(*failedQueries, name)
}
Comment thread
peatey marked this conversation as resolved.
return result
Comment thread
peatey marked this conversation as resolved.
}

// SnapshotProvider is an interface that defines a prototype for generating `ClusterSnapshot` instances
// leveraging the agent `DataSource`
type SnapshotProvider interface {
Expand Down Expand Up @@ -410,112 +435,116 @@ func snapshotMetrics(mq source.MetricsQuerier, start, end time.Time) (*MetricsSn
resourceQuotaStatusUsedRamLimitAvgFuture := source.WithGroup(grp, mq.QueryResourceQuotaStatusUsedRAMLimitAverage(start, end))
resourceQuotaStatusUsedRamLimitMaxFuture := source.WithGroup(grp, mq.QueryResourceQuotaStatusUsedRAMLimitMax(start, end))

pvActiveMinutes, _ := pvActiveMinutesFuture.Await()
pvUsedAverage, _ := pvUsedAverageFuture.Await()
pvUsedMax, _ := pvUsedMaxFuture.Await()
localStorageActiveMinutes, _ := localStorageActiveMinutesFuture.Await()
localStorageCost, _ := localStorageCostFuture.Await()
localStorageUsedCost, _ := localStorageUsedCostFuture.Await()
localStorageUsedAvg, _ := localStorageUsedAvgFuture.Await()
localStorageUsedMax, _ := localStorageUsedMaxFuture.Await()
localStorageBytes, _ := localStorageBytesFuture.Await()
nodeActiveMinutes, _ := nodeActiveMinutesFuture.Await()
nodeCPUCoresCapacity, _ := nodeCPUCoresCapacityFuture.Await()
nodeCPUCoresAllocatable, _ := nodeCPUCoresAllocatableFuture.Await()
nodeRAMBytesCapacity, _ := nodeRAMBytesCapacityFuture.Await()
nodeRAMBytesAllocatable, _ := nodeRAMBytesAllocatableFuture.Await()
nodeGPUCount, _ := nodeGPUCountFuture.Await()
nodeCPUModeTotal, _ := nodeCPUModeTotalFuture.Await()
nodeIsSpot, _ := nodeIsSpotFuture.Await()
nodeRAMSystemPercent, _ := nodeRAMSystemPercentFuture.Await()
nodeRAMUserPercent, _ := nodeRAMUserPercentFuture.Await()
lbActiveMinutes, _ := lbActiveMinutesFuture.Await()
lbPricePerHr, _ := lbPricePerHrFuture.Await()
clusterUptime, _ := clusterUptimeFuture.Await()
clusterManagementDuration, _ := clusterManagementDurationFuture.Await()
clusterManagementPricePerHr, _ := clusterManagementPricePerHrFuture.Await()
pods, _ := podsFuture.Await()
podsUID, _ := podsUIDFuture.Await()
ramBytesAllocated, _ := ramBytesAllocatedFuture.Await()
ramRequests, _ := ramRequestsFuture.Await()
ramLimits, _ := ramLimitsFuture.Await()
ramUsageAvg, _ := ramUsageAvgFuture.Await()
ramUsageMax, _ := ramUsageMaxFuture.Await()
nodeRAMPricePerGiBHr, _ := nodeRAMPricePerGiBHrFuture.Await()
cpuCoresAllocated, _ := cpuCoresAllocatedFuture.Await()
cpuRequests, _ := cpuRequestsFuture.Await()
cpuLimits, _ := cpuLimitsFuture.Await()
cpuUsageAvg, _ := cpuUsageAvgFuture.Await()
cpuUsageMax, _ := cpuUsageMaxFuture.Await()
nodeCPUPricePerHr, _ := nodeCPUPricePerHrFuture.Await()
gpusAllocated, _ := gpusAllocatedFuture.Await()
gpusRequested, _ := gpusRequestedFuture.Await()
gpusUsageAvg, _ := gpusUsageAvgFuture.Await()
gpusUsageMax, _ := gpusUsageMaxFuture.Await()
nodeGPUPricePerHr, _ := nodeGPUPricePerHrFuture.Await()
gpuInfo, _ := gpuInfoFuture.Await()
isGPUShared, _ := isGPUSharedFuture.Await()
podPVCAllocation, _ := podPVCAllocationFuture.Await()
pvcBytesRequested, _ := pvcBytesRequestedFuture.Await()
pvcInfo, _ := pvcInfoFuture.Await()
pvBytes, _ := pvBytesFuture.Await()
pvPricePerGiBHour, _ := pvPricePerGiBHourFuture.Await()
pvInfo, _ := pvInfoFuture.Await()
netZoneGiB, _ := netZoneGiBFuture.Await()
netZonePricePerGiB, _ := netZonePricePerGiBFuture.Await()
netRegionGiB, _ := netRegionGiBFuture.Await()
netRegionPricePerGiB, _ := netRegionPricePerGiBFuture.Await()
netInternetGiB, _ := netInternetGiBFuture.Await()
netInternetPricePerGiB, _ := netInternetPricePerGiBFuture.Await()
netInternetServiceGiB, _ := netInternetServiceGiBFuture.Await()
netNatGatewayPricePerGiB, _ := netNatGatewayPricePerGiBFuture.Await()
netNatGatewayGiB, _ := netNatGatewayGiBFuture.Await()
netTransferBytes, _ := netTransferBytesFuture.Await()
netZoneIngressGiB, _ := netZoneIngressGiBFuture.Await()
netRegionIngressGiB, _ := netRegionIngressGiBFuture.Await()
netInternetIngressGiB, _ := netInternetIngressGiBFuture.Await()
netInternetServiceIngressGiB, _ := netInternetServiceIngressGiBFuture.Await()
netNatGatewayIngressPricePerGiB, _ := netNatGatewayIngressPricePerGiBFuture.Await()
netNatGatewayIngressGiB, _ := netNatGatewayIngressGiBFuture.Await()
netReceiveBytes, _ := netReceiveBytesFuture.Await()
namespaceUptime, _ := namespaceUptimeFuture.Await()
namespaceAnnotations, _ := namespaceAnnotationsFuture.Await()
podAnnotations, _ := podAnnotationsFuture.Await()
nodeLabels, _ := nodeLabelsFuture.Await()
namespaceLabels, _ := namespaceLabelsFuture.Await()
podLabels, _ := podLabelsFuture.Await()
serviceLabels, _ := serviceLabelsFuture.Await()
deploymentLabels, _ := deploymentLabelsFuture.Await()
statefulSetLabels, _ := statefulSetLabelsFuture.Await()
daemonSetLabels, _ := daemonSetLabelsFuture.Await()
jobLabels, _ := jobLabelsFuture.Await()
podsWithReplicaSetOwner, _ := podsWithReplicaSetOwnerFuture.Await()
replicaSetsWithoutOwners, _ := replicaSetsWithoutOwnersFuture.Await()
replicaSetsWithRollout, _ := replicaSetsWithRolloutFuture.Await()
resourceQuotaUptime, _ := resourceQuotaUptimeFuture.Await()
resourceQuotaSpecCpuRequestAvg, _ := resourceQuotaSpecCpuRequestAvgFuture.Await()
resourceQuotaSpecCpuRequestMax, _ := resourceQuotaSpecCpuRequestMaxFuture.Await()
resourceQuotaSpecRamRequestAvg, _ := resourceQuotaSpecRamRequestAvgFuture.Await()
resourceQuotaSpecRamRequestMax, _ := resourceQuotaSpecRamRequestMaxFuture.Await()
resourceQuotaSpecCpuLimitAvg, _ := resourceQuotaSpecCpuLimitAvgFuture.Await()
resourceQuotaSpecCpuLimitMax, _ := resourceQuotaSpecCpuLimitMaxFuture.Await()
resourceQuotaSpecRamLimitAvg, _ := resourceQuotaSpecRamLimitAvgFuture.Await()
resourceQuotaSpecRamLimitMax, _ := resourceQuotaSpecRamLimitMaxFuture.Await()
resourceQuotaStatusUsedCpuRequestAvg, _ := resourceQuotaStatusUsedCpuRequestAvgFuture.Await()
resourceQuotaStatusUsedCpuRequestMax, _ := resourceQuotaStatusUsedCpuRequestMaxFuture.Await()
resourceQuotaStatusUsedRamRequestAvg, _ := resourceQuotaStatusUsedRamRequestAvgFuture.Await()
resourceQuotaStatusUsedRamRequestMax, _ := resourceQuotaStatusUsedRamRequestMaxFuture.Await()
resourceQuotaStatusUsedCpuLimitAvg, _ := resourceQuotaStatusUsedCpuLimitAvgFuture.Await()
resourceQuotaStatusUsedCpuLimitMax, _ := resourceQuotaStatusUsedCpuLimitMaxFuture.Await()
resourceQuotaStatusUsedRamLimitAvg, _ := resourceQuotaStatusUsedRamLimitAvgFuture.Await()
resourceQuotaStatusUsedRamLimitMax, _ := resourceQuotaStatusUsedRamLimitMaxFuture.Await()
// Track failed queries for observability
var failedQueries []string

pvActiveMinutes := awaitWithLog("PVActiveMinutes", pvActiveMinutesFuture, &failedQueries)
pvUsedAverage := awaitWithLog("pvUsedAverage", pvUsedAverageFuture, &failedQueries)
pvUsedMax := awaitWithLog("pvUsedMax", pvUsedMaxFuture, &failedQueries)
Comment thread
peatey marked this conversation as resolved.
Outdated
localStorageActiveMinutes := awaitWithLog("localStorageActiveMinutes", localStorageActiveMinutesFuture, &failedQueries)
localStorageCost := awaitWithLog("localStorageCost", localStorageCostFuture, &failedQueries)
localStorageUsedCost := awaitWithLog("localStorageUsedCost", localStorageUsedCostFuture, &failedQueries)
localStorageUsedAvg := awaitWithLog("localStorageUsedAvg", localStorageUsedAvgFuture, &failedQueries)
localStorageUsedMax := awaitWithLog("localStorageUsedMax", localStorageUsedMaxFuture, &failedQueries)
localStorageBytes := awaitWithLog("localStorageBytes", localStorageBytesFuture, &failedQueries)
nodeActiveMinutes := awaitWithLog("nodeActiveMinutes", nodeActiveMinutesFuture, &failedQueries)
nodeCPUCoresCapacity := awaitWithLog("nodeCPUCoresCapacity", nodeCPUCoresCapacityFuture, &failedQueries)
nodeCPUCoresAllocatable := awaitWithLog("nodeCPUCoresAllocatable", nodeCPUCoresAllocatableFuture, &failedQueries)
nodeRAMBytesCapacity := awaitWithLog("nodeRAMBytesCapacity", nodeRAMBytesCapacityFuture, &failedQueries)
nodeRAMBytesAllocatable := awaitWithLog("nodeRAMBytesAllocatable", nodeRAMBytesAllocatableFuture, &failedQueries)
nodeGPUCount := awaitWithLog("nodeGPUCount", nodeGPUCountFuture, &failedQueries)
nodeCPUModeTotal := awaitWithLog("nodeCPUModeTotal", nodeCPUModeTotalFuture, &failedQueries)
nodeIsSpot := awaitWithLog("nodeIsSpot", nodeIsSpotFuture, &failedQueries)
nodeRAMSystemPercent := awaitWithLog("nodeRAMSystemPercent", nodeRAMSystemPercentFuture, &failedQueries)
nodeRAMUserPercent := awaitWithLog("nodeRAMUserPercent", nodeRAMUserPercentFuture, &failedQueries)
lbActiveMinutes := awaitWithLog("lbActiveMinutes", lbActiveMinutesFuture, &failedQueries)
lbPricePerHr := awaitWithLog("lbPricePerHr", lbPricePerHrFuture, &failedQueries)
clusterUptime := awaitWithLog("clusterUptime", clusterUptimeFuture, &failedQueries)
clusterManagementDuration := awaitWithLog("clusterManagementDuration", clusterManagementDurationFuture, &failedQueries)
clusterManagementPricePerHr := awaitWithLog("clusterManagementPricePerHr", clusterManagementPricePerHrFuture, &failedQueries)
pods := awaitWithLog("pods", podsFuture, &failedQueries)
podsUID := awaitWithLog("podsUID", podsUIDFuture, &failedQueries)
ramBytesAllocated := awaitWithLog("ramBytesAllocated", ramBytesAllocatedFuture, &failedQueries)
ramRequests := awaitWithLog("ramRequests", ramRequestsFuture, &failedQueries)
ramLimits := awaitWithLog("ramLimits", ramLimitsFuture, &failedQueries)
ramUsageAvg := awaitWithLog("ramUsageAvg", ramUsageAvgFuture, &failedQueries)
ramUsageMax := awaitWithLog("ramUsageMax", ramUsageMaxFuture, &failedQueries)
nodeRAMPricePerGiBHr := awaitWithLog("nodeRAMPricePerGiBHr", nodeRAMPricePerGiBHrFuture, &failedQueries)
cpuCoresAllocated := awaitWithLog("cpuCoresAllocated", cpuCoresAllocatedFuture, &failedQueries)
cpuRequests := awaitWithLog("cpuRequests", cpuRequestsFuture, &failedQueries)
cpuLimits := awaitWithLog("cpuLimits", cpuLimitsFuture, &failedQueries)
cpuUsageAvg := awaitWithLog("cpuUsageAvg", cpuUsageAvgFuture, &failedQueries)
cpuUsageMax := awaitWithLog("cpuUsageMax", cpuUsageMaxFuture, &failedQueries)
nodeCPUPricePerHr := awaitWithLog("nodeCPUPricePerHr", nodeCPUPricePerHrFuture, &failedQueries)
gpusAllocated := awaitWithLog("gpusAllocated", gpusAllocatedFuture, &failedQueries)
gpusRequested := awaitWithLog("gpusRequested", gpusRequestedFuture, &failedQueries)
gpusUsageAvg := awaitWithLog("gpusUsageAvg", gpusUsageAvgFuture, &failedQueries)
gpusUsageMax := awaitWithLog("gpusUsageMax", gpusUsageMaxFuture, &failedQueries)
nodeGPUPricePerHr := awaitWithLog("nodeGPUPricePerHr", nodeGPUPricePerHrFuture, &failedQueries)
gpuInfo := awaitWithLog("gpuInfo", gpuInfoFuture, &failedQueries)
isGPUShared := awaitWithLog("isGPUShared", isGPUSharedFuture, &failedQueries)
podPVCAllocation := awaitWithLog("podPVCAllocation", podPVCAllocationFuture, &failedQueries)
pvcBytesRequested := awaitWithLog("pvcBytesRequested", pvcBytesRequestedFuture, &failedQueries)
pvcInfo := awaitWithLog("pvcInfo", pvcInfoFuture, &failedQueries)
pvBytes := awaitWithLog("pvBytes", pvBytesFuture, &failedQueries)
pvPricePerGiBHour := awaitWithLog("pvPricePerGiBHour", pvPricePerGiBHourFuture, &failedQueries)
pvInfo := awaitWithLog("pvInfo", pvInfoFuture, &failedQueries)
netZoneGiB := awaitWithLog("netZoneGiB", netZoneGiBFuture, &failedQueries)
netZonePricePerGiB := awaitWithLog("netZonePricePerGiB", netZonePricePerGiBFuture, &failedQueries)
netRegionGiB := awaitWithLog("netRegionGiB", netRegionGiBFuture, &failedQueries)
netRegionPricePerGiB := awaitWithLog("netRegionPricePerGiB", netRegionPricePerGiBFuture, &failedQueries)
netInternetGiB := awaitWithLog("netInternetGiB", netInternetGiBFuture, &failedQueries)
netInternetPricePerGiB := awaitWithLog("netInternetPricePerGiB", netInternetPricePerGiBFuture, &failedQueries)
netInternetServiceGiB := awaitWithLog("netInternetServiceGiB", netInternetServiceGiBFuture, &failedQueries)
netNatGatewayPricePerGiB := awaitWithLog("netNatGatewayPricePerGiB", netNatGatewayPricePerGiBFuture, &failedQueries)
netNatGatewayGiB := awaitWithLog("netNatGatewayGiB", netNatGatewayGiBFuture, &failedQueries)
netTransferBytes := awaitWithLog("netTransferBytes", netTransferBytesFuture, &failedQueries)
netZoneIngressGiB := awaitWithLog("netZoneIngressGiB", netZoneIngressGiBFuture, &failedQueries)
netRegionIngressGiB := awaitWithLog("netRegionIngressGiB", netRegionIngressGiBFuture, &failedQueries)
netInternetIngressGiB := awaitWithLog("netInternetIngressGiB", netInternetIngressGiBFuture, &failedQueries)
netInternetServiceIngressGiB := awaitWithLog("netInternetServiceIngressGiB", netInternetServiceIngressGiBFuture, &failedQueries)
netNatGatewayIngressPricePerGiB := awaitWithLog("netNatGatewayIngressPricePerGiB", netNatGatewayIngressPricePerGiBFuture, &failedQueries)
netNatGatewayIngressGiB := awaitWithLog("netNatGatewayIngressGiB", netNatGatewayIngressGiBFuture, &failedQueries)
netReceiveBytes := awaitWithLog("netReceiveBytes", netReceiveBytesFuture, &failedQueries)
namespaceUptime := awaitWithLog("namespaceUptime", namespaceUptimeFuture, &failedQueries)
namespaceAnnotations := awaitWithLog("namespaceAnnotations", namespaceAnnotationsFuture, &failedQueries)
podAnnotations := awaitWithLog("podAnnotations", podAnnotationsFuture, &failedQueries)
nodeLabels := awaitWithLog("nodeLabels", nodeLabelsFuture, &failedQueries)
namespaceLabels := awaitWithLog("namespaceLabels", namespaceLabelsFuture, &failedQueries)
podLabels := awaitWithLog("podLabels", podLabelsFuture, &failedQueries)
serviceLabels := awaitWithLog("serviceLabels", serviceLabelsFuture, &failedQueries)
deploymentLabels := awaitWithLog("deploymentLabels", deploymentLabelsFuture, &failedQueries)
statefulSetLabels := awaitWithLog("statefulSetLabels", statefulSetLabelsFuture, &failedQueries)
daemonSetLabels := awaitWithLog("daemonSetLabels", daemonSetLabelsFuture, &failedQueries)
jobLabels := awaitWithLog("jobLabels", jobLabelsFuture, &failedQueries)
podsWithReplicaSetOwner := awaitWithLog("podsWithReplicaSetOwner", podsWithReplicaSetOwnerFuture, &failedQueries)
replicaSetsWithoutOwners := awaitWithLog("replicaSetsWithoutOwners", replicaSetsWithoutOwnersFuture, &failedQueries)
replicaSetsWithRollout := awaitWithLog("replicaSetsWithRollout", replicaSetsWithRolloutFuture, &failedQueries)
resourceQuotaUptime := awaitWithLog("resourceQuotaUptime", resourceQuotaUptimeFuture, &failedQueries)
resourceQuotaSpecCpuRequestAvg := awaitWithLog("resourceQuotaSpecCpuRequestAvg", resourceQuotaSpecCpuRequestAvgFuture, &failedQueries)
resourceQuotaSpecCpuRequestMax := awaitWithLog("resourceQuotaSpecCpuRequestMax", resourceQuotaSpecCpuRequestMaxFuture, &failedQueries)
resourceQuotaSpecRamRequestAvg := awaitWithLog("resourceQuotaSpecRamRequestAvg", resourceQuotaSpecRamRequestAvgFuture, &failedQueries)
resourceQuotaSpecRamRequestMax := awaitWithLog("resourceQuotaSpecRamRequestMax", resourceQuotaSpecRamRequestMaxFuture, &failedQueries)
resourceQuotaSpecCpuLimitAvg := awaitWithLog("resourceQuotaSpecCpuLimitAvg", resourceQuotaSpecCpuLimitAvgFuture, &failedQueries)
resourceQuotaSpecCpuLimitMax := awaitWithLog("resourceQuotaSpecCpuLimitMax", resourceQuotaSpecCpuLimitMaxFuture, &failedQueries)
resourceQuotaSpecRamLimitAvg := awaitWithLog("resourceQuotaSpecRamLimitAvg", resourceQuotaSpecRamLimitAvgFuture, &failedQueries)
resourceQuotaSpecRamLimitMax := awaitWithLog("resourceQuotaSpecRamLimitMax", resourceQuotaSpecRamLimitMaxFuture, &failedQueries)
resourceQuotaStatusUsedCpuRequestAvg := awaitWithLog("resourceQuotaStatusUsedCpuRequestAvg", resourceQuotaStatusUsedCpuRequestAvgFuture, &failedQueries)
resourceQuotaStatusUsedCpuRequestMax := awaitWithLog("resourceQuotaStatusUsedCpuRequestMax", resourceQuotaStatusUsedCpuRequestMaxFuture, &failedQueries)
resourceQuotaStatusUsedRamRequestAvg := awaitWithLog("resourceQuotaStatusUsedRamRequestAvg", resourceQuotaStatusUsedRamRequestAvgFuture, &failedQueries)
resourceQuotaStatusUsedRamRequestMax := awaitWithLog("resourceQuotaStatusUsedRamRequestMax", resourceQuotaStatusUsedRamRequestMaxFuture, &failedQueries)
resourceQuotaStatusUsedCpuLimitAvg := awaitWithLog("resourceQuotaStatusUsedCpuLimitAvg", resourceQuotaStatusUsedCpuLimitAvgFuture, &failedQueries)
resourceQuotaStatusUsedCpuLimitMax := awaitWithLog("resourceQuotaStatusUsedCpuLimitMax", resourceQuotaStatusUsedCpuLimitMaxFuture, &failedQueries)
resourceQuotaStatusUsedRamLimitAvg := awaitWithLog("resourceQuotaStatusUsedRamLimitAvg", resourceQuotaStatusUsedRamLimitAvgFuture, &failedQueries)
resourceQuotaStatusUsedRamLimitMax := awaitWithLog("resourceQuotaStatusUsedRamLimitMax", resourceQuotaStatusUsedRamLimitMaxFuture, &failedQueries)

if grp.HasErrors() {
return nil, grp.Error()
}

return &MetricsSnapshot{
Window: opencost.NewClosedWindow(start, end),
FailedQueries: failedQueries,
PVActiveMinutes: pvActiveMinutes,
PVUsedAverage: pvUsedAverage,
PVUsedMax: pvUsedMax,
Expand Down
Loading