summaryrefslogtreecommitdiff
path: root/internal/mapr/logformat/parser.go
diff options
context:
space:
mode:
authorPaul Buetow <paul@buetow.org>2026-07-22 23:51:18 +0300
committerPaul Buetow <paul@buetow.org>2026-07-22 23:51:18 +0300
commit849951be1d1a7ee9f9302006ccb187bf5b4e36f3 (patch)
tree496c924a03a9ea6212e29bb4699e268066ebad81 /internal/mapr/logformat/parser.go
parentbf78b3abffee6d49c08ca2980156afc455994969 (diff)
feat: DTail fork — server/client feature development
Squashed development of the snonux/dtail fork's product code (internal/, cmd/) since diverging from mimecast/dtail. Major areas: - Read/output path: the former "turbo" channel-less path is now the single, default server-side read/output path for cat/grep/tail and MapReduce; the old channel-based path and its config/env toggles were removed. - MapReduce: single aggregate implementation (server + serverless) fed directly by a processor pipeline, with input-exhausted finalization via the shutdown coordinator; high-concurrency and data-race fixes. - Journal source reads (journal:unit.service) via journalctl, Linux-gated behind a journal-v1 capability. - Auth-key fast reconnect: in-memory per-user public-key cache with TTL/max-keys, registered over an authenticated session (AUTHKEY), checked before authorized_keys. - Interactive query reload (--interactive-query) with SESSION START/UPDATE generation boundaries and capability negotiation. - Client-side deadlines: --timeout / --shutdownAfter as context deadlines; follow shutdown handling. - Client logging: diagnostics-only daily log by default, opt-in payload tee via --log-payload. - Numerous correctness fixes (buffer-pool double-recycle races, EOF-sentinel leaks, glob-expansion cap, TOCTOU in CSV parsing) with accompanying unit tests. Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Diffstat (limited to 'internal/mapr/logformat/parser.go')
-rw-r--r--internal/mapr/logformat/parser.go123
1 files changed, 96 insertions, 27 deletions
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)
+ }
+}
+
+func wrapParserFactory[T Parser](factory func(string, string, int) (T, error)) ParserFactory {
+ return func(hostname, timeZoneName string, timeZoneOffset int) (Parser, error) {
+ return factory(hostname, timeZoneName, timeZoneOffset)
+ }
}
// NewParser returns a new log parser.
@@ -27,30 +93,33 @@ func NewParser(logFormatName string, query *mapr.Query) (Parser, error) {
now := time.Now()
timeZoneName, timeZoneOffset := now.Zone()
- // Extend this for adding more log formats!
- switch logFormatName {
- case "generic":
- return newGenericParser(hostname, timeZoneName, timeZoneOffset)
- case "generickv":
- return newGenericKVParser(hostname, timeZoneName, timeZoneOffset)
- case "csv":
- return newCSVParser(hostname, timeZoneName, timeZoneOffset)
- case "mimecast":
- return newMimecastParser(hostname, timeZoneName, timeZoneOffset)
- case "mimecastgeneric":
- return newMimecastGenericParser(hostname, timeZoneName, timeZoneOffset)
- case "default":
- return newDefaultParser(hostname, timeZoneName, timeZoneOffset)
- case "custom1":
- return newCustom1Parser(hostname, timeZoneName, timeZoneOffset)
- case "custom2":
- return newCustom2Parser(hostname, timeZoneName, timeZoneOffset)
- default:
- p, err := newDefaultParser(hostname, timeZoneName, timeZoneOffset)
- if err != nil {
- return p, fmt.Errorf("No '%s' mapr log format and problem creating default one: %v",
- logFormatName, err)
- }
- return p, fmt.Errorf("No '%s' mapr log format", logFormatName)
+ if parserFactory, found := getParserFactory(logFormatName); found {
+ parser, err := parserFactory(hostname, timeZoneName, timeZoneOffset)
+ configureParserQuery(parser, query)
+ return parser, err
+ }
+
+ defaultFactory, found := getParserFactory("default")
+ if !found {
+ return nil, fmt.Errorf("No '%s' mapr log format and no default parser registered", logFormatName)
+ }
+
+ p, err := defaultFactory(hostname, timeZoneName, timeZoneOffset)
+ if err != nil {
+ return p, fmt.Errorf("No '%s' mapr log format and problem creating default one: %v",
+ logFormatName, err)
+ }
+ configureParserQuery(p, query)
+ return p, fmt.Errorf("No '%s' mapr log format", logFormatName)
+}
+
+func configureParserQuery(parser Parser, query *mapr.Query) {
+ if parser == nil {
+ return
+ }
+ queryAware, ok := parser.(queryAwareParser)
+ if !ok {
+ return
}
+ queryAware.setQuery(query)
}