summaryrefslogtreecommitdiff
path: root/internal/mapr/logformat/csv.go
blob: d82b23865e94fcc4ed6d4ebcfb3c5b180556e132 (plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
package logformat

import (
	"fmt"
	"sync"

	"github.com/mimecast/dtail/internal/protocol"
)

// csvParser parses CSV log lines. The first line encountered for a given
// sourceID is treated as the column header and stored so that subsequent
// lines from the same source can be mapped to named fields. State is kept
// per sourceID because a single parser instance is shared across every
// file/stream processed within a mapreduce session; without this, the
// header row of every file after the first one would silently be mapped
// as a data row, corrupting aggregates.
type csvParser struct {
	defaultParser
	mu      sync.RWMutex
	headers map[string][]string
}

var _ Parser = (*csvParser)(nil)

func newCSVParser(hostname, timeZoneName string, timeZoneOffset int) (*csvParser, error) {
	defaultParser, err := newDefaultParser(hostname, timeZoneName, timeZoneOffset)
	if err != nil {
		return &csvParser{}, err
	}
	return &csvParser{
		defaultParser: *defaultParser,
		headers:       make(map[string][]string),
	}, nil
}

func (p *csvParser) MakeFields(maprLine, sourceID string) (map[string]string, error) {
	header, installed := p.ensureHeader(sourceID, maprLine)
	if installed {
		return nil, ErrIgnoreFields
	}

	fields := make(map[string]string, p.fieldsCapacity)
	p.addDefaultFields(fields, maprLine)
	start := 0
	column := 0
	delimiter := protocol.CSVDelimiter[0]

	for {
		value, next, done := scanDelimitedField(maprLine, start, delimiter)
		if column >= len(header) {
			return fields, fmt.Errorf("CSV file seems corrupted, more fields than header values?")
		}
		p.addDynamicField(fields, header[column], value)
		column++
		if done {
			break
		}
		start = next
	}

	return fields, nil
}

// ensureHeader atomically checks for, and if necessary installs, the header
// for sourceID. It returns the effective header for the source and whether
// this call was the one that installed it. Only the goroutine that actually
// installs the header should tell its caller to ignore the current line
// (i.e. return ErrIgnoreFields); any racing goroutine on the same sourceID
// sees installed=false and proceeds to map its line against the installed
// header. The previous implementation split the check (RLock) from the
// install (Lock), so two goroutines could both observe "missing" and both
// report ErrIgnoreFields, silently dropping the loser's data row.
func (p *csvParser) ensureHeader(sourceID, maprLine string) ([]string, bool) {
	p.mu.RLock()
	if header, ok := p.headers[sourceID]; ok {
		p.mu.RUnlock()
		return header, false
	}
	p.mu.RUnlock()

	p.mu.Lock()
	defer p.mu.Unlock()
	if header, ok := p.headers[sourceID]; ok {
		return header, false
	}
	header := parseHeaderLine(maprLine)
	p.headers[sourceID] = header
	return header, true
}

func parseHeaderLine(maprLine string) []string {
	var header []string
	start := 0
	delimiter := protocol.CSVDelimiter[0]
	for {
		field, next, done := scanDelimitedField(maprLine, start, delimiter)
		header = append(header, field)
		if done {
			break
		}
		start = next
	}
	return header
}