| 1 | // SPDX-License-Identifier: GPL-3.0-or-later |
| 2 | |
| 3 | package httpsd |
| 4 | |
| 5 | import ( |
| 6 | "context" |
| 7 | "errors" |
| 8 | "fmt" |
| 9 | "io" |
| 10 | "log/slog" |
| 11 | "net/http" |
| 12 | "time" |
| 13 | |
| 14 | "github.com/netdata/netdata/go/plugins/logger" |
| 15 | "github.com/netdata/netdata/go/plugins/pkg/web" |
| 16 | "github.com/netdata/netdata/go/plugins/plugin/agent/discovery/sd/model" |
| 17 | ) |
| 18 | |
| 19 | const ( |
| 20 | shortName = "http" |
| 21 | fullName = "sd:http" |
| 22 | ) |
| 23 | |
| 24 | func NewDiscoverer(cfg Config) (*Discoverer, error) { |
| 25 | if err := cfg.validate(); err != nil { |
| 26 | return nil, err |
| 27 | } |
| 28 | |
| 29 | client, err := web.NewHTTPClient(cfg.clientConfig()) |
| 30 | if err != nil { |
| 31 | return nil, err |
| 32 | } |
| 33 | |
| 34 | d := &Discoverer{ |
| 35 | Logger: logger.New().With( |
| 36 | slog.String("component", "service discovery"), |
| 37 | slog.String("discoverer", shortName), |
| 38 | ), |
| 39 | client: client, |
| 40 | request: cfg.RequestConfig, |
| 41 | interval: cfg.interval(), |
| 42 | parser: responseParser{format: cfg.format()}, |
| 43 | source: sourceString(cfg), |
| 44 | } |
| 45 | |
| 46 | return d, nil |
| 47 | } |
| 48 | |
| 49 | type Discoverer struct { |
| 50 | *logger.Logger |
| 51 | model.Base |
| 52 | |
| 53 | client *http.Client |
| 54 | request web.RequestConfig |
| 55 | |
| 56 | interval time.Duration |
| 57 | parser responseParser |
| 58 | source string |
| 59 | } |
| 60 | |
| 61 | func (d *Discoverer) String() string { |
| 62 | return fullName |
| 63 | } |
| 64 | |
| 65 | func (d *Discoverer) Discover(ctx context.Context, in chan<- []model.TargetGroup) { |
| 66 | d.Info("instance is started") |
| 67 | d.Debugf("used config: interval: %s, response body limit: %d, source: %s", d.interval, responseBodyLimit, d.source) |
| 68 | defer func() { d.Info("instance is stopped") }() |
| 69 | |
| 70 | d.discover(ctx, in) |
| 71 | |
| 72 | if d.interval <= 0 { |
| 73 | return |
| 74 | } |
| 75 | |
| 76 | tk := time.NewTicker(d.interval) |
| 77 | defer tk.Stop() |
| 78 | |
| 79 | for { |
| 80 | select { |
| 81 | case <-ctx.Done(): |
| 82 | return |
| 83 | case <-tk.C: |
| 84 | d.discover(ctx, in) |
| 85 | } |
| 86 | } |
| 87 | } |
| 88 | |
| 89 | func (d *Discoverer) discover(ctx context.Context, in chan<- []model.TargetGroup) { |
| 90 | tgg, err := d.fetchTargetGroup(ctx) |
| 91 | if err != nil { |
| 92 | if !errors.Is(err, context.Canceled) { |
| 93 | d.Warning(err) |
| 94 | } |
| 95 | return |
| 96 | } |
| 97 | |
| 98 | model.SendTargetGroup(ctx, in, tgg) |
| 99 | } |
| 100 | |
| 101 | func (d *Discoverer) fetchTargetGroup(ctx context.Context) (model.TargetGroup, error) { |
| 102 | req, err := web.NewHTTPRequest(d.request) |
| 103 | if err != nil { |
| 104 | return nil, fmt.Errorf("create HTTP request: %w", err) |
| 105 | } |
| 106 | req = req.WithContext(ctx) |
| 107 | safeURL := sanitizedURL(req.URL.String()) |
| 108 | |
| 109 | resp, err := d.client.Do(req) |
| 110 | if err != nil { |
| 111 | return nil, fmt.Errorf("HTTP request to %q failed: %w", safeURL, err) |
| 112 | } |
| 113 | if resp.Body != nil { |
| 114 | defer func() { _ = resp.Body.Close() }() |
| 115 | } |
| 116 | |
| 117 | if resp.StatusCode != http.StatusOK { |
| 118 | return nil, fmt.Errorf("%s %q returned HTTP status code: %d", req.Method, safeURL, resp.StatusCode) |
| 119 | } |
| 120 | |
| 121 | bs, err := readResponseBody(resp.Body, responseBodyLimit) |
| 122 | if err != nil { |
| 123 | return nil, fmt.Errorf("read response from %q: %w", safeURL, err) |
| 124 | } |
| 125 | |
| 126 | items, err := d.parser.parse(bs, resp.Header.Get("Content-Type")) |
| 127 | if err != nil { |
| 128 | return nil, fmt.Errorf("parse response from %q: %w", safeURL, err) |
| 129 | } |
| 130 | |
| 131 | targets, err := targetsFromItems(d.source, items) |
| 132 | if err != nil { |
| 133 | return nil, err |
| 134 | } |
| 135 | |
| 136 | return &targetGroup{ |
| 137 | source: d.source, |
| 138 | targets: targets, |
| 139 | }, nil |
| 140 | } |
| 141 | |
| 142 | func readResponseBody(r io.Reader, maxBytes int64) ([]byte, error) { |
| 143 | lr := io.LimitReader(r, maxBytes+1) |
| 144 | bs, err := io.ReadAll(lr) |
| 145 | if err != nil { |
| 146 | return nil, err |
| 147 | } |
| 148 | if int64(len(bs)) > maxBytes { |
| 149 | return nil, fmt.Errorf("response body exceeds limit (%d bytes)", maxBytes) |
| 150 | } |
| 151 | return bs, nil |
| 152 | } |