master
go 432 lines 9.26 KB
Raw
1 // SPDX-License-Identifier: GPL-3.0-or-later
2
3 package metricsaudit
4
5 import (
6 "encoding/json"
7 "fmt"
8 "maps"
9 "os"
10 "path/filepath"
11 "sort"
12 "time"
13
14 "github.com/netdata/netdata/go/plugins/plugin/framework/collectorapi"
15 )
16
17 type manifestJob struct {
18 Name string `json:"name"`
19 Module string `json:"module"`
20 Directory string `json:"directory"`
21 Collections int `json:"collections"`
22 }
23
24 type manifestPayload struct {
25 GeneratedAt time.Time `json:"generated_at"`
26 Jobs []manifestJob `json:"jobs"`
27 }
28
29 func (da *Auditor) EnableDataCapture(dir string, onComplete func()) {
30 da.mu.Lock()
31 defer da.mu.Unlock()
32
33 da.dataDir = dir
34 da.onComplete = onComplete
35 if dir == "" || da.writeCh != nil {
36 return
37 }
38
39 da.writeCh = make(chan writeTask, writeQueueSize)
40 go da.runWriter(da.writeCh)
41 }
42
43 func (da *Auditor) runWriter(ch <-chan writeTask) {
44 for task := range ch {
45 if task.flush != nil {
46 close(task.flush)
47 continue
48 }
49
50 if task.run != nil {
51 if err := task.run(); err != nil {
52 da.recordWriteError(fmt.Errorf("%s: %w", task.label, err))
53 }
54 }
55
56 if task.after != nil {
57 task.after()
58 }
59 }
60 }
61
62 // RegisterJob registers directory info for a job.
63 func (da *Auditor) RegisterJob(jobName, moduleName, dir string) {
64 id := newJobID(moduleName, jobName)
65
66 da.mu.Lock()
67 if da.jobDirs == nil {
68 da.jobDirs = make(map[JobID]string)
69 }
70 da.jobDirs[id] = dir
71 if da.jobDone == nil {
72 da.jobDone = make(map[JobID]bool)
73 }
74 da.jobDone[id] = false
75 da.mu.Unlock()
76
77 if dir == "" {
78 return
79 }
80
81 for _, sub := range []string{"queries", "rows", "metrics", "meta"} {
82 if err := os.MkdirAll(filepath.Join(dir, sub), 0o755); err != nil {
83 da.recordWriteError(fmt.Errorf("prepare %s directory for %s[%s]: %w", sub, moduleName, jobName, err))
84 }
85 }
86 }
87
88 // RecordJobStructure records the initial chart structure for a job.
89 func (da *Auditor) RecordJobStructure(jobName, moduleName string, charts *collectorapi.Charts) {
90 if charts == nil {
91 return
92 }
93
94 job := &JobAnalysis{
95 Name: jobName,
96 Module: moduleName,
97 Charts: make([]ChartAnalysis, 0, len(*charts)),
98 AllSeenMetrics: make(map[string]bool),
99 }
100
101 for _, chart := range *charts {
102 ca := ChartAnalysis{
103 Chart: chart,
104 CollectedValues: make(map[string][]int64),
105 SeenDimensions: make(map[string]bool),
106 }
107 for _, dim := range chart.Dims {
108 ca.CollectedValues[dim.ID] = make([]int64, 0)
109 ca.SeenDimensions[dim.ID] = false
110 }
111 job.Charts = append(job.Charts, ca)
112 }
113
114 id := newJobID(moduleName, jobName)
115 var dir string
116 var captureEnabled bool
117
118 da.mu.Lock()
119 da.jobs[id] = job
120 dir = da.jobDirs[id]
121 captureEnabled = da.dataDir != "" && dir != ""
122 da.mu.Unlock()
123
124 if !captureEnabled {
125 return
126 }
127
128 meta := struct {
129 Job string `json:"job"`
130 Module string `json:"module"`
131 Created time.Time `json:"created_at"`
132 Metadata map[string]string `json:"metadata"`
133 }{
134 Job: jobName,
135 Module: moduleName,
136 Created: time.Now(),
137 Metadata: map[string]string{
138 "module": moduleName,
139 },
140 }
141 path := filepath.Join(dir, "meta", "job.json")
142 da.enqueueJSONWrite(
143 fmt.Sprintf("write metadata for %s[%s]", moduleName, jobName),
144 path,
145 meta,
146 nil,
147 )
148 }
149
150 // UpdateJobStructure updates the chart structure for a job with current charts.
151 // This is needed for collectors that create charts dynamically during collection.
152 func (da *Auditor) UpdateJobStructure(jobName, moduleName string, charts *collectorapi.Charts) {
153 if charts == nil {
154 return
155 }
156
157 id := newJobID(moduleName, jobName)
158
159 da.mu.Lock()
160 defer da.mu.Unlock()
161
162 job, exists := da.jobs[id]
163 if !exists {
164 return
165 }
166
167 existingCharts := make(map[string]*ChartAnalysis)
168 for i := range job.Charts {
169 existingCharts[job.Charts[i].Chart.ID] = &job.Charts[i]
170 }
171
172 job.Charts = make([]ChartAnalysis, 0, len(*charts))
173 for _, chart := range *charts {
174 var ca ChartAnalysis
175 if existing, ok := existingCharts[chart.ID]; ok {
176 ca = *existing
177 ca.Chart = chart
178 for _, dim := range chart.Dims {
179 if _, tracked := ca.CollectedValues[dim.ID]; !tracked {
180 ca.CollectedValues[dim.ID] = make([]int64, 0)
181 ca.SeenDimensions[dim.ID] = false
182 }
183 }
184 } else {
185 ca = ChartAnalysis{
186 Chart: chart,
187 CollectedValues: make(map[string][]int64),
188 SeenDimensions: make(map[string]bool),
189 }
190 for _, dim := range chart.Dims {
191 ca.CollectedValues[dim.ID] = make([]int64, 0)
192 ca.SeenDimensions[dim.ID] = false
193 }
194 }
195 job.Charts = append(job.Charts, ca)
196 }
197 }
198
199 // RecordCollection records collected metrics directly from structured data.
200 func (da *Auditor) RecordCollection(jobName, moduleName string, mx map[string]int64) {
201 if mx == nil {
202 return
203 }
204
205 id := newJobID(moduleName, jobName)
206
207 var seq int
208 var metricsDir string
209 var metricsPath string
210 var captureEnabled bool
211 var manifest *manifestPayload
212 var onComplete func()
213
214 da.mu.Lock()
215 job, exists := da.jobs[id]
216 if !exists {
217 da.mu.Unlock()
218 return
219 }
220
221 job.CollectionCount++
222 job.LastCollection = time.Now()
223 seq = job.CollectionCount
224
225 for metricID := range mx {
226 job.AllSeenMetrics[metricID] = true
227 }
228
229 for i := range job.Charts {
230 ca := &job.Charts[i]
231 for _, dim := range ca.Chart.Dims {
232 if value, collected := mx[dim.ID]; collected {
233 ca.SeenDimensions[dim.ID] = true
234 ca.CollectedValues[dim.ID] = append(ca.CollectedValues[dim.ID], value)
235 }
236 }
237 }
238
239 if dir := da.jobDirs[id]; da.dataDir != "" && dir != "" {
240 captureEnabled = true
241 metricsDir = filepath.Join(dir, "metrics")
242 metricsPath = filepath.Join(metricsDir, fmt.Sprintf("metrics-%04d.json", seq))
243 }
244
245 manifest, onComplete = da.markJobCollectedLocked(id)
246 da.mu.Unlock()
247
248 if captureEnabled {
249 payload := struct {
250 CollectedAt time.Time `json:"collected_at"`
251 Metrics map[string]int64 `json:"metrics"`
252 }{
253 CollectedAt: time.Now(),
254 Metrics: cloneIntMetrics(mx),
255 }
256
257 da.enqueueWriteTask(writeTask{
258 label: fmt.Sprintf("write metrics snapshot for %s[%s]", moduleName, jobName),
259 run: func() error {
260 if err := os.MkdirAll(metricsDir, 0o755); err != nil {
261 return err
262 }
263 return writeJSON(metricsPath, payload)
264 },
265 })
266 }
267
268 da.handleCompletionWrite(manifest, onComplete)
269 }
270
271 func (da *Auditor) handleCompletionWrite(manifest *manifestPayload, onComplete func()) {
272 if onComplete == nil {
273 return
274 }
275 if manifest == nil {
276 go onComplete()
277 return
278 }
279
280 da.mu.RLock()
281 manifestPath := filepath.Join(da.dataDir, "manifest.json")
282 da.mu.RUnlock()
283
284 enqueued := da.enqueueJSONWrite("write audit manifest", manifestPath, manifest, func() {
285 go onComplete()
286 })
287 if !enqueued {
288 go onComplete()
289 }
290 }
291
292 func (da *Auditor) markJobCollectedLocked(id JobID) (*manifestPayload, func()) {
293 if da.jobDone == nil {
294 return nil, nil
295 }
296 if _, tracked := da.jobDone[id]; !tracked {
297 return nil, nil
298 }
299
300 da.jobDone[id] = true
301 for jobID, dir := range da.jobDirs {
302 if dir == "" {
303 continue
304 }
305 if !da.jobDone[jobID] {
306 return nil, nil
307 }
308 }
309
310 if da.completed {
311 return nil, nil
312 }
313 da.completed = true
314 if da.dataDir == "" {
315 return nil, da.onComplete
316 }
317
318 manifest := da.buildManifestLocked()
319 return &manifest, da.onComplete
320 }
321
322 func (da *Auditor) buildManifestLocked() manifestPayload {
323 jobs := make([]manifestJob, 0, len(da.jobs))
324 for id, job := range da.jobs {
325 dir := da.jobDirs[id]
326 jobs = append(jobs, manifestJob{
327 Name: id.Name,
328 Module: id.Module,
329 Directory: dir,
330 Collections: job.CollectionCount,
331 })
332 }
333
334 sort.Slice(jobs, func(i, j int) bool {
335 if jobs[i].Module == jobs[j].Module {
336 return jobs[i].Name < jobs[j].Name
337 }
338 return jobs[i].Module < jobs[j].Module
339 })
340
341 return manifestPayload{
342 GeneratedAt: time.Now(),
343 Jobs: jobs,
344 }
345 }
346
347 func (da *Auditor) enqueueJSONWrite(label, path string, payload any, after func()) bool {
348 return da.enqueueWriteTask(writeTask{
349 label: label,
350 run: func() error {
351 return writeJSON(path, payload)
352 },
353 after: after,
354 })
355 }
356
357 func (da *Auditor) enqueueWriteTask(task writeTask) bool {
358 da.mu.RLock()
359 ch := da.writeCh
360 da.mu.RUnlock()
361 if ch == nil {
362 return false
363 }
364
365 select {
366 case ch <- task:
367 return true
368 default:
369 da.recordWriteError(fmt.Errorf("%s: write queue is full", task.label))
370 return false
371 }
372 }
373
374 func (da *Auditor) flushWriteQueue(timeout time.Duration) bool {
375 da.mu.RLock()
376 ch := da.writeCh
377 da.mu.RUnlock()
378 if ch == nil {
379 return true
380 }
381
382 ack := make(chan struct{})
383 task := writeTask{flush: ack}
384 timer := time.NewTimer(timeout)
385 defer timer.Stop()
386
387 select {
388 case ch <- task:
389 case <-timer.C:
390 da.recordWriteError(fmt.Errorf("flush write queue: timed out while enqueueing sentinel"))
391 return false
392 }
393
394 timer.Reset(timeout)
395 select {
396 case <-ack:
397 return true
398 case <-timer.C:
399 da.recordWriteError(fmt.Errorf("flush write queue: timed out waiting for sentinel"))
400 return false
401 }
402 }
403
404 func (da *Auditor) recordWriteError(err error) {
405 if err == nil {
406 return
407 }
408
409 da.mu.Lock()
410 defer da.mu.Unlock()
411 da.writeErrorCount++
412 if len(da.writeErrors) < maxWriteErrorSamples {
413 da.writeErrors = append(da.writeErrors, err.Error())
414 }
415 }
416
417 func cloneIntMetrics(mx map[string]int64) map[string]int64 {
418 if len(mx) == 0 {
419 return map[string]int64{}
420 }
421 out := make(map[string]int64, len(mx))
422 maps.Copy(out, mx)
423 return out
424 }
425
426 func writeJSON(path string, payload any) error {
427 data, err := json.MarshalIndent(payload, "", " ")
428 if err != nil {
429 return err
430 }
431 return os.WriteFile(path, data, 0o644)
432 }