master
go 118 lines 2.25 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package beanstalk
4
5 import (
6 "context"
7 "fmt"
8 "maps"
9 "slices"
10 "time"
11
12 "github.com/netdata/netdata/go/plugins/pkg/stm"
13 )
14
15 func (c *Collector) collect() (map[string]int64, error) {
16 if c.conn == nil {
17 conn, err := c.establishConn()
18 if err != nil {
19 return nil, err
20 }
21 c.conn = conn
22 }
23
24 mx := make(map[string]int64)
25
26 if err := c.collectStats(mx); err != nil {
27 c.Cleanup(context.Background())
28 return nil, err
29 }
30 if err := c.collectTubesStats(mx); err != nil {
31 return mx, err
32 }
33
34 return mx, nil
35 }
36
37 func (c *Collector) collectStats(mx map[string]int64) error {
38 stats, err := c.conn.queryStats()
39 if err != nil {
40 return err
41 }
42 maps.Copy(mx, stm.ToMap(stats))
43 return nil
44 }
45
46 func (c *Collector) collectTubesStats(mx map[string]int64) error {
47 now := time.Now()
48
49 if now.Sub(c.lastDiscoverTubesTime) > c.discoverTubesEvery {
50 tubes, err := c.conn.queryListTubes()
51 if err != nil {
52 return err
53 }
54
55 c.Debugf("discovered tubes (%d): %v", len(tubes), tubes)
56 v := slices.DeleteFunc(tubes, func(s string) bool { return !c.tubeSr.MatchString(s) })
57 if len(tubes) != len(v) {
58 c.Debugf("discovered tubes after filtering (%d): %v", len(v), v)
59 }
60
61 c.discoveredTubes = v
62 c.lastDiscoverTubesTime = now
63 }
64
65 seen := make(map[string]bool)
66
67 for i, tube := range c.discoveredTubes {
68 if tube == "" {
69 continue
70 }
71
72 stats, err := c.conn.queryStatsTube(tube)
73 if err != nil {
74 return err
75 }
76
77 if stats == nil {
78 c.Infof("tube '%s' stats object not found (tube does not exist)", tube)
79 c.discoveredTubes[i] = ""
80 continue
81 }
82 if stats.Name == "" {
83 c.Debugf("tube '%s' stats object has an empty name, ignoring it", tube)
84 c.discoveredTubes[i] = ""
85 continue
86 }
87
88 seen[stats.Name] = true
89 if !c.seenTubes[stats.Name] {
90 c.seenTubes[stats.Name] = true
91 c.addTubeCharts(stats.Name)
92 }
93
94 px := fmt.Sprintf("tube_%s_", stats.Name)
95 for k, v := range stm.ToMap(stats) {
96 mx[px+k] = v
97 }
98 }
99
100 for tube := range c.seenTubes {
101 if !seen[tube] {
102 delete(c.seenTubes, tube)
103 c.removeTubeCharts(tube)
104 }
105 }
106
107 return nil
108 }
109
110 func (c *Collector) establishConn() (beanstalkConn, error) {
111 conn := c.newConn(c.Config, c.Logger)
112
113 if err := conn.connect(); err != nil {
114 return nil, err
115 }
116
117 return conn, nil
118 }