diff options
| author | Paul Buetow <paul@buetow.org> | 2026-07-22 23:51:18 +0300 |
|---|---|---|
| committer | Paul Buetow <paul@buetow.org> | 2026-07-22 23:51:18 +0300 |
| commit | 849951be1d1a7ee9f9302006ccb187bf5b4e36f3 (patch) | |
| tree | 496c924a03a9ea6212e29bb4699e268066ebad81 /internal/mapr/logformat/parser.go | |
| parent | bf78b3abffee6d49c08ca2980156afc455994969 (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.go | 123 |
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) } |
