go.d.plugin: execute local-listeners periodically (#17160)
Ilya Mashchenko committed
Mar 14, 2024 at 16:47 UTC
ddb61900a313b1c17be6a039b939be0ae720c7fc
3 files changed
+283
-111
src/go/collectors/go.d.plugin/agent/discovery/sd/discoverer/netlisteners/netlisteners.go
+73
-28
@@ -47,12 +47,16 @@ func NewDiscoverer(cfg Config) (*Discoverer, error) {
47
slog.String("discoverer", shortName),
48
),
49
cfgSource: cfg.Source,
50
- interval: time.Second * 60,
50
ll: &localListenersExec{
51
binPath: filepath.Join(dir, "local-listeners"),
52
timeout: time.Second * 5,
53
},
54
+ interval: time.Minute * 2,
55
+ expiryTime: time.Minute * 10,
56
+ cache: make(map[uint64]*cacheItem),
57
+ started: make(chan struct{}),
58
}
59
+
60
d.Tags().Merge(tags)
61
62
return d, nil
@@ -72,6 +76,15 @@ type (
76
77
interval time.Duration
78
ll localListeners
79
+
80
+ expiryTime time.Duration
81
+ cache map[uint64]*cacheItem // [target.Hash]
82
+
83
+ started chan struct{}
84
+ }
85
+ cacheItem struct {
86
+ lastSeenTime time.Time
87
+ tgt model.Target
88
}
89
localListeners interface {
90
discover(ctx context.Context) ([]byte, error)
@@ -83,25 +96,30 @@ func (d *Discoverer) String() string {
96
}
97
98
func (d *Discoverer) Discover(ctx context.Context, in chan<- []model.TargetGroup) {
99
+ d.Info("instance is started")
100
+ defer func() { d.Info("instance is stopped") }()
101
+
102
+ close(d.started)
103
+
104
if err := d.discoverLocalListeners(ctx, in); err != nil {
105
d.Error(err)
106
return
107
}
108
91
- //tk := time.NewTicker(d.interval)
92
- //defer tk.Stop()
93
- //
94
- //for {
95
- // select {
96
- // case <-ctx.Done():
97
- // return
98
- // case <-tk.C:
99
- // if err := d.discoverLocalListeners(ctx, in); err != nil {
100
- // d.Error(err)
101
- // return
102
- // }
103
- // }
104
- //}
109
+ tk := time.NewTicker(d.interval)
110
+ defer tk.Stop()
111
+
112
+ for {
113
+ select {
114
+ case <-ctx.Done():
115
+ return
116
+ case <-tk.C:
117
+ if err := d.discoverLocalListeners(ctx, in); err != nil {
118
+ d.Warning(err)
119
+ return
120
+ }
121
+ }
122
+ }
123
}
124
125
func (d *Discoverer) discoverLocalListeners(ctx context.Context, in chan<- []model.TargetGroup) error {
@@ -113,11 +131,13 @@ func (d *Discoverer) discoverLocalListeners(ctx context.Context, in chan<- []mod
131
return err
132
}
133
116
- tggs, err := d.parseLocalListeners(bs)
134
+ tgts, err := d.parseLocalListeners(bs)
135
if err != nil {
136
return err
137
}
138
139
+ tggs := d.processTargets(tgts)
140
+
141
select {
142
case <-ctx.Done():
143
case in <- tggs:
@@ -126,7 +146,42 @@ func (d *Discoverer) discoverLocalListeners(ctx context.Context, in chan<- []mod
146
return nil
147
}
148
129
-func (d *Discoverer) parseLocalListeners(bs []byte) ([]model.TargetGroup, error) {
149
+func (d *Discoverer) processTargets(tgts []model.Target) []model.TargetGroup {
150
+ tgg := &targetGroup{
151
+ provider: fullName,
152
+ source: fmt.Sprintf("discoverer=%s,host=localhost", shortName),
153
+ }
154
+ if d.cfgSource != "" {
155
+ tgg.source += fmt.Sprintf(",%s", d.cfgSource)
156
+ }
157
+
158
+ if d.expiryTime.Milliseconds() == 0 {
159
+ tgg.targets = tgts
160
+ return []model.TargetGroup{tgg}
161
+ }
162
+
163
+ now := time.Now()
164
+
165
+ for _, tgt := range tgts {
166
+ hash := tgt.Hash()
167
+ if _, ok := d.cache[hash]; !ok {
168
+ d.cache[hash] = &cacheItem{tgt: tgt}
169
+ }
170
+ d.cache[hash].lastSeenTime = now
171
+ }
172
+
173
+ for k, v := range d.cache {
174
+ if now.Sub(v.lastSeenTime) > d.expiryTime {
175
+ delete(d.cache, k)
176
+ continue
177
+ }
178
+ tgg.targets = append(tgg.targets, v.tgt)
179
+ }
180
+
181
+ return []model.TargetGroup{tgg}
182
+}
183
+
184
+func (d *Discoverer) parseLocalListeners(bs []byte) ([]model.Target, error) {
185
var tgts []model.Target
186
187
sc := bufio.NewScanner(bytes.NewReader(bs))
@@ -161,17 +216,7 @@ func (d *Discoverer) parseLocalListeners(bs []byte) ([]model.TargetGroup, error)
216
tgts = append(tgts, &tgt)
217
}
218
164
- tgg := &targetGroup{
165
- provider: fullName,
166
- source: fmt.Sprintf("discoverer=%s,host=localhost", shortName),
167
- targets: tgts,
168
- }
169
-
170
- if d.cfgSource != "" {
171
- tgg.source += fmt.Sprintf(",%s", d.cfgSource)
172
- }
173
-
174
- return []model.TargetGroup{tgg}, nil
219
+ return tgts, nil
220
}
221
222
type localListenersExec struct {
src/go/collectors/go.d.plugin/agent/discovery/sd/discoverer/netlisteners/netlisteners_test.go
+84
-48
@@ -3,28 +3,23 @@
3
package netlisteners
4
5
import (
6
- "context"
7
- "errors"
6
"testing"
7
+ "time"
8
9
"github.com/netdata/netdata/go/go.d.plugin/agent/discovery/sd/model"
10
)
11
13
-var (
14
- localListenersOutputSample = []byte(`
15
-UDP6|::1|8125|/opt/netdata/usr/sbin/netdata -P /run/netdata/netdata.pid -D
16
-TCP6|::1|8125|/opt/netdata/usr/sbin/netdata -P /run/netdata/netdata.pid -D
17
-TCP|127.0.0.1|8125|/opt/netdata/usr/sbin/netdata -P /run/netdata/netdata.pid -D
18
-UDP|127.0.0.1|53768|/opt/netdata/usr/libexec/netdata/plugins.d/go.d.plugin 1
19
-`)
20
-)
21
-
12
func TestDiscoverer_Discover(t *testing.T) {
13
tests := map[string]discoverySim{
24
- "valid response": {
25
- mock: &mockLocalListenersExec{},
26
- wantDoneBeforeCancel: false,
27
- wantTargetGroups: []model.TargetGroup{&targetGroup{
14
+ "add listeners": {
15
+ listenersCli: func(cli listenersCli, interval, expiry time.Duration) {
16
+ cli.addListener("UDP6|::1|8125|/opt/netdata/usr/sbin/netdata -P /run/netdata/netdata.pid -D")
17
+ cli.addListener("TCP6|::1|8125|/opt/netdata/usr/sbin/netdata -P /run/netdata/netdata.pid -D")
18
+ cli.addListener("TCP|127.0.0.1|8125|/opt/netdata/usr/sbin/netdata -P /run/netdata/netdata.pid -D")
19
+ cli.addListener("UDP|127.0.0.1|53768|/opt/netdata/usr/libexec/netdata/plugins.d/go.d.plugin 1")
20
+ time.Sleep(interval * 2)
21
+ },
22
+ wantGroups: []model.TargetGroup{&targetGroup{
23
provider: "sd:net_listeners",
24
source: "discoverer=net_listeners,host=localhost",
25
targets: []model.Target{
@@ -59,23 +54,83 @@ func TestDiscoverer_Discover(t *testing.T) {
54
},
55
}},
56
},
62
- "empty response": {
63
- mock: &mockLocalListenersExec{emptyResponse: true},
64
- wantDoneBeforeCancel: false,
65
- wantTargetGroups: []model.TargetGroup{&targetGroup{
57
+ "remove listeners; not expired": {
58
+ listenersCli: func(cli listenersCli, interval, expiry time.Duration) {
59
+ cli.addListener("UDP6|::1|8125|/opt/netdata/usr/sbin/netdata -P /run/netdata/netdata.pid -D")
60
+ cli.addListener("TCP6|::1|8125|/opt/netdata/usr/sbin/netdata -P /run/netdata/netdata.pid -D")
61
+ cli.addListener("TCP|127.0.0.1|8125|/opt/netdata/usr/sbin/netdata -P /run/netdata/netdata.pid -D")
62
+ cli.addListener("UDP|127.0.0.1|53768|/opt/netdata/usr/libexec/netdata/plugins.d/go.d.plugin 1")
63
+ time.Sleep(interval * 2)
64
+ cli.removeListener("UDP6|::1|8125|/opt/netdata/usr/sbin/netdata -P /run/netdata/netdata.pid -D")
65
+ cli.removeListener("UDP|127.0.0.1|53768|/opt/netdata/usr/libexec/netdata/plugins.d/go.d.plugin 1")
66
+ time.Sleep(interval * 2)
67
+ },
68
+ wantGroups: []model.TargetGroup{&targetGroup{
69
provider: "sd:net_listeners",
70
source: "discoverer=net_listeners,host=localhost",
71
+ targets: []model.Target{
72
+ withHash(&target{
73
+ Protocol: "UDP6",
74
+ Address: "::1",
75
+ Port: "8125",
76
+ Comm: "netdata",
77
+ Cmdline: "/opt/netdata/usr/sbin/netdata -P /run/netdata/netdata.pid -D",
78
+ }),
79
+ withHash(&target{
80
+ Protocol: "TCP6",
81
+ Address: "::1",
82
+ Port: "8125",
83
+ Comm: "netdata",
84
+ Cmdline: "/opt/netdata/usr/sbin/netdata -P /run/netdata/netdata.pid -D",
85
+ }),
86
+ withHash(&target{
87
+ Protocol: "TCP",
88
+ Address: "127.0.0.1",
89
+ Port: "8125",
90
+ Comm: "netdata",
91
+ Cmdline: "/opt/netdata/usr/sbin/netdata -P /run/netdata/netdata.pid -D",
92
+ }),
93
+ withHash(&target{
94
+ Protocol: "UDP",
95
+ Address: "127.0.0.1",
96
+ Port: "53768",
97
+ Comm: "go.d.plugin",
98
+ Cmdline: "/opt/netdata/usr/libexec/netdata/plugins.d/go.d.plugin 1",
99
+ }),
100
+ },
101
}},
102
},
70
- "error on exec": {
71
- mock: &mockLocalListenersExec{err: true},
72
- wantDoneBeforeCancel: true,
73
- wantTargetGroups: nil,
74
- },
75
- "invalid data": {
76
- mock: &mockLocalListenersExec{invalidResponse: true},
77
- wantDoneBeforeCancel: true,
78
- wantTargetGroups: nil,
103
+ "remove listeners; expired": {
104
+ listenersCli: func(cli listenersCli, interval, expiry time.Duration) {
105
+ cli.addListener("UDP6|::1|8125|/opt/netdata/usr/sbin/netdata -P /run/netdata/netdata.pid -D")
106
+ cli.addListener("TCP6|::1|8125|/opt/netdata/usr/sbin/netdata -P /run/netdata/netdata.pid -D")
107
+ cli.addListener("TCP|127.0.0.1|8125|/opt/netdata/usr/sbin/netdata -P /run/netdata/netdata.pid -D")
108
+ cli.addListener("UDP|127.0.0.1|53768|/opt/netdata/usr/libexec/netdata/plugins.d/go.d.plugin 1")
109
+ time.Sleep(interval * 2)
110
+ cli.removeListener("UDP6|::1|8125|/opt/netdata/usr/sbin/netdata -P /run/netdata/netdata.pid -D")
111
+ cli.removeListener("UDP|127.0.0.1|53768|/opt/netdata/usr/libexec/netdata/plugins.d/go.d.plugin 1")
112
+ time.Sleep(expiry * 2)
113
+ },
114
+ wantGroups: []model.TargetGroup{&targetGroup{
115
+ provider: "sd:net_listeners",
116
+ source: "discoverer=net_listeners,host=localhost",
117
+ targets: []model.Target{
118
+ withHash(&target{
119
+ Protocol: "TCP6",
120
+ Address: "::1",
121
+ Port: "8125",
122
+ Comm: "netdata",
123
+ Cmdline: "/opt/netdata/usr/sbin/netdata -P /run/netdata/netdata.pid -D",
124
+ }),
125
+ withHash(&target{
126
+ Protocol: "TCP",
127
+ Address: "127.0.0.1",
128
+ Port: "8125",
129
+ Comm: "netdata",
130
+ Cmdline: "/opt/netdata/usr/sbin/netdata -P /run/netdata/netdata.pid -D",
131
+ }),
132
+ },
133
+ }},
134
},
135
}
136
@@ -88,26 +143,7 @@ func TestDiscoverer_Discover(t *testing.T) {
143
144
func withHash(l *target) *target {
145
l.hash, _ = calcHash(l)
91
- tags, _ := model.ParseTags("hostnetsocket")
146
+ tags, _ := model.ParseTags("netlisteners")
147
l.Tags().Merge(tags)
148
return l
149
}
95
-
96
-type mockLocalListenersExec struct {
97
- err bool
98
- emptyResponse bool
99
- invalidResponse bool
100
-}
101
-
102
-func (m *mockLocalListenersExec) discover(context.Context) ([]byte, error) {
103
- if m.err {
104
- return nil, errors.New("mock discover() error")
105
- }
106
- if m.emptyResponse {
107
- return nil, nil
108
- }
109
- if m.invalidResponse {
110
- return []byte("this is very incorrect data"), nil
111
- }
112
- return localListenersOutputSample, nil
113
-}
src/go/collectors/go.d.plugin/agent/discovery/sd/discoverer/netlisteners/sim_test.go
+126
-35
@@ -4,6 +4,10 @@ package netlisteners
4
5
import (
6
"context"
7
+ "errors"
8
+ "sort"
9
+ "strings"
10
+ "sync"
11
"testing"
12
"time"
13
@@ -13,63 +17,150 @@ import (
17
"github.com/stretchr/testify/require"
18
)
19
20
+type listenersCli interface {
21
+ addListener(s string)
22
+ removeListener(s string)
23
+}
24
+
25
type discoverySim struct {
17
- mock *mockLocalListenersExec
18
- wantDoneBeforeCancel bool
19
- wantTargetGroups []model.TargetGroup
26
+ listenersCli func(cli listenersCli, interval, expiry time.Duration)
27
+ wantGroups []model.TargetGroup
28
}
29
30
func (sim *discoverySim) run(t *testing.T) {
23
- d, err := NewDiscoverer(Config{Tags: "hostnetsocket"})
31
+ d, err := NewDiscoverer(Config{
32
+ Source: "",
33
+ Tags: "netlisteners",
34
+ })
35
require.NoError(t, err)
36
26
- d.ll = sim.mock
37
+ mock := newMockLocalListenersExec()
38
+
39
+ d.ll = mock
40
+
41
+ d.interval = time.Millisecond * 100
42
+ d.expiryTime = time.Second * 1
43
44
+ seen := make(map[string]model.TargetGroup)
45
ctx, cancel := context.WithCancel(context.Background())
29
- tggs, done := sim.collectTargetGroups(t, ctx, d)
46
+ in := make(chan []model.TargetGroup)
47
+ var wg sync.WaitGroup
48
+
49
+ wg.Add(1)
50
+ go func() {
51
+ defer wg.Done()
52
+ d.Discover(ctx, in)
53
+ }()
54
31
- if sim.wantDoneBeforeCancel {
32
- select {
33
- case <-done:
34
- default:
35
- assert.Fail(t, "discovery hasn't finished before cancel")
55
+ wg.Add(1)
56
+ go func() {
57
+ defer wg.Done()
58
+ for {
59
+ select {
60
+ case <-ctx.Done():
61
+ return
62
+ case tggs := <-in:
63
+ for _, tgg := range tggs {
64
+ seen[tgg.Source()] = tgg
65
+ }
66
+ }
67
}
68
+ }()
69
+
70
+ done := make(chan struct{})
71
+ go func() {
72
+ defer close(done)
73
+ wg.Wait()
74
+ }()
75
+
76
+ select {
77
+ case <-d.started:
78
+ case <-time.After(time.Second * 3):
79
+ require.Fail(t, "discovery failed to start")
80
}
38
- assert.Equal(t, sim.wantTargetGroups, tggs)
81
+
82
+ sim.listenersCli(mock, d.interval, d.expiryTime)
83
84
cancel()
85
+
86
select {
87
case <-done:
88
case <-time.After(time.Second * 3):
44
- assert.Fail(t, "discovery hasn't finished after cancel")
89
+ require.Fail(t, "discovery hasn't finished after cancel")
90
+ }
91
+
92
+ var tggs []model.TargetGroup
93
+ for _, tgg := range seen {
94
+ tggs = append(tggs, tgg)
95
+ }
96
+
97
+ sortTargetGroups(tggs)
98
+ sortTargetGroups(sim.wantGroups)
99
+
100
+ wantLen, gotLen := calcTargets(sim.wantGroups), calcTargets(tggs)
101
+ assert.Equalf(t, wantLen, gotLen, "different len (want %d got %d)", wantLen, gotLen)
102
+ assert.Equal(t, sim.wantGroups, tggs)
103
+}
104
+
105
+func newMockLocalListenersExec() *mockLocalListenersExec {
106
+ return &mockLocalListenersExec{
107
+ listeners: make(map[string]bool),
108
}
109
}
110
48
-func (sim *discoverySim) collectTargetGroups(t *testing.T, ctx context.Context, d *Discoverer) ([]model.TargetGroup, chan struct{}) {
111
+type mockLocalListenersExec struct {
112
+ errResponse bool
113
+ mux sync.Mutex
114
+ listeners map[string]bool
115
+}
116
50
- in := make(chan []model.TargetGroup)
51
- done := make(chan struct{})
117
+func (m *mockLocalListenersExec) addListener(s string) {
118
+ m.mux.Lock()
119
+ defer m.mux.Unlock()
120
53
- go func() { defer close(done); d.Discover(ctx, in) }()
121
+ m.listeners[s] = true
122
+}
123
55
- timeout := time.Second * 5
56
- var tggs []model.TargetGroup
124
+func (m *mockLocalListenersExec) removeListener(s string) {
125
+ m.mux.Lock()
126
+ defer m.mux.Unlock()
127
58
- func() {
59
- for {
60
- select {
61
- case groups := <-in:
62
- if tggs = append(tggs, groups...); len(tggs) == len(sim.wantTargetGroups) {
63
- return
64
- }
65
- case <-done:
66
- return
67
- case <-time.After(timeout):
68
- t.Logf("discovery timed out after %s", timeout)
69
- return
70
- }
71
- }
72
- }()
128
+ delete(m.listeners, s)
129
+}
130
+
131
+func (m *mockLocalListenersExec) discover(context.Context) ([]byte, error) {
132
+ if m.errResponse {
133
+ return nil, errors.New("mock discover() error")
134
+ }
135
+
136
+ m.mux.Lock()
137
+ defer m.mux.Unlock()
138
+
139
+ var buf strings.Builder
140
+ for s := range m.listeners {
141
+ buf.WriteString(s)
142
+ buf.WriteByte('\n')
143
+ }
144
+
145
+ return []byte(buf.String()), nil
146
+}
147
+
148
+func calcTargets(tggs []model.TargetGroup) int {
149
+ var n int
150
+ for _, tgg := range tggs {
151
+ n += len(tgg.Targets())
152
+ }
153
+ return n
154
+}
155
74
- return tggs, done
156
+func sortTargetGroups(tggs []model.TargetGroup) {
157
+ if len(tggs) == 0 {
158
+ return
159
+ }
160
+ sort.Slice(tggs, func(i, j int) bool { return tggs[i].Source() < tggs[j].Source() })
161
+
162
+ for idx := range tggs {
163
+ tgts := tggs[idx].Targets()
164
+ sort.Slice(tgts, func(i, j int) bool { return tgts[i].Hash() < tgts[j].Hash() })
165
+ }
166
}