| 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 | } |