master
go 222 lines 4.99 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package gearman
4
5 import (
6 "bufio"
7 "bytes"
8 "context"
9 "errors"
10 "fmt"
11 "strconv"
12 "strings"
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 status, err := c.conn.queryStatus()
25 if err != nil {
26 c.Cleanup(context.Background())
27 return nil, fmt.Errorf("couldn't query status: %v", err)
28 }
29
30 prioStatus, err := c.conn.queryPriorityStatus()
31 if err != nil {
32 c.Cleanup(context.Background())
33 return nil, fmt.Errorf("couldn't query priority status: %v", err)
34 }
35
36 mx := make(map[string]int64)
37
38 if err := c.collectStatus(mx, status); err != nil {
39 return nil, fmt.Errorf("couldn't collect status: %v", err)
40 }
41 if err := c.collectPriorityStatus(mx, prioStatus); err != nil {
42 return nil, fmt.Errorf("couldn't collect priority status: %v", err)
43 }
44
45 return mx, nil
46
47 }
48
49 func (c *Collector) collectStatus(mx map[string]int64, statusData []byte) error {
50 /*
51 Same output as the "gearadmin --status" command:
52
53 FUNCTION\tTOTAL\tRUNNING\tAVAILABLE_WORKERS
54
55 E.g.:
56
57 prefix generic_worker4 78 78 500
58 generic_worker2 78 78 500
59 generic_worker3 0 0 760
60 generic_worker1 0 0 500
61 */
62
63 seen := make(map[string]bool)
64 var foundEnd bool
65 sc := bufio.NewScanner(bytes.NewReader(statusData))
66
67 mx["total_jobs_queued"] = 0
68 mx["total_jobs_running"] = 0
69 mx["total_jobs_waiting"] = 0
70 mx["total_workers_avail"] = 0
71
72 for sc.Scan() {
73 line := strings.TrimSpace(sc.Text())
74
75 if foundEnd = line == "."; foundEnd {
76 break
77 }
78
79 parts := strings.Fields(line)
80
81 // Gearman does not remove old tasks. We are only interested in tasks that have stats.
82 if len(parts) < 4 {
83 continue
84 }
85
86 name := strings.Join(parts[:len(parts)-3], "_")
87 metrics := parts[len(parts)-3:]
88
89 var queued, running, availWorkers int64
90 var err error
91
92 if queued, err = strconv.ParseInt(metrics[0], 10, 64); err != nil {
93 return fmt.Errorf("couldn't parse queued count: %v", err)
94 }
95 if running, err = strconv.ParseInt(metrics[1], 10, 64); err != nil {
96 return fmt.Errorf("couldn't parse running count: %v", err)
97 }
98 if availWorkers, err = strconv.ParseInt(metrics[2], 10, 64); err != nil {
99 return fmt.Errorf("couldn't parse available count: %v", err)
100 }
101
102 px := fmt.Sprintf("function_%s_", name)
103
104 waiting := queued - running
105
106 mx[px+"jobs_queued"] = queued
107 mx[px+"jobs_running"] = running
108 mx[px+"jobs_waiting"] = waiting
109 mx[px+"workers_available"] = availWorkers
110
111 mx["total_jobs_queued"] += queued
112 mx["total_jobs_running"] += running
113 mx["total_jobs_waiting"] += waiting
114 mx["total_workers_available"] += availWorkers
115
116 seen[name] = true
117 }
118
119 if !foundEnd {
120 return errors.New("unexpected status response")
121 }
122
123 for name := range seen {
124 if !c.seenTasks[name] {
125 c.seenTasks[name] = true
126 c.addFunctionStatusCharts(name)
127 }
128 }
129 for name := range c.seenTasks {
130 if !seen[name] {
131 delete(c.seenTasks, name)
132 c.removeFunctionStatusCharts(name)
133 }
134 }
135
136 return nil
137 }
138
139 func (c *Collector) collectPriorityStatus(mx map[string]int64, prioStatusData []byte) error {
140 /*
141 Same output as the "gearadmin --priority-status" command:
142
143 FUNCTION\tHIGH\tNORMAL\tLOW\tAVAILABLE_WORKERS
144 */
145
146 seen := make(map[string]bool)
147 var foundEnd bool
148 sc := bufio.NewScanner(bytes.NewReader(prioStatusData))
149
150 mx["total_high_priority_jobs"] = 0
151 mx["total_normal_priority_jobs"] = 0
152 mx["total_low_priority_jobs"] = 0
153
154 for sc.Scan() {
155 line := strings.TrimSpace(sc.Text())
156
157 if foundEnd = line == "."; foundEnd {
158 break
159 }
160
161 parts := strings.Fields(line)
162 if len(parts) < 5 {
163 continue
164 }
165
166 name := strings.Join(parts[:len(parts)-4], "_")
167 metrics := parts[len(parts)-4:]
168
169 var high, normal, low int64
170 var err error
171
172 if high, err = strconv.ParseInt(metrics[0], 10, 64); err != nil {
173 return fmt.Errorf("couldn't parse high count: %v", err)
174 }
175 if normal, err = strconv.ParseInt(metrics[1], 10, 64); err != nil {
176 return fmt.Errorf("couldn't parse normal count: %v", err)
177 }
178 if low, err = strconv.ParseInt(metrics[2], 10, 64); err != nil {
179 return fmt.Errorf("couldn't parse low count: %v", err)
180 }
181
182 px := fmt.Sprintf("function_%s_", name)
183
184 mx[px+"high_priority_jobs"] = high
185 mx[px+"normal_priority_jobs"] = normal
186 mx[px+"low_priority_jobs"] = low
187 mx["total_high_priority_jobs"] += high
188 mx["total_normal_priority_jobs"] += normal
189 mx["total_low_priority_jobs"] += low
190
191 seen[name] = true
192 }
193
194 if !foundEnd {
195 return errors.New("unexpected priority status response")
196 }
197
198 for name := range seen {
199 if !c.seenPriorityTasks[name] {
200 c.seenPriorityTasks[name] = true
201 c.addFunctionPriorityStatusCharts(name)
202 }
203 }
204 for name := range c.seenPriorityTasks {
205 if !seen[name] {
206 delete(c.seenPriorityTasks, name)
207 c.removeFunctionPriorityStatusCharts(name)
208 }
209 }
210
211 return nil
212 }
213
214 func (c *Collector) establishConn() (gearmanConn, error) {
215 conn := c.newConn(c.Config)
216
217 if err := conn.connect(); err != nil {
218 return nil, err
219 }
220
221 return conn, nil
222 }