Skip to content

Commit 7e73df4

Browse files
committed
feat: enhance process exporter metrics parsing and update sender to handle new structure
1 parent d0e488c commit 7e73df4

3 files changed

Lines changed: 84 additions & 55 deletions

File tree

internal/exporters/process_exporter.go

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,8 @@ type ProcessExporter struct {
1616
client *http.Client
1717
}
1818

19+
var _ Exporter = (*ProcessExporter)(nil)
20+
1921
// NewProcessExporter creates a new ProcessExporter instance
2022
func NewProcessExporter(endpoint string, timeout time.Duration) *ProcessExporter {
2123
// Use defaults if not specified

internal/prometheus/process_exporter_parser.go

Lines changed: 34 additions & 30 deletions
Original file line numberDiff line numberDiff line change
@@ -8,38 +8,37 @@ import (
88
"time"
99
)
1010

11-
// ProcessExporterMetricSnapshot represents parsed process metrics from process_exporter
12-
// This captures per-process group metrics (grouped by command name)
11+
// ProcessExporterMetricSnapshot represents a single process group metric snapshot
12+
// This is a flat structure matching NodeExporterMetricSnapshot pattern
13+
// Each process group (e.g., "nginx", "postgres") becomes one snapshot
1314
type ProcessExporterMetricSnapshot struct {
14-
Timestamp time.Time `json:"timestamp"`
15-
Processes []ProcessMetric `json:"processes"`
15+
Timestamp time.Time `json:"timestamp"`
16+
Name string `json:"name"` // Process name (groupname from process_exporter)
17+
NumProcs int `json:"num_procs"` // Number of processes in this group
18+
CPUSecondsTotal float64 `json:"cpu_seconds_total"` // Total CPU time consumed (counter)
19+
MemoryBytes int64 `json:"memory_bytes"` // Resident memory (RSS) in bytes
1620
}
1721

18-
// ProcessMetric represents metrics for a single process group (e.g., all "nginx" processes)
19-
type ProcessMetric struct {
20-
Name string `json:"name"` // Process name (groupname from process_exporter)
21-
NumProcs int `json:"num_procs"` // Number of processes in this group
22-
CPUSecondsTotal float64 `json:"cpu_seconds_total"` // Total CPU time consumed (counter)
23-
MemoryBytes int64 `json:"memory_bytes"` // Resident memory (RSS) in bytes
22+
// processData is a temporary struct used during parsing
23+
type processData struct {
24+
numProcs int
25+
cpuSecondsTotal float64
26+
memoryBytes int64
2427
}
2528

2629
// ParseProcessExporterMetrics parses Prometheus process_exporter text format
27-
// Extracts per-process group metrics (CPU, memory, count)
30+
// Returns a slice of ProcessExporterMetricSnapshot (one per process group)
2831
//
2932
// Expected metrics from process_exporter:
3033
// - namedprocess_namegroup_num_procs{groupname="nginx"} 4
3134
// - namedprocess_namegroup_cpu_seconds_total{groupname="nginx"} 1234.56
3235
// - namedprocess_namegroup_memory_bytes{groupname="nginx",memtype="resident"} 104857600
33-
func ParseProcessExporterMetrics(data []byte) (*ProcessExporterMetricSnapshot, error) {
34-
snapshot := &ProcessExporterMetricSnapshot{
35-
Timestamp: time.Now().UTC(),
36-
Processes: []ProcessMetric{},
37-
}
38-
36+
func ParseProcessExporterMetrics(data []byte) ([]ProcessExporterMetricSnapshot, error) {
37+
timestamp := time.Now().UTC()
3938
scanner := bufio.NewScanner(bytes.NewReader(data))
4039

4140
// Track metrics per process group (groupname)
42-
processMetrics := make(map[string]*ProcessMetric)
41+
processMetrics := make(map[string]*processData)
4342

4443
for scanner.Scan() {
4544
line := scanner.Text()
@@ -60,18 +59,25 @@ func ParseProcessExporterMetrics(data []byte) (*ProcessExporterMetricSnapshot, e
6059
return nil, fmt.Errorf("scanner error: %w", err)
6160
}
6261

63-
// Convert map to slice
64-
for _, pm := range processMetrics {
62+
// Convert map to slice of flat snapshots
63+
snapshots := []ProcessExporterMetricSnapshot{}
64+
for name, data := range processMetrics {
6565
// Only include processes that have at least 1 running instance
66-
if pm.NumProcs > 0 {
67-
snapshot.Processes = append(snapshot.Processes, *pm)
66+
if data.numProcs > 0 {
67+
snapshots = append(snapshots, ProcessExporterMetricSnapshot{
68+
Timestamp: timestamp,
69+
Name: name,
70+
NumProcs: data.numProcs,
71+
CPUSecondsTotal: data.cpuSecondsTotal,
72+
MemoryBytes: data.memoryBytes,
73+
})
6874
}
6975
}
7076

71-
return snapshot, nil
77+
return snapshots, nil
7278
}
7379

74-
func parseProcessLine(line string, processMetrics map[string]*ProcessMetric) error {
80+
func parseProcessLine(line string, processMetrics map[string]*processData) error {
7581
// Split metric name and value
7682
parts := strings.Fields(line)
7783
if len(parts) < 2 {
@@ -107,26 +113,24 @@ func parseProcessLine(line string, processMetrics map[string]*ProcessMetric) err
107113

108114
// Ensure process metric entry exists
109115
if processMetrics[groupname] == nil {
110-
processMetrics[groupname] = &ProcessMetric{
111-
Name: groupname,
112-
}
116+
processMetrics[groupname] = &processData{}
113117
}
114118

115119
pm := processMetrics[groupname]
116120

117121
// Parse specific metrics
118122
switch metricName {
119123
case "namedprocess_namegroup_num_procs":
120-
pm.NumProcs = int(value)
124+
pm.numProcs = int(value)
121125

122126
case "namedprocess_namegroup_cpu_seconds_total":
123-
pm.CPUSecondsTotal = value
127+
pm.cpuSecondsTotal = value
124128

125129
case "namedprocess_namegroup_memory_bytes":
126130
// Only use resident memory (RSS)
127131
memtype, ok := labels["memtype"]
128132
if ok && memtype == "resident" {
129-
pm.MemoryBytes = int64(value)
133+
pm.memoryBytes = int64(value)
130134
}
131135
}
132136

internal/report/sender.go

Lines changed: 48 additions & 25 deletions
Original file line numberDiff line numberDiff line change
@@ -178,14 +178,15 @@ func (s *Sender) drainLoop() {
178178

179179
// processBatch loads and sends buffered files grouped by exporter
180180
// Returns error if send fails (files are kept for retry)
181-
// Payload format: { "node_exporter": [...], "postgres_exporter": [...] }
181+
// Payload format: { "node_exporter": [...], "process_exporter": [...] }
182182
func (s *Sender) processBatch(filePaths []string) error {
183183
if len(filePaths) == 0 {
184184
return nil
185185
}
186186

187-
// Group entries by exporter name
188-
exporterMetrics := make(map[string][]prometheus.NodeExporterMetricSnapshot)
187+
// Group entries by exporter name - use separate maps for type safety
188+
nodeExporterMetrics := []prometheus.NodeExporterMetricSnapshot{}
189+
processExporterMetrics := []prometheus.ProcessExporterMetricSnapshot{}
189190
processedFiles := []string{}
190191
var serverID string
191192

@@ -216,38 +217,60 @@ func (s *Sender) processBatch(filePaths []string) error {
216217
serverID = entry.ServerID
217218
}
218219

219-
// Parse Prometheus text to structured metrics
220-
// Note: Currently only node_exporter is parsed, other exporters need their own parsers
221-
snapshot, err := prometheus.ParseNodeExporterMetrics(entry.Data)
222-
if err != nil {
223-
logger.Warn("Failed to parse node_exporter metrics, using zero values",
224-
logger.String("exporter", entry.ExporterName),
225-
logger.String("file", filePath),
226-
logger.Err(err))
227-
// Use zero-value snapshot
228-
snapshot = &prometheus.NodeExporterMetricSnapshot{
229-
Timestamp: time.Now().UTC(),
220+
// Parse Prometheus text to structured metrics based on exporter type
221+
switch entry.ExporterName {
222+
case "node_exporter":
223+
snapshot, err := prometheus.ParseNodeExporterMetrics(entry.Data)
224+
if err != nil {
225+
logger.Warn("Failed to parse node_exporter metrics, using zero values",
226+
logger.String("exporter", entry.ExporterName),
227+
logger.String("file", filePath),
228+
logger.Err(err))
229+
// Use zero-value snapshot
230+
snapshot = &prometheus.NodeExporterMetricSnapshot{
231+
Timestamp: time.Now().UTC(),
232+
}
230233
}
231-
}
234+
nodeExporterMetrics = append(nodeExporterMetrics, *snapshot)
235+
236+
case "process_exporter":
237+
snapshots, err := prometheus.ParseProcessExporterMetrics(entry.Data)
238+
if err != nil {
239+
logger.Warn("Failed to parse process_exporter metrics, skipping",
240+
logger.String("exporter", entry.ExporterName),
241+
logger.String("file", filePath),
242+
logger.Err(err))
243+
continue
244+
}
245+
// Append all process snapshots (one per process group)
246+
processExporterMetrics = append(processExporterMetrics, snapshots...)
232247

233-
// Add to exporter's array
234-
exporterMetrics[entry.ExporterName] = append(
235-
exporterMetrics[entry.ExporterName],
236-
*snapshot,
237-
)
248+
default:
249+
logger.Warn("Unknown exporter type, skipping",
250+
logger.String("exporter", entry.ExporterName),
251+
logger.String("file", filePath))
252+
continue
253+
}
238254

239255
processedFiles = append(processedFiles, filePath)
240256
}
241257

242258
// Nothing to send
243-
if len(exporterMetrics) == 0 {
259+
if len(nodeExporterMetrics) == 0 && len(processExporterMetrics) == 0 {
244260
return nil
245261
}
246262

247-
// Build payload: { "node_exporter": [...], "mysql_exporter": [...] }
263+
// Build payload: { "node_exporter": [...], "process_exporter": [...] }
264+
// Only include exporters that have data
248265
payload := make(map[string]interface{})
249-
for exporterName, snapshots := range exporterMetrics {
250-
payload[exporterName] = snapshots
266+
exporterCount := 0
267+
if len(nodeExporterMetrics) > 0 {
268+
payload["node_exporter"] = nodeExporterMetrics
269+
exporterCount++
270+
}
271+
if len(processExporterMetrics) > 0 {
272+
payload["process_exporter"] = processExporterMetrics
273+
exporterCount++
251274
}
252275

253276
// Convert to JSON
@@ -280,7 +303,7 @@ func (s *Sender) processBatch(filePaths []string) error {
280303
if successCount > 0 {
281304
logger.Info("Successfully sent buffered data",
282305
logger.Int("files", successCount),
283-
logger.Int("exporters", len(exporterMetrics)))
306+
logger.Int("exporters", exporterCount))
284307

285308
// Periodically clean up old buffer files
286309
if err := s.buffer.Cleanup(); err != nil {

0 commit comments

Comments
 (0)