| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106 |
- package scrape
- import (
- "io"
- "sync"
- "github.com/kubecost/events"
- "github.com/opencost/opencost/core/pkg/log"
- "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/parser"
- "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, enrich UpdateEnricher) *TargetScraper {
- metricSet := make(map[string]struct{})
- for _, metricName := range metricNames {
- metricSet[metricName] = struct{}{}
- }
- return &TargetScraper{
- name: name,
- targetProvider: provider,
- metricNames: metricSet,
- includeMetrics: includeMetrics,
- enrich: enrich,
- }
- }
- func (s *TargetScraper) Scrape() []metric.Update {
- targets := s.targetProvider.GetTargets()
- var errLock sync.Mutex
- var errors []error
- var scrapeFuncs []ScrapeFunc
- for i := range targets {
- target := targets[i]
- fn := func() []metric.Update {
- var scrapeResults []metric.Update
- f, err := target.Load()
- if err != nil {
- errLock.Lock()
- errors = append(errors, err)
- errLock.Unlock()
- log.Errorf("failed to scrape target: %s", err.Error())
- return scrapeResults
- }
- if closer, ok := f.(io.ReadCloser); ok {
- defer closer.Close()
- }
- results, err := parser.Parse(f)
- if err != nil {
- errLock.Lock()
- errors = append(errors, err)
- errLock.Unlock()
- log.Errorf("failed to parse target: %s", err.Error())
- return scrapeResults
- }
- for _, result := range results {
- // filter metrics to be processed by name
- if _, ok := s.metricNames[result.Name]; ok != s.includeMetrics {
- continue
- }
- update := metric.Update{
- Name: result.Name,
- Labels: result.Labels,
- Value: result.Value,
- }
- scrapeResults = append(scrapeResults, update)
- }
- return scrapeResults
- }
- scrapeFuncs = append(scrapeFuncs, fn)
- }
- 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,
- Targets: len(targets),
- Errors: errors,
- })
- return updates
- }
|