scrapecontroller.go 5.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170
  1. package scrape
  2. import (
  3. "fmt"
  4. "time"
  5. "github.com/opencost/opencost/core/pkg/clustercache"
  6. "github.com/opencost/opencost/core/pkg/clusters"
  7. coreenv "github.com/opencost/opencost/core/pkg/env"
  8. "github.com/opencost/opencost/core/pkg/external"
  9. "github.com/opencost/opencost/core/pkg/log"
  10. "github.com/opencost/opencost/core/pkg/nodestats"
  11. "github.com/opencost/opencost/core/pkg/util/atomic"
  12. "github.com/opencost/opencost/modules/collector-source/pkg/metric"
  13. "github.com/opencost/opencost/modules/collector-source/pkg/util"
  14. )
  15. // ScrapeController initializes and holds the scrapers in addition to running the loop that triggers scrapes
  16. type ScrapeController struct {
  17. scrapeInterval util.Interval
  18. runState atomic.AtomicRunState
  19. scrapers []Scraper
  20. updater metric.Updater
  21. }
  22. // getDefaultMetricFilter builds a MetricFilter from environment variable configuration.
  23. // Metrics whose corresponding env var is false (the default) are added to the deny set.
  24. func getDefaultMetricFilter() MetricFilter {
  25. f := MetricFilter{}
  26. deny := func(name string) { f[name] = struct{}{} }
  27. if !coreenv.IsEmitPodAnnotationsMetric() {
  28. deny(metric.KubePodAnnotations)
  29. }
  30. if !coreenv.IsEmitNamespaceAnnotationsMetric() {
  31. deny(metric.KubeNamespaceAnnotations)
  32. }
  33. if !coreenv.IsEmitDeploymentLabelsMetric() {
  34. deny(metric.DeploymentLabels)
  35. }
  36. if !coreenv.IsEmitDeploymentAnnotationsMetric() {
  37. deny(metric.DeploymentAnnotations)
  38. }
  39. if !coreenv.IsEmitStatefulSetLabelsMetric() {
  40. deny(metric.StatefulSetLabels)
  41. }
  42. if !coreenv.IsEmitStatefulSetAnnotationsMetric() {
  43. deny(metric.StatefulSetAnnotations)
  44. }
  45. if !coreenv.IsEmitDaemonSetLabelsMetric() {
  46. deny(metric.DaemonSetLabels)
  47. }
  48. if !coreenv.IsEmitDaemonSetAnnotationsMetric() {
  49. deny(metric.DaemonSetAnnotations)
  50. }
  51. if !coreenv.IsEmitJobLabelsMetric() {
  52. deny(metric.JobLabels)
  53. }
  54. if !coreenv.IsEmitJobAnnotationsMetric() {
  55. deny(metric.JobAnnotations)
  56. }
  57. if !coreenv.IsEmitCronJobLabelsMetric() {
  58. deny(metric.CronJobLabels)
  59. }
  60. if !coreenv.IsEmitCronJobAnnotationsMetric() {
  61. deny(metric.CronJobAnnotations)
  62. }
  63. if !coreenv.IsEmitReplicaSetLabelsMetric() {
  64. deny(metric.ReplicaSetLabels)
  65. }
  66. if !coreenv.IsEmitReplicaSetAnnotationsMetric() {
  67. deny(metric.ReplicaSetAnnotations)
  68. }
  69. return f
  70. }
  71. func NewScrapeController(
  72. clusterUID string,
  73. scrapeInterval string,
  74. networkPort int,
  75. updater metric.Updater,
  76. clusterInfoProvider clusters.ClusterInfoProvider,
  77. clusterCache clustercache.ClusterCache,
  78. statSummaryClient nodestats.StatSummaryClient,
  79. externalLabelProvider external.LabelProvider,
  80. ) *ScrapeController {
  81. // Start with env-driven defaults, then layer in any caller-supplied entries.
  82. filter := getDefaultMetricFilter()
  83. var scrapers []Scraper
  84. clusterInfoScrapper := withFilter(newClusterInfoScrapper(clusterUID, clusterInfoProvider), filter)
  85. scrapers = append(scrapers, clusterInfoScrapper)
  86. clusterCacheScraper := withFilter(newClusterCacheScraper(clusterCache, externalLabelProvider), filter)
  87. scrapers = append(scrapers, clusterCacheScraper)
  88. opencostScraper := withFilter(newOpenCostScraper(), filter)
  89. scrapers = append(scrapers, opencostScraper)
  90. statSummaryScraper := withFilter(newStatSummaryScraper(statSummaryClient, clusterCache), filter)
  91. scrapers = append(scrapers, statSummaryScraper)
  92. networkScraper := withFilter(newNetworkScraper(networkPort, clusterCache), filter)
  93. scrapers = append(scrapers, networkScraper)
  94. dcgmScraper := withFilter(newDCGMScrapper(clusterCache), filter)
  95. scrapers = append(scrapers, dcgmScraper)
  96. si, err := util.NewInterval(scrapeInterval)
  97. if err != nil {
  98. panic(fmt.Errorf("scrapecontroller failed to create scrape interval: %w", err))
  99. }
  100. sc := &ScrapeController{
  101. scrapeInterval: si,
  102. scrapers: scrapers,
  103. updater: updater,
  104. }
  105. return sc
  106. }
  107. func (sc *ScrapeController) Start() {
  108. // Before we attempt to start, we must ensure we are not in a stopping state
  109. sc.runState.WaitForReset()
  110. // This will atomically check the current state to ensure we can run, then advances the state.
  111. // If the state is already started, it will return false.
  112. if !sc.runState.Start() {
  113. log.Info("metric already running")
  114. return
  115. }
  116. go func() {
  117. nextScrape := time.Now().UTC()
  118. timer := time.NewTimer(time.Duration(0))
  119. for {
  120. select {
  121. case <-sc.runState.OnStop():
  122. sc.runState.Reset()
  123. timer.Stop()
  124. return // exit go routine
  125. case <-timer.C:
  126. sc.Scrape(nextScrape)
  127. nextScrape = sc.scrapeInterval.Add(sc.scrapeInterval.Truncate(time.Now().UTC()), 1)
  128. timer.Reset(time.Until(nextScrape))
  129. }
  130. }
  131. }()
  132. }
  133. func (sc *ScrapeController) Stop() {
  134. sc.runState.Stop()
  135. }
  136. func (sc *ScrapeController) Scrape(timestamp time.Time) {
  137. // Run scrapes concurrently to minimize time from call to data collection
  138. var scrapeFuncs []ScrapeFunc
  139. for i := range sc.scrapers {
  140. scraper := sc.scrapers[i]
  141. scrapeFuncs = append(scrapeFuncs, scraper.Scrape)
  142. }
  143. scrapeResults := concurrentScrape(scrapeFuncs...)
  144. // once all results are returned run updates all at once with the same timestamp
  145. sc.updater.Update(&metric.UpdateSet{
  146. Timestamp: timestamp,
  147. Updates: scrapeResults,
  148. })
  149. }