master
go 152 lines 3.19 KB
Raw
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 }