master
go 385 lines 9.47 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package beanstalk
4
5 import (
6 "bufio"
7 "context"
8 "errors"
9 "fmt"
10 "net"
11 "os"
12 "strings"
13 "testing"
14 "time"
15
16 "github.com/netdata/netdata/go/plugins/plugin/go.d/pkg/collecttest"
17
18 "github.com/stretchr/testify/assert"
19 "github.com/stretchr/testify/require"
20 )
21
22 var (
23 dataConfigJSON, _ = os.ReadFile("testdata/config.json")
24 dataConfigYAML, _ = os.ReadFile("testdata/config.yaml")
25
26 dataStats, _ = os.ReadFile("testdata/stats.txt")
27 dataListTubes, _ = os.ReadFile("testdata/list-tubes.txt")
28 dataStatsTubeDefault, _ = os.ReadFile("testdata/stats-tube-default.txt")
29 )
30
31 func Test_testDataIsValid(t *testing.T) {
32 for name, data := range map[string][]byte{
33 "dataConfigJSON": dataConfigJSON,
34 "dataConfigYAML": dataConfigYAML,
35 "dataStats": dataStats,
36 "dataListTubes": dataListTubes,
37 "dataStatsTubeDefault": dataStatsTubeDefault,
38 } {
39 require.NotNil(t, data, name)
40 }
41 }
42
43 func TestCollector_ConfigurationSerialize(t *testing.T) {
44 collecttest.TestConfigurationSerialize(t, &Collector{}, dataConfigJSON, dataConfigYAML)
45 }
46
47 func TestCollector_Init(t *testing.T) {
48 tests := map[string]struct {
49 config Config
50 wantFail bool
51 }{
52 "success with default config": {
53 wantFail: false,
54 config: New().Config,
55 },
56 "fails if address not set": {
57 wantFail: true,
58 config: func() Config {
59 conf := New().Config
60 conf.Address = ""
61 return conf
62 }(),
63 },
64 }
65
66 for name, test := range tests {
67 t.Run(name, func(t *testing.T) {
68 collr := New()
69 collr.Config = test.config
70
71 if test.wantFail {
72 assert.Error(t, collr.Init(context.Background()))
73 } else {
74 assert.NoError(t, collr.Init(context.Background()))
75 }
76 })
77 }
78 }
79
80 func TestCollector_Charts(t *testing.T) {
81 assert.NotNil(t, New().Charts())
82 }
83
84 func TestCollector_Check(t *testing.T) {
85 tests := map[string]struct {
86 prepare func() (*Collector, *mockBeanstalkDaemon)
87 wantFail bool
88 }{
89 "success on valid response": {
90 wantFail: false,
91 prepare: prepareCaseOk,
92 },
93 "fails on unexpected response": {
94 wantFail: true,
95 prepare: prepareCaseUnexpectedResponse,
96 },
97 "fails on connection refused": {
98 wantFail: true,
99 prepare: prepareCaseConnectionRefused,
100 },
101 }
102 for name, test := range tests {
103 t.Run(name, func(t *testing.T) {
104 collr, daemon := test.prepare()
105
106 defer func() {
107 assert.NoError(t, daemon.Close(), "daemon.Close()")
108 }()
109 go func() {
110 assert.NoError(t, daemon.Run(), "daemon.Run()")
111 }()
112
113 select {
114 case <-daemon.started:
115 case <-time.After(time.Second * 3):
116 t.Errorf("mock collr daemon start timed out")
117 }
118
119 require.NoError(t, collr.Init(context.Background()))
120
121 if test.wantFail {
122 assert.Error(t, collr.Check(context.Background()))
123 } else {
124 assert.NoError(t, collr.Check(context.Background()))
125 }
126
127 collr.Cleanup(context.Background())
128
129 select {
130 case <-daemon.stopped:
131 case <-time.After(time.Second * 3):
132 t.Errorf("mock collr daemon stop timed out")
133 }
134 })
135 }
136 }
137
138 func TestCollector_Collect(t *testing.T) {
139 tests := map[string]struct {
140 prepare func() (*Collector, *mockBeanstalkDaemon)
141 wantMetrics map[string]int64
142 wantCharts int
143 }{
144 "success on valid response": {
145 prepare: prepareCaseOk,
146 wantMetrics: map[string]int64{
147 "binlog-records-migrated": 0,
148 "binlog-records-written": 0,
149 "cmd-bury": 0,
150 "cmd-delete": 0,
151 "cmd-ignore": 0,
152 "cmd-kick": 0,
153 "cmd-list-tube-used": 0,
154 "cmd-list-tubes": 317,
155 "cmd-list-tubes-watched": 0,
156 "cmd-pause-tube": 0,
157 "cmd-peek": 0,
158 "cmd-peek-buried": 0,
159 "cmd-peek-delayed": 0,
160 "cmd-peek-ready": 0,
161 "cmd-put": 0,
162 "cmd-release": 0,
163 "cmd-reserve": 0,
164 "cmd-reserve-with-timeout": 0,
165 "cmd-stats": 23619,
166 "cmd-stats-job": 0,
167 "cmd-stats-tube": 18964,
168 "cmd-touch": 0,
169 "cmd-use": 0,
170 "cmd-watch": 0,
171 "current-connections": 2,
172 "current-jobs-buried": 0,
173 "current-jobs-delayed": 0,
174 "current-jobs-ready": 0,
175 "current-jobs-reserved": 0,
176 "current-jobs-urgent": 0,
177 "current-producers": 0,
178 "current-tubes": 1,
179 "current-waiting": 0,
180 "current-workers": 0,
181 "job-timeouts": 0,
182 "rusage-stime": 3922,
183 "rusage-utime": 1602,
184 "total-connections": 72,
185 "total-jobs": 0,
186 "tube_default_cmd-delete": 0,
187 "tube_default_cmd-pause-tube": 0,
188 "tube_default_current-jobs-buried": 0,
189 "tube_default_current-jobs-delayed": 0,
190 "tube_default_current-jobs-ready": 0,
191 "tube_default_current-jobs-reserved": 0,
192 "tube_default_current-jobs-urgent": 0,
193 "tube_default_current-using": 2,
194 "tube_default_current-waiting": 0,
195 "tube_default_current-watching": 2,
196 "tube_default_pause": 0,
197 "tube_default_pause-time-left": 0,
198 "tube_default_total-jobs": 0,
199 "uptime": 105881,
200 },
201 wantCharts: len(statsCharts) + len(tubeChartsTmpl)*1,
202 },
203 "fails on unexpected response": {
204 prepare: prepareCaseUnexpectedResponse,
205 wantCharts: len(statsCharts),
206 },
207 "fails on connection refused": {
208 prepare: prepareCaseConnectionRefused,
209 wantCharts: len(statsCharts),
210 },
211 }
212
213 for name, test := range tests {
214 t.Run(name, func(t *testing.T) {
215 collr, daemon := test.prepare()
216
217 defer func() {
218 assert.NoError(t, daemon.Close(), "daemon.Close()")
219 }()
220 go func() {
221 assert.NoError(t, daemon.Run(), "daemon.Run()")
222 }()
223
224 select {
225 case <-daemon.started:
226 case <-time.After(time.Second * 3):
227 t.Errorf("mock collr daemon start timed out")
228 }
229
230 require.NoError(t, collr.Init(context.Background()))
231
232 mx := collr.Collect(context.Background())
233
234 require.Equal(t, test.wantMetrics, mx)
235
236 assert.Equal(t, test.wantCharts, len(*collr.Charts()), "want charts")
237
238 if len(test.wantMetrics) > 0 {
239 collecttest.TestMetricsHasAllChartsDims(t, collr.Charts(), mx)
240 }
241
242 collr.Cleanup(context.Background())
243
244 select {
245 case <-daemon.stopped:
246 case <-time.After(time.Second * 3):
247 t.Errorf("mock collr daemon stop timed out")
248 }
249 })
250 }
251 }
252
253 func prepareCaseOk() (*Collector, *mockBeanstalkDaemon) {
254 daemon := &mockBeanstalkDaemon{
255 addr: "127.0.0.1:65001",
256 started: make(chan struct{}),
257 stopped: make(chan struct{}),
258 dataStats: dataStats,
259 dataListTubes: dataListTubes,
260 dataStatsTube: dataStatsTubeDefault,
261 }
262
263 collr := New()
264 collr.Address = daemon.addr
265
266 return collr, daemon
267 }
268
269 func prepareCaseUnexpectedResponse() (*Collector, *mockBeanstalkDaemon) {
270 daemon := &mockBeanstalkDaemon{
271 addr: "127.0.0.1:65001",
272 started: make(chan struct{}),
273 stopped: make(chan struct{}),
274 dataStats: []byte("INTERNAL_ERROR\n"),
275 dataListTubes: []byte("INTERNAL_ERROR\n"),
276 dataStatsTube: []byte("INTERNAL_ERROR\n"),
277 }
278
279 collr := New()
280 collr.Address = daemon.addr
281
282 return collr, daemon
283 }
284
285 func prepareCaseConnectionRefused() (*Collector, *mockBeanstalkDaemon) {
286 ch := make(chan struct{})
287 close(ch)
288 daemon := &mockBeanstalkDaemon{
289 addr: "127.0.0.1:65001",
290 dontStart: true,
291 started: ch,
292 stopped: ch,
293 }
294
295 collr := New()
296 collr.Address = daemon.addr
297
298 return collr, daemon
299 }
300
301 type mockBeanstalkDaemon struct {
302 addr string
303 srv net.Listener
304 started chan struct{}
305 stopped chan struct{}
306 dontStart bool
307
308 dataStats []byte
309 dataListTubes []byte
310 dataStatsTube []byte
311 }
312
313 func (m *mockBeanstalkDaemon) Run() error {
314 if m.dontStart {
315 return nil
316 }
317
318 srv, err := net.Listen("tcp", m.addr)
319 if err != nil {
320 return err
321 }
322
323 m.srv = srv
324
325 close(m.started)
326 defer close(m.stopped)
327
328 return m.handleConnections()
329 }
330
331 func (m *mockBeanstalkDaemon) Close() error {
332 if m.srv != nil {
333 err := m.srv.Close()
334 m.srv = nil
335 return err
336 }
337 return nil
338 }
339
340 func (m *mockBeanstalkDaemon) handleConnections() error {
341 conn, err := m.srv.Accept()
342 if err != nil || conn == nil {
343 return errors.New("could not accept connection")
344 }
345 return m.handleConnection(conn)
346 }
347
348 func (m *mockBeanstalkDaemon) handleConnection(conn net.Conn) error {
349 defer func() { _ = conn.Close() }()
350
351 rw := bufio.NewReadWriter(bufio.NewReader(conn), bufio.NewWriter(conn))
352 var line string
353 var err error
354
355 for {
356 if line, err = rw.ReadString('\n'); err != nil {
357 return fmt.Errorf("error reading from connection: %v", err)
358 }
359
360 line = strings.TrimSpace(line)
361
362 cmd, param, _ := strings.Cut(line, " ")
363
364 switch cmd {
365 case cmdQuit:
366 return nil
367 case cmdStats:
368 _, err = rw.Write(m.dataStats)
369 case cmdListTubes:
370 _, err = rw.Write(m.dataListTubes)
371 case cmdStatsTube:
372 if param == "default" {
373 _, err = rw.Write(m.dataStatsTube)
374 } else {
375 _, err = rw.WriteString("NOT_FOUND\n")
376 }
377 default:
378 return fmt.Errorf("unexpected command: %s", line)
379 }
380 _ = rw.Flush()
381 if err != nil {
382 return err
383 }
384 }
385 }