targetscraper.go 2.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106
  1. package scrape
  2. import (
  3. "io"
  4. "sync"
  5. "github.com/kubecost/events"
  6. "github.com/opencost/opencost/core/pkg/log"
  7. "github.com/opencost/opencost/modules/collector-source/pkg/event"
  8. "github.com/opencost/opencost/modules/collector-source/pkg/metric"
  9. "github.com/opencost/opencost/modules/collector-source/pkg/scrape/parser"
  10. "github.com/opencost/opencost/modules/collector-source/pkg/scrape/target"
  11. )
  12. // UpdateEnricher optionally batch transforms entire set of updates before they returned from
  13. // Scrape().
  14. type UpdateEnricher func(update []metric.Update)
  15. type TargetScraper struct {
  16. name string // identifier for the scraper
  17. targetProvider target.TargetProvider
  18. metricNames map[string]struct{} // filter for which metrics will be processed
  19. includeMetrics bool // toggle to make metrics an include or exclude list
  20. enrich UpdateEnricher // optional per-update enrichment, nil means no-op
  21. }
  22. func newTargetScrapper(name string, provider target.TargetProvider, metricNames []string, includeMetrics bool, enrich UpdateEnricher) *TargetScraper {
  23. metricSet := make(map[string]struct{})
  24. for _, metricName := range metricNames {
  25. metricSet[metricName] = struct{}{}
  26. }
  27. return &TargetScraper{
  28. name: name,
  29. targetProvider: provider,
  30. metricNames: metricSet,
  31. includeMetrics: includeMetrics,
  32. enrich: enrich,
  33. }
  34. }
  35. func (s *TargetScraper) Scrape() []metric.Update {
  36. targets := s.targetProvider.GetTargets()
  37. var errLock sync.Mutex
  38. var errors []error
  39. var scrapeFuncs []ScrapeFunc
  40. for i := range targets {
  41. target := targets[i]
  42. fn := func() []metric.Update {
  43. var scrapeResults []metric.Update
  44. f, err := target.Load()
  45. if err != nil {
  46. errLock.Lock()
  47. errors = append(errors, err)
  48. errLock.Unlock()
  49. log.Errorf("failed to scrape target: %s", err.Error())
  50. return scrapeResults
  51. }
  52. if closer, ok := f.(io.ReadCloser); ok {
  53. defer closer.Close()
  54. }
  55. results, err := parser.Parse(f)
  56. if err != nil {
  57. errLock.Lock()
  58. errors = append(errors, err)
  59. errLock.Unlock()
  60. log.Errorf("failed to parse target: %s", err.Error())
  61. return scrapeResults
  62. }
  63. for _, result := range results {
  64. // filter metrics to be processed by name
  65. if _, ok := s.metricNames[result.Name]; ok != s.includeMetrics {
  66. continue
  67. }
  68. update := metric.Update{
  69. Name: result.Name,
  70. Labels: result.Labels,
  71. Value: result.Value,
  72. }
  73. scrapeResults = append(scrapeResults, update)
  74. }
  75. return scrapeResults
  76. }
  77. scrapeFuncs = append(scrapeFuncs, fn)
  78. }
  79. updates := concurrentScrape(scrapeFuncs...)
  80. if s.enrich != nil {
  81. s.enrich(updates)
  82. }
  83. // dispatch a scrape event for this specific scrape
  84. events.Dispatch(event.ScrapeEvent{
  85. ScraperName: s.name,
  86. Targets: len(targets),
  87. Errors: errors,
  88. })
  89. return updates
  90. }