summaryrefslogtreecommitdiff
path: root/internal/mapr/logformat
diff options
context:
space:
mode:
Diffstat (limited to 'internal/mapr/logformat')
-rw-r--r--internal/mapr/logformat/csv.go97
-rw-r--r--internal/mapr/logformat/csv_test.go170
-rw-r--r--internal/mapr/logformat/custom1.go5
-rw-r--r--internal/mapr/logformat/custom2.go5
-rw-r--r--internal/mapr/logformat/default.go239
-rw-r--r--internal/mapr/logformat/default_benchmark_test.go44
-rw-r--r--internal/mapr/logformat/default_test.go42
-rw-r--r--internal/mapr/logformat/delimited.go12
-rw-r--r--internal/mapr/logformat/generic.go15
-rw-r--r--internal/mapr/logformat/generickv.go36
-rw-r--r--internal/mapr/logformat/mimecast.go5
-rw-r--r--internal/mapr/logformat/parser.go123
-rw-r--r--internal/mapr/logformat/parser_test.go69
-rw-r--r--internal/mapr/logformat/variables.go107
-rw-r--r--internal/mapr/logformat/variables_test.go145
15 files changed, 981 insertions, 133 deletions
diff --git a/internal/mapr/logformat/csv.go b/internal/mapr/logformat/csv.go
index ea85ca9..d82b238 100644
--- a/internal/mapr/logformat/csv.go
+++ b/internal/mapr/logformat/csv.go
@@ -2,52 +2,103 @@ package logformat
import (
"fmt"
- "strings"
+ "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
- header []string
- hasHeader bool
+ 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}, nil
+ return &csvParser{
+ defaultParser: *defaultParser,
+ headers: make(map[string][]string),
+ }, nil
}
-func (p *csvParser) MakeFields(maprLine string) (map[string]string, error) {
- if !p.hasHeader {
- p.parseHeader(maprLine)
+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, 7+len(p.header))
- fields["*"] = "*"
- fields["$hostname"] = p.hostname
- fields["$server"] = p.hostname
- fields["$line"] = maprLine
- fields["$empty"] = ""
- fields["$timezone"] = p.timeZoneName
- fields["$timeoffset"] = p.timeZoneOffset
-
- splitted := strings.Split(maprLine, protocol.CSVDelimiter)
- for i, value := range splitted {
- if i >= len(p.header) {
+ 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?")
}
- fields[p.header[i]] = value
+ p.addDynamicField(fields, header[column], value)
+ column++
+ if done {
+ break
+ }
+ start = next
}
return fields, nil
}
-func (p *csvParser) parseHeader(maprLine string) {
- p.header = strings.Split(maprLine, protocol.CSVDelimiter)
- p.hasHeader = true
+// 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
}
diff --git a/internal/mapr/logformat/csv_test.go b/internal/mapr/logformat/csv_test.go
index 1baf032..fa85a99 100644
--- a/internal/mapr/logformat/csv_test.go
+++ b/internal/mapr/logformat/csv_test.go
@@ -2,6 +2,7 @@ package logformat
import (
"strings"
+ "sync"
"testing"
"github.com/mimecast/dtail/internal/protocol"
@@ -23,13 +24,15 @@ func TestCSVLogFormat(t *testing.T) {
strings.Join(dataLine2, protocol.CSVDelimiter),
}
+ const sourceID = "file-a"
+
// First line is the header!
- if _, err := parser.MakeFields(inputs[0]); err != ErrIgnoreFields {
+ if _, err := parser.MakeFields(inputs[0], sourceID); err != ErrIgnoreFields {
t.Errorf("Unable to parse the CSV header")
}
// First data line
- fields, err := parser.MakeFields(inputs[1])
+ fields, err := parser.MakeFields(inputs[1], sourceID)
if err != nil {
t.Errorf("Unable to parse first CSV data line: %s", err.Error())
}
@@ -41,7 +44,7 @@ func TestCSVLogFormat(t *testing.T) {
}
// Second data line
- fields, err = parser.MakeFields(inputs[2])
+ fields, err = parser.MakeFields(inputs[2], sourceID)
if err != nil {
t.Errorf("Unable to parse first CSV data line: %s", err.Error())
}
@@ -52,3 +55,164 @@ func TestCSVLogFormat(t *testing.T) {
t.Errorf("Expected 'color' to be 'Black' but got '%s'", val)
}
}
+
+// TestCSVLogFormatMultiFileHeaders reproduces the bug where a single
+// csvParser instance (as used by the Aggregate for every file
+// in a mapreduce session) was treating the header row of every file after
+// the first as a data row, silently corrupting aggregates.
+func TestCSVLogFormatMultiFileHeaders(t *testing.T) {
+ parser, err := NewParser("csv", nil)
+ if err != nil {
+ t.Fatalf("Unable to create parser: %s", err.Error())
+ }
+
+ headersA := []string{"name", "value"}
+ headersB := []string{"color", "count"}
+
+ fileA := []string{
+ strings.Join(headersA, protocol.CSVDelimiter),
+ strings.Join([]string{"alpha", "1"}, protocol.CSVDelimiter),
+ strings.Join([]string{"beta", "2"}, protocol.CSVDelimiter),
+ }
+ fileB := []string{
+ strings.Join(headersB, protocol.CSVDelimiter),
+ strings.Join([]string{"orange", "3"}, protocol.CSVDelimiter),
+ strings.Join([]string{"black", "4"}, protocol.CSVDelimiter),
+ }
+
+ const sourceA = "file-a"
+ const sourceB = "file-b"
+
+ // First line of file A is its header.
+ if _, err := parser.MakeFields(fileA[0], sourceA); err != ErrIgnoreFields {
+ t.Fatalf("Expected header line of file A to be ignored, got err=%v", err)
+ }
+ for _, line := range fileA[1:] {
+ fields, err := parser.MakeFields(line, sourceA)
+ if err != nil {
+ t.Fatalf("Unable to parse data line %q of file A: %s", line, err.Error())
+ }
+ if _, ok := fields["name"]; !ok {
+ t.Errorf("Expected file A field 'name' for line %q, got %v", line, fields)
+ }
+ }
+
+ // First line of file B MUST also be treated as a header, not a data row.
+ if _, err := parser.MakeFields(fileB[0], sourceB); err != ErrIgnoreFields {
+ t.Fatalf("Expected header line of file B to be ignored (bug: header is being consumed as a data row), got err=%v", err)
+ }
+
+ // Data lines of file B must be mapped against file B's headers, not
+ // file A's.
+ for _, line := range fileB[1:] {
+ fields, err := parser.MakeFields(line, sourceB)
+ if err != nil {
+ t.Fatalf("Unable to parse data line %q of file B: %s", line, err.Error())
+ }
+ if _, ok := fields["color"]; !ok {
+ t.Errorf("Expected file B field 'color' for line %q, got %v", line, fields)
+ }
+ if _, ok := fields["name"]; ok {
+ t.Errorf("File B line %q should not carry file A field 'name'; got %v",
+ line, fields)
+ }
+ }
+}
+
+// TestCSVLogFormatConcurrentSameSourceInstall reproduces a TOCTOU bug in
+// csvParser.MakeFields: the original code first called headerFor under
+// RLock, and only if the header was missing did it call parseHeader under
+// Lock. Two goroutines racing on the same sourceID could both observe
+// "missing" and both return ErrIgnoreFields, silently dropping the loser's
+// data row (the installer wrote the header; the non-installer still
+// signalled "this line was a header" to the caller).
+//
+// With the fix, the check-and-install is a single critical section and
+// exactly one of the two concurrent calls reports ErrIgnoreFields; the
+// other maps its line against the installed header.
+func TestCSVLogFormatConcurrentSameSourceInstall(t *testing.T) {
+ header := strings.Join([]string{"name", "value"}, protocol.CSVDelimiter)
+ data := strings.Join([]string{"alpha", "1"}, protocol.CSVDelimiter)
+
+ const attempts = 500
+ const sourceID = "file-race"
+
+ for attempt := 0; attempt < attempts; attempt++ {
+ parser, err := NewParser("csv", nil)
+ if err != nil {
+ t.Fatalf("attempt %d: unable to create parser: %s", attempt, err.Error())
+ }
+
+ lines := [2]string{header, data}
+ var results [2]struct {
+ fields map[string]string
+ err error
+ }
+
+ start := make(chan struct{})
+ var wg sync.WaitGroup
+ wg.Add(2)
+ for i := 0; i < 2; i++ {
+ i := i
+ go func() {
+ defer wg.Done()
+ <-start
+ results[i].fields, results[i].err = parser.MakeFields(lines[i], sourceID)
+ }()
+ }
+ close(start)
+ wg.Wait()
+
+ ignored := 0
+ for _, r := range results {
+ if r.err == ErrIgnoreFields {
+ ignored++
+ }
+ }
+ if ignored != 1 {
+ t.Fatalf("attempt %d: expected exactly one ErrIgnoreFields across two racing calls on the same sourceID, got %d; results=%+v",
+ attempt, ignored, results)
+ }
+ }
+}
+
+// TestCSVLogFormatConcurrentSources ensures the per-source header store is
+// safe for concurrent access across multiple sourceIDs, matching how the
+// aggregator drives the parser from batched lines across files.
+func TestCSVLogFormatConcurrentSources(t *testing.T) {
+ parser, err := NewParser("csv", nil)
+ if err != nil {
+ t.Fatalf("Unable to create parser: %s", err.Error())
+ }
+
+ header := strings.Join([]string{"name", "value"}, protocol.CSVDelimiter)
+ data := strings.Join([]string{"alpha", "1"}, protocol.CSVDelimiter)
+
+ const workers = 16
+ const iterations = 64
+
+ var wg sync.WaitGroup
+ wg.Add(workers)
+ for w := 0; w < workers; w++ {
+ go func(id int) {
+ defer wg.Done()
+ sourceID := "source-" + string(rune('a'+id))
+ if _, err := parser.MakeFields(header, sourceID); err != ErrIgnoreFields {
+ t.Errorf("worker %d: expected header to be ignored, got err=%v", id, err)
+ return
+ }
+ for i := 0; i < iterations; i++ {
+ fields, err := parser.MakeFields(data, sourceID)
+ if err != nil {
+ t.Errorf("worker %d: parse err=%v", id, err)
+ return
+ }
+ if fields["name"] != "alpha" {
+ t.Errorf("worker %d: expected name=alpha, got %q", id, fields["name"])
+ return
+ }
+ }
+ }(w)
+ }
+ wg.Wait()
+}
diff --git a/internal/mapr/logformat/custom1.go b/internal/mapr/logformat/custom1.go
index 7229f3e..a159328 100644
--- a/internal/mapr/logformat/custom1.go
+++ b/internal/mapr/logformat/custom1.go
@@ -2,15 +2,18 @@ package logformat
import "errors"
+// ErrCustom1NotImplemented indicates custom1 parser is only a template.
var ErrCustom1NotImplemented error = errors.New("custom1 log format is not implemented")
// Template for creating a custom log format.
type custom1Parser struct{}
+var _ Parser = (*custom1Parser)(nil)
+
func newCustom1Parser(hostname, timeZoneName string, timeZoneOffset int) (*custom1Parser, error) {
return &custom1Parser{}, ErrCustom1NotImplemented
}
-func (p *custom1Parser) MakeFields(maprLine string) (map[string]string, error) {
+func (p *custom1Parser) MakeFields(maprLine, _ string) (map[string]string, error) {
return nil, ErrCustom1NotImplemented
}
diff --git a/internal/mapr/logformat/custom2.go b/internal/mapr/logformat/custom2.go
index 262c721..a1a732b 100644
--- a/internal/mapr/logformat/custom2.go
+++ b/internal/mapr/logformat/custom2.go
@@ -2,15 +2,18 @@ package logformat
import "errors"
+// ErrCustom2NotImplemented indicates custom2 parser is only a template.
var ErrCustom2NotImplemented error = errors.New("custom2 log format is not implemented")
// Template for creating a custom log format.
type custom2Parser struct{}
+var _ Parser = (*custom2Parser)(nil)
+
func newCustom2Parser(hostname, timeZoneName string, timeZoneOffset int) (*custom2Parser, error) {
return &custom2Parser{}, ErrCustom2NotImplemented
}
-func (p *custom2Parser) MakeFields(maprLine string) (map[string]string, error) {
+func (p *custom2Parser) MakeFields(maprLine, _ string) (map[string]string, error) {
return nil, ErrCustom2NotImplemented
}
diff --git a/internal/mapr/logformat/default.go b/internal/mapr/logformat/default.go
index a499bc5..c1e6e44 100644
--- a/internal/mapr/logformat/default.go
+++ b/internal/mapr/logformat/default.go
@@ -4,6 +4,7 @@ import (
"fmt"
"strings"
+ "github.com/mimecast/dtail/internal/mapr"
"github.com/mimecast/dtail/internal/protocol"
)
@@ -11,62 +12,216 @@ type defaultParser struct {
hostname string
timeZoneName string
timeZoneOffset string
+ fieldsCapacity int
+
+ wantStar bool
+ wantLine bool
+ wantEmpty bool
+ wantHostname bool
+ wantServer bool
+ wantTimezone bool
+ wantTimeOffset bool
+ wantSeverity bool
+ wantLogLevel bool
+ wantTime bool
+ wantDate bool
+ wantHour bool
+ wantMinute bool
+ wantSecond bool
+ wantPID bool
+ wantCaller bool
+ wantCPUs bool
+ wantGoroutines bool
+ wantCGOCalls bool
+ wantLoadAvg bool
+ wantUptime bool
+
+ allDynamicFields bool
+ dynamicFields map[string]struct{}
}
func newDefaultParser(hostname, timeZoneName string, timeZoneOffset int) (*defaultParser, error) {
- return &defaultParser{
+ parser := &defaultParser{
hostname: hostname,
timeZoneName: timeZoneName,
timeZoneOffset: fmt.Sprintf("%d", timeZoneOffset),
- }, nil
+ }
+ parser.configureFieldPlan(mapr.ParserFieldPlan{AllFields: true})
+ return parser, nil
+}
+
+func (p *defaultParser) setQuery(query *mapr.Query) {
+ p.configureFieldPlan(query.ParserFieldPlan())
}
-func (p *defaultParser) MakeFields(maprLine string) (map[string]string, error) {
- splitted := strings.Split(maprLine, protocol.FieldDelimiter)
+func (p *defaultParser) MakeFields(maprLine, _ string) (map[string]string, error) {
+ fields := make(map[string]string, p.fieldsCapacity)
+ tokenIndex := 0
+ start := 0
+ delimiter := protocol.FieldDelimiter[0]
- if len(splitted) < 11 || !strings.HasPrefix(splitted[9], "MAPREDUCE:") ||
- !strings.HasPrefix(splitted[0], "INFO") {
+ for {
+ token, next, done := scanDelimitedField(maprLine, start, delimiter)
+ switch {
+ case tokenIndex == 0:
+ if !strings.HasPrefix(token, "INFO") {
+ return nil, ErrIgnoreFields
+ }
+ p.addDefaultFields(fields, maprLine)
+ if p.wantSeverity {
+ fields["$severity"] = token
+ }
+ if p.wantLogLevel {
+ fields["$loglevel"] = token
+ }
+ case tokenIndex == 1:
+ if p.wantTime {
+ fields["$time"] = token
+ }
+ if len(token) == 15 {
+ // Example: 20211002-071209
+ if p.wantDate {
+ fields["$date"] = token[0:8]
+ }
+ if p.wantHour {
+ fields["$hour"] = token[9:11]
+ }
+ if p.wantMinute {
+ fields["$minute"] = token[11:13]
+ }
+ if p.wantSecond {
+ fields["$second"] = token[13:]
+ }
+ }
+ case tokenIndex == 2:
+ if p.wantPID {
+ fields["$pid"] = token
+ }
+ case tokenIndex == 3:
+ if p.wantCaller {
+ fields["$caller"] = token
+ }
+ case tokenIndex == 4:
+ if p.wantCPUs {
+ fields["$cpus"] = token
+ }
+ case tokenIndex == 5:
+ if p.wantGoroutines {
+ fields["$goroutines"] = token
+ }
+ case tokenIndex == 6:
+ if p.wantCGOCalls {
+ fields["$cgocalls"] = token
+ }
+ case tokenIndex == 7:
+ if p.wantLoadAvg {
+ fields["$loadavg"] = token
+ }
+ case tokenIndex == 8:
+ if p.wantUptime {
+ fields["$uptime"] = token
+ }
+ case tokenIndex == 9:
+ if !strings.HasPrefix(token, "MAPREDUCE:") {
+ return nil, ErrIgnoreFields
+ }
+ default:
+ if err := p.addKeyValueField(fields, token); err != nil {
+ return fields, err
+ }
+ }
+
+ tokenIndex++
+ if done {
+ break
+ }
+ start = next
+ }
+
+ if tokenIndex < 11 {
// Not a DTail mapreduce log line.
return nil, ErrIgnoreFields
}
- fields := make(map[string]string, len(splitted)+8)
-
- fields["*"] = "*"
- fields["$line"] = maprLine
- fields["$empty"] = ""
- fields["$hostname"] = p.hostname
- fields["$server"] = p.hostname
- fields["$timezone"] = p.timeZoneName
- fields["$timeoffset"] = p.timeZoneOffset
-
- fields["$severity"] = splitted[0]
- fields["$loglevel"] = splitted[0]
-
- time := splitted[1]
- fields["$time"] = time
- if len(time) == 15 {
- // Example: 20211002-071209
- fields["$date"] = time[0:8]
- fields["$hour"] = time[9:11]
- fields["$minute"] = time[11:13]
- fields["$second"] = time[13:]
+ return fields, nil
+}
+
+func (p *defaultParser) addDefaultFields(fields map[string]string, maprLine string) {
+ if p.wantStar {
+ fields["*"] = "*"
}
- fields["$pid"] = splitted[2]
- fields["$caller"] = splitted[3]
- fields["$cpus"] = splitted[4]
- fields["$goroutines"] = splitted[5]
- fields["$cgocalls"] = splitted[6]
- fields["$loadavg"] = splitted[7]
- fields["$uptime"] = splitted[8]
-
- for _, kv := range splitted[10:] {
- keyAndValue := strings.SplitN(kv, "=", 2)
- if len(keyAndValue) != 2 {
- return fields, fmt.Errorf("Unable to parse key-value token '%s'", kv)
- }
- fields[keyAndValue[0]] = keyAndValue[1]
+ if p.wantLine {
+ fields["$line"] = maprLine
+ }
+ if p.wantEmpty {
+ fields["$empty"] = ""
+ }
+ if p.wantHostname {
+ fields["$hostname"] = p.hostname
+ }
+ if p.wantServer {
+ fields["$server"] = p.hostname
}
+ if p.wantTimezone {
+ fields["$timezone"] = p.timeZoneName
+ }
+ if p.wantTimeOffset {
+ fields["$timeoffset"] = p.timeZoneOffset
+ }
+}
- return fields, nil
+func (p *defaultParser) addDynamicField(fields map[string]string, key string, value string) {
+ if p.allDynamicFields {
+ fields[key] = value
+ return
+ }
+ if _, ok := p.dynamicFields[key]; ok {
+ fields[key] = value
+ }
+}
+
+func (p *defaultParser) addKeyValueField(fields map[string]string, token string) error {
+ keyAndValueIndex := strings.IndexByte(token, '=')
+ if keyAndValueIndex < 0 {
+ return fmt.Errorf("Unable to parse key-value token '%s'", token)
+ }
+ p.addDynamicField(fields, token[:keyAndValueIndex], token[keyAndValueIndex+1:])
+ return nil
+}
+
+func (p *defaultParser) configureFieldPlan(plan mapr.ParserFieldPlan) {
+ p.fieldsCapacity = plan.Capacity()
+ p.dynamicFields = nil
+ p.allDynamicFields = plan.AllFields
+
+ p.wantStar = plan.Needs("*")
+ p.wantLine = plan.Needs("$line")
+ p.wantEmpty = plan.Needs("$empty")
+ p.wantHostname = plan.Needs("$hostname")
+ p.wantServer = plan.Needs("$server")
+ p.wantTimezone = plan.Needs("$timezone")
+ p.wantTimeOffset = plan.Needs("$timeoffset")
+ p.wantSeverity = plan.Needs("$severity")
+ p.wantLogLevel = plan.Needs("$loglevel")
+ p.wantTime = plan.Needs("$time")
+ p.wantDate = plan.Needs("$date")
+ p.wantHour = plan.Needs("$hour")
+ p.wantMinute = plan.Needs("$minute")
+ p.wantSecond = plan.Needs("$second")
+ p.wantPID = plan.Needs("$pid")
+ p.wantCaller = plan.Needs("$caller")
+ p.wantCPUs = plan.Needs("$cpus")
+ p.wantGoroutines = plan.Needs("$goroutines")
+ p.wantCGOCalls = plan.Needs("$cgocalls")
+ p.wantLoadAvg = plan.Needs("$loadavg")
+ p.wantUptime = plan.Needs("$uptime")
+
+ if plan.AllFields {
+ return
+ }
+
+ p.dynamicFields = make(map[string]struct{}, len(plan.Fields))
+ for field := range plan.Fields {
+ p.dynamicFields[field] = struct{}{}
+ }
}
diff --git a/internal/mapr/logformat/default_benchmark_test.go b/internal/mapr/logformat/default_benchmark_test.go
new file mode 100644
index 0000000..2eb468e
--- /dev/null
+++ b/internal/mapr/logformat/default_benchmark_test.go
@@ -0,0 +1,44 @@
+package logformat
+
+import (
+ "testing"
+
+ "github.com/mimecast/dtail/internal/mapr"
+)
+
+func BenchmarkDefaultParserMakeFields(b *testing.B) {
+ input := "INFO|20211002-072342|1|default_benchmark_test.go:0|8|14|7|0.21|471h0m21s|" +
+ "MAPREDUCE:STATS|foo=bar|bar=baz|qux=quux|alpha=beta|gamma=delta"
+
+ b.Run("all_fields", func(b *testing.B) {
+ parser, err := NewParser("default", nil)
+ if err != nil {
+ b.Fatalf("Unable to create parser: %s", err.Error())
+ }
+
+ b.ReportAllocs()
+ for i := 0; i < b.N; i++ {
+ if _, err := parser.MakeFields(input, ""); err != nil {
+ b.Fatalf("Unable to parse input: %s", err.Error())
+ }
+ }
+ })
+
+ b.Run("query_specific", func(b *testing.B) {
+ q, err := mapr.NewQuery(`select count(foo) from STATS where bar eq "baz"`)
+ if err != nil {
+ b.Fatalf("Unable to create query: %s", err.Error())
+ }
+ parser, err := NewParser("default", q)
+ if err != nil {
+ b.Fatalf("Unable to create parser: %s", err.Error())
+ }
+
+ b.ReportAllocs()
+ for i := 0; i < b.N; i++ {
+ if _, err := parser.MakeFields(input, ""); err != nil {
+ b.Fatalf("Unable to parse input: %s", err.Error())
+ }
+ }
+ })
+}
diff --git a/internal/mapr/logformat/default_test.go b/internal/mapr/logformat/default_test.go
index 4eae81b..992a052 100644
--- a/internal/mapr/logformat/default_test.go
+++ b/internal/mapr/logformat/default_test.go
@@ -3,6 +3,8 @@ package logformat
import (
"fmt"
"testing"
+
+ "github.com/mimecast/dtail/internal/mapr"
)
func TestDefaultLogFormat(t *testing.T) {
@@ -23,7 +25,7 @@ func TestDefaultLogFormat(t *testing.T) {
}
for _, input := range inputs {
- fields, err := parser.MakeFields(input)
+ fields, err := parser.MakeFields(input, "")
if err != nil {
t.Errorf("Parser unable to make fields: %s", err.Error())
@@ -86,12 +88,46 @@ func TestDefaultLogFormat(t *testing.T) {
}
}
- fields, err := parser.MakeFields("foozoo=bar|bazbay")
+ fields, err := parser.MakeFields("foozoo=bar|bazbay", "")
if err != nil && err != ErrIgnoreFields {
- t.Errorf(err.Error())
+ t.Errorf("%s", err.Error())
}
if _, ok := fields["foo"]; ok {
t.Errorf("Expected fiending field 'foo', but found it\n")
}
}
+
+func TestDefaultLogFormatQuerySpecificFields(t *testing.T) {
+ q, err := mapr.NewQuery(`select count(foo) from STATS where $hostname eq "testhost"`)
+ if err != nil {
+ t.Fatalf("Unable to create query: %s", err.Error())
+ }
+
+ parser, err := NewParser("default", q)
+ if err != nil {
+ t.Fatalf("Unable to create parser: %s", err.Error())
+ }
+
+ fields, err := parser.MakeFields(
+ "INFO|20211002-072342|1|default_test.go:0|8|14|7|0.21|471h0m21s|MAPREDUCE:STATS|foo=bar|bar=baz",
+ "",
+ )
+ if err != nil {
+ t.Fatalf("Parser unable to make fields: %s", err.Error())
+ }
+
+ requiredFields := []string{"foo", "$hostname"}
+ for _, field := range requiredFields {
+ if _, ok := fields[field]; !ok {
+ t.Errorf("Expected query-specific field '%s' to be present", field)
+ }
+ }
+
+ omittedFields := []string{"bar", "$time", "$pid", "$line"}
+ for _, field := range omittedFields {
+ if _, ok := fields[field]; ok {
+ t.Errorf("Expected query-specific field '%s' to be omitted", field)
+ }
+ }
+}
diff --git a/internal/mapr/logformat/delimited.go b/internal/mapr/logformat/delimited.go
new file mode 100644
index 0000000..2fa0639
--- /dev/null
+++ b/internal/mapr/logformat/delimited.go
@@ -0,0 +1,12 @@
+package logformat
+
+import "strings"
+
+func scanDelimitedField(input string, start int, delimiter byte) (token string, next int, done bool) {
+ index := strings.IndexByte(input[start:], delimiter)
+ if index < 0 {
+ return input[start:], len(input), true
+ }
+ end := start + index
+ return input[start:end], end + 1, false
+}
diff --git a/internal/mapr/logformat/generic.go b/internal/mapr/logformat/generic.go
index 32d9b4a..4f6d58b 100644
--- a/internal/mapr/logformat/generic.go
+++ b/internal/mapr/logformat/generic.go
@@ -4,6 +4,8 @@ type genericParser struct {
defaultParser
}
+var _ Parser = (*genericParser)(nil)
+
func newGenericParser(hostname, timeZoneName string, timeZoneOffset int) (*genericParser, error) {
defaultParser, err := newDefaultParser(hostname, timeZoneName, timeZoneOffset)
if err != nil {
@@ -12,16 +14,9 @@ func newGenericParser(hostname, timeZoneName string, timeZoneOffset int) (*gener
return &genericParser{defaultParser: *defaultParser}, nil
}
-func (p *genericParser) MakeFields(maprLine string) (map[string]string, error) {
- fields := make(map[string]string, 3)
-
- fields["*"] = "*"
- fields["$hostname"] = p.hostname
- fields["$server"] = p.hostname
- fields["$line"] = maprLine
- fields["$empty"] = ""
- fields["$timezone"] = p.timeZoneName
- fields["$timeoffset"] = p.timeZoneOffset
+func (p *genericParser) MakeFields(maprLine, _ string) (map[string]string, error) {
+ fields := make(map[string]string, p.fieldsCapacity)
+ p.addDefaultFields(fields, maprLine)
return fields, nil
}
diff --git a/internal/mapr/logformat/generickv.go b/internal/mapr/logformat/generickv.go
index 9c3de92..5a0595d 100644
--- a/internal/mapr/logformat/generickv.go
+++ b/internal/mapr/logformat/generickv.go
@@ -1,15 +1,13 @@
package logformat
-import (
- "strings"
-
- "github.com/mimecast/dtail/internal/protocol"
-)
+import "github.com/mimecast/dtail/internal/protocol"
type genericKVParser struct {
defaultParser
}
+var _ Parser = (*genericKVParser)(nil)
+
func newGenericKVParser(hostname, timeZoneName string, timeZoneOffset int) (*genericKVParser, error) {
defaultParser, err := newDefaultParser(hostname, timeZoneName, timeZoneOffset)
if err != nil {
@@ -18,25 +16,21 @@ func newGenericKVParser(hostname, timeZoneName string, timeZoneOffset int) (*gen
return &genericKVParser{defaultParser: *defaultParser}, nil
}
-func (p *genericKVParser) MakeFields(maprLine string) (map[string]string, error) {
- splitted := strings.Split(maprLine, protocol.FieldDelimiter)
- fields := make(map[string]string, len(splitted))
+func (p *genericKVParser) MakeFields(maprLine, _ string) (map[string]string, error) {
+ fields := make(map[string]string, p.fieldsCapacity)
+ p.addDefaultFields(fields, maprLine)
+ start := 0
+ delimiter := protocol.FieldDelimiter[0]
- fields["*"] = "*"
- fields["$line"] = maprLine
- fields["$empty"] = ""
- fields["$hostname"] = p.hostname
- fields["$server"] = p.hostname
- fields["$timezone"] = p.timeZoneName
- fields["$timeoffset"] = p.timeZoneOffset
-
- for _, kv := range splitted[0:] {
- keyAndValue := strings.SplitN(kv, "=", 2)
- if len(keyAndValue) != 2 {
- //dlog.Common.Debug("Unable to parse key-value token, ignoring it", kv)
+ for {
+ token, next, done := scanDelimitedField(maprLine, start, delimiter)
+ if err := p.addKeyValueField(fields, token); err != nil {
continue
}
- fields[keyAndValue[0]] = keyAndValue[1]
+ if done {
+ break
+ }
+ start = next
}
return fields, nil
diff --git a/internal/mapr/logformat/mimecast.go b/internal/mapr/logformat/mimecast.go
index cf6b333..249e2db 100644
--- a/internal/mapr/logformat/mimecast.go
+++ b/internal/mapr/logformat/mimecast.go
@@ -1,5 +1,4 @@
//go:build !proprietary
-// +build !proprietary
package logformat
@@ -10,6 +9,8 @@ var ErrMimecastNotAvailable error = errors.New("The mimecast logformat is not av
type mimecastParser struct{}
+var _ Parser = (*mimecastParser)(nil)
+
func newMimecastParser(hostname, timeZoneName string, timeZoneOffset int) (*mimecastParser, error) {
return &mimecastParser{}, ErrMimecastNotAvailable
}
@@ -18,6 +19,6 @@ func newMimecastGenericParser(hostname, timeZoneName string, timeZoneOffset int)
return &mimecastParser{}, ErrMimecastNotAvailable
}
-func (p *mimecastParser) MakeFields(maprLine string) (map[string]string, error) {
+func (p *mimecastParser) MakeFields(maprLine, _ string) (map[string]string, error) {
return nil, ErrMimecastNotAvailable
}
diff --git a/internal/mapr/logformat/parser.go b/internal/mapr/logformat/parser.go
index 37d7a63..0556c8b 100644
--- a/internal/mapr/logformat/parser.go
+++ b/internal/mapr/logformat/parser.go
@@ -3,6 +3,8 @@ package logformat
import (
"errors"
"fmt"
+ "strings"
+ "sync"
"time"
"github.com/mimecast/dtail/internal/config"
@@ -14,8 +16,72 @@ var ErrIgnoreFields error = errors.New("Ignore this field set")
// Parser is used to parse the mapreduce information from the server log files.
type Parser interface {
- // MakeFields creates a field map from an input log line.
- MakeFields(string) (map[string]string, error)
+ // MakeFields creates a field map from an input log line. The sourceID
+ // identifies the log file (or stream) the line belongs to so that
+ // stateful parsers (e.g. CSV with per-file headers) can key their
+ // state per source instead of smearing it across every file in a
+ // session.
+ MakeFields(maprLine, sourceID string) (map[string]string, error)
+}
+
+type queryAwareParser interface {
+ setQuery(*mapr.Query)
+}
+
+// ParserFactory builds a Parser for a specific log format.
+type ParserFactory func(hostname, timeZoneName string, timeZoneOffset int) (Parser, error)
+
+var parserFactories = make(map[string]ParserFactory)
+var parserFactoriesMu sync.RWMutex
+
+func init() {
+ registerBuiltInParsers()
+}
+
+// RegisterParser registers or replaces a parser factory for a log format name.
+func RegisterParser(logFormatName string, factory ParserFactory) error {
+ name := strings.TrimSpace(logFormatName)
+ if name == "" {
+ return errors.New("log format name cannot be empty")
+ }
+ if factory == nil {
+ return errors.New("parser factory cannot be nil")
+ }
+
+ parserFactoriesMu.Lock()
+ defer parserFactoriesMu.Unlock()
+ parserFactories[name] = factory
+ return nil
+}
+
+func getParserFactory(logFormatName string) (ParserFactory, bool) {
+ parserFactoriesMu.RLock()
+ defer parserFactoriesMu.RUnlock()
+ factory, found := parserFactories[logFormatName]
+ return factory, found
+}
+
+func registerBuiltInParsers() {
+ mustRegisterParser("generic", wrapParserFactory(newGenericParser))
+ mustRegisterParser("generickv", wrapParserFactory(newGenericKVParser))
+ mustRegisterParser("csv", wrapParserFactory(newCSVParser))
+ mustRegisterParser("mimecast", wrapParserFactory(newMimecastParser))
+ mustRegisterParser("mimecastgeneric", wrapParserFactory(newMimecastGenericParser))
+ mustRegisterParser("default", wrapParserFactory(newDefaultParser))
+ mustRegisterParser("custom1", wrapParserFactory(newCustom1Parser))
+ mustRegisterParser("custom2", wrapParserFactory(newCustom2Parser))
+}
+
+func mustRegisterParser(logFormatName string, factory ParserFactory) {
+ if err := RegisterParser(logFormatName, factory); err != nil {
+ panic(err)
+ }