diff --git a/common/log_parser_test.go b/common/log_parser_test.go index af2f6857..5f72c8cf 100644 --- a/common/log_parser_test.go +++ b/common/log_parser_test.go @@ -12,7 +12,7 @@ import ( func TestLogParser(t *testing.T) { const defaultPatternsPerLevel = 256 ch := make(chan logparser.LogEntry) - parser := logparser.NewParser(ch, nil, nil, 1*time.Second, defaultPatternsPerLevel, logparser.SensitiveConfig{ + parser := logparser.NewParser(ch, nil, nil, 1*time.Second, defaultPatternsPerLevel, false, nil, logparser.SensitiveConfig{ Enabled: true, MinConfidence: "high", MaxDetections: 100, diff --git a/containers/container.go b/containers/container.go index e9e0eb6c..43493309 100644 --- a/containers/container.go +++ b/containers/container.go @@ -682,7 +682,8 @@ func (c *Container) onFileOpen(pid uint32, fd uint64, mnt uint64, log bool) { return } } - mntId, logPath := resolveFd(pid, fd) + info := proc.GetFdInfo(pid, fd) + mntId, logPath := resolveFd(info) func() { if mntId == "" { return @@ -706,8 +707,17 @@ func (c *Container) onFileOpen(pid uint32, fd uint64, mnt uint64, log bool) { c.lock.Unlock() } }() - if logPath != "" { - if *flags.EnableDynamicLogTailing { + if *flags.EnableDynamicLogTailing { + var logPaths []string + switch { + case logPath != "": + logPaths = []string{logPath} + case log && (info == nil || !strings.HasPrefix(info.Dest, "/var/log/")): + // The kernel saw a /var/log/ file opened, but the fd now points + // elsewhere (freopen): find the log files the process holds. + logPaths = findLogFiles(pid) + } + for _, logPath := range logPaths { c.lock.Lock() c.runLogParser(logPath) c.lock.Unlock() @@ -2112,7 +2122,7 @@ func (c *Container) runLogParser(logPath string) { return } ch := make(chan logparser.LogEntry) - parser := logparser.NewParser(ch, nil, logs.OtelLogEmitter(containerId), multilineCollectorTimeout, *flags.LogPatternsPerContainer, sensitiveCfg) + parser := logparser.NewParser(ch, nil, logs.OtelLogEmitter(containerId), multilineCollectorTimeout, *flags.LogPatternsPerContainer, !*flags.DisableJsonLogParsing, logs.PatternExtractionRateLimiter(), sensitiveCfg) reader, err := logs.NewTailReader(proc.HostPath(logPath), ch) if err != nil { klog.Warningln(err) @@ -2132,7 +2142,7 @@ func (c *Container) runLogParser(logPath string) { klog.Warningln(err) return } - parser := logparser.NewParser(ch, nil, logs.OtelLogEmitter(containerId), multilineCollectorTimeout, *flags.LogPatternsPerContainer, sensitiveCfg) + parser := logparser.NewParser(ch, nil, logs.OtelLogEmitter(containerId), multilineCollectorTimeout, *flags.LogPatternsPerContainer, !*flags.DisableJsonLogParsing, logs.PatternExtractionRateLimiter(), sensitiveCfg) stop := func() { JournaldUnsubscribe(c.metadata.systemd.Unit) } @@ -2149,7 +2159,7 @@ func (c *Container) runLogParser(logPath string) { delete(c.logParsers, "stdout/stderr") } ch := make(chan logparser.LogEntry) - parser := logparser.NewParser(ch, c.metadata.logDecoder, logs.OtelLogEmitter(containerId), multilineCollectorTimeout, *flags.LogPatternsPerContainer, sensitiveCfg) + parser := logparser.NewParser(ch, c.metadata.logDecoder, logs.OtelLogEmitter(containerId), multilineCollectorTimeout, *flags.LogPatternsPerContainer, !*flags.DisableJsonLogParsing, logs.PatternExtractionRateLimiter(), sensitiveCfg) reader, err := logs.NewTailReader(proc.HostPath(c.metadata.logPath), ch) if err != nil { klog.Warningln(err) @@ -2445,8 +2455,7 @@ func countTLSAttach(lib string, result ebpftracer.TLSAttachResult) { } } -func resolveFd(pid uint32, fd uint64) (mntId string, logPath string) { - info := proc.GetFdInfo(pid, fd) +func resolveFd(info *proc.FdInfo) (mntId string, logPath string) { if info == nil { return } @@ -2498,3 +2507,20 @@ func sampleString(v interface{}) string { } return "" } + +func findLogFiles(pid uint32) []string { + fds, err := proc.ReadFds(pid) + if err != nil { + return nil + } + var res []string + for _, fd := range fds { + if !strings.HasPrefix(fd.Dest, "/var/log/") { + continue + } + if _, logPath := resolveFd(proc.GetFdInfo(pid, fd.Fd)); logPath != "" { + res = append(res, logPath) + } + } + return res +} diff --git a/flags/flags.go b/flags/flags.go index 477ad893..1d04212c 100644 --- a/flags/flags.go +++ b/flags/flags.go @@ -9,12 +9,13 @@ import ( ) var ( - ListenAddress = kingpin.Flag("listen", "Listen address - ip:port or :port (default 0.0.0.0:80, or 127.0.0.1:10300 when --metrics-endpoint is set)").Envar("LISTEN").String() - CgroupRoot = kingpin.Flag("cgroupfs-root", "The mount point of the host cgroupfs root").Default("/sys/fs/cgroup").Envar("CGROUPFS_ROOT").String() - DisableLogParsing = kingpin.Flag("disable-log-parsing", "Disable container log parsing").Default("false").Envar("DISABLE_LOG_PARSING").Bool() - DisablePinger = kingpin.Flag("disable-pinger", "Don't ping upstreams").Default("true").Envar("DISABLE_PINGER").Bool() - DisableL7Tracing = kingpin.Flag("disable-l7-tracing", "Disable L7 tracing").Default("false").Envar("DISABLE_L7_TRACING").Bool() - EnableDotNetTracing = kingpin.Flag("enable-dotnet-tracing", "Enable .NET CLR tracing").Default("false").Envar("ENABLE_DOTNET_TRACING").Bool() + ListenAddress = kingpin.Flag("listen", "Listen address - ip:port or :port (default 0.0.0.0:80, or 127.0.0.1:10300 when --metrics-endpoint is set)").Envar("LISTEN").String() + CgroupRoot = kingpin.Flag("cgroupfs-root", "The mount point of the host cgroupfs root").Default("/sys/fs/cgroup").Envar("CGROUPFS_ROOT").String() + DisableLogParsing = kingpin.Flag("disable-log-parsing", "Disable container log parsing").Default("false").Envar("DISABLE_LOG_PARSING").Bool() + DisableJsonLogParsing = kingpin.Flag("disable-json-log-parsing", "Disable extracting the message, severity, and attributes from JSON-formatted logs").Default("false").Envar("DISABLE_JSON_LOG_PARSING").Bool() + DisablePinger = kingpin.Flag("disable-pinger", "Don't ping upstreams").Default("true").Envar("DISABLE_PINGER").Bool() + DisableL7Tracing = kingpin.Flag("disable-l7-tracing", "Disable L7 tracing").Default("false").Envar("DISABLE_L7_TRACING").Bool() + EnableDotNetTracing = kingpin.Flag("enable-dotnet-tracing", "Enable .NET CLR tracing").Default("false").Envar("ENABLE_DOTNET_TRACING").Bool() // Off by default: the only thing it produces is // container_nodejs_event_loop_blocked_time_seconds_total, which nothing // currently consumes, and attaching the probes reads the whole ELF symbol @@ -44,16 +45,17 @@ var ( Strings() EphemeralPortRange = kingpin.Flag("ephemeral-port-range", `Destination and Listen TCP ports from these ranges will be skipped, e.g. "32768-60999" or "1024-23768 30000-65535"`).Default("32768-60999").Envar("EPHEMERAL_PORT_RANGE").String() - Provider = kingpin.Flag("provider", "`provider` label for `node_cloud_info` metric").Envar("PROVIDER").String() - Region = kingpin.Flag("region", "`region` label for `node_cloud_info` metric").Envar("REGION").String() - AvailabilityZone = kingpin.Flag("availability-zone", "`availability_zone` label for `node_cloud_info` metric").Envar("AVAILABILITY_ZONE").String() - AccountId = kingpin.Flag("account-id", "`account_id` label for `node_cloud_info` metric").Envar("ACCOUNT_ID").String() - InstanceType = kingpin.Flag("instance-type", "`instance_type` label for `node_cloud_info` metric").Envar("INSTANCE_TYPE").String() - InstanceLifeCycle = kingpin.Flag("instance-life-cycle", "`instance_life_cycle` label for `node_cloud_info` metric").Envar("INSTANCE_LIFE_CYCLE").String() - LogPerSecond = kingpin.Flag("log-per-second", "The number of logs per second").Default("10.0").Envar("LOG_PER_SECOND").Float64() - LogBurst = kingpin.Flag("log-burst", "The maximum number of tokens that can be consumed in a single call to allow").Default("100").Envar("LOG_BURST").Int() - LogPatternsPerContainer = kingpin.Flag("log-patterns-per-container", "Max unique log patterns per container per level").Default("256").Envar("LOG_PATTERNS_PER_CONTAINER").Int() - MaxFQDNsPerContainer = kingpin.Flag("max-fqdns-per-container", "Max unique FQDN values per container, extras are bucketed under '~other'").Default("50").Envar("MAX_FQDNS_PER_CONTAINER").Int() + Provider = kingpin.Flag("provider", "`provider` label for `node_cloud_info` metric").Envar("PROVIDER").String() + Region = kingpin.Flag("region", "`region` label for `node_cloud_info` metric").Envar("REGION").String() + AvailabilityZone = kingpin.Flag("availability-zone", "`availability_zone` label for `node_cloud_info` metric").Envar("AVAILABILITY_ZONE").String() + AccountId = kingpin.Flag("account-id", "`account_id` label for `node_cloud_info` metric").Envar("ACCOUNT_ID").String() + InstanceType = kingpin.Flag("instance-type", "`instance_type` label for `node_cloud_info` metric").Envar("INSTANCE_TYPE").String() + InstanceLifeCycle = kingpin.Flag("instance-life-cycle", "`instance_life_cycle` label for `node_cloud_info` metric").Envar("INSTANCE_LIFE_CYCLE").String() + LogPerSecond = kingpin.Flag("log-per-second", "The number of logs per second").Default("10.0").Envar("LOG_PER_SECOND").Float64() + LogBurst = kingpin.Flag("log-burst", "The maximum number of tokens that can be consumed in a single call to allow").Default("100").Envar("LOG_BURST").Int() + LogPatternsPerContainer = kingpin.Flag("log-patterns-per-container", "Max unique log patterns per container per level").Default("256").Envar("LOG_PATTERNS_PER_CONTAINER").Int() + LogPatternExtractionLimit = kingpin.Flag("log-pattern-extraction-limit", "Max warning/error log messages per second per container for which patterns are extracted. Over-limit messages are still exported and counted, under a dedicated 'event was sampled' pattern (0 - unlimited)").Default("0").Envar("LOG_PATTERN_EXTRACTION_LIMIT").Float64() + MaxFQDNsPerContainer = kingpin.Flag("max-fqdns-per-container", "Max unique FQDN values per container, extras are bucketed under '~other'").Default("50").Envar("MAX_FQDNS_PER_CONTAINER").Int() MaxLabelLength = kingpin.Flag("max-label-length", "Maximum length of a metric label value").Default("4096").Envar("MAX_LABEL_LENGTH").Int() diff --git a/go.mod b/go.mod index 93edd7ce..c48cae3b 100644 --- a/go.mod +++ b/go.mod @@ -20,7 +20,7 @@ require ( github.com/hashicorp/golang-lru/v2 v2.0.7 github.com/jpillora/backoff v1.0.0 github.com/mdlayher/taskstats v0.0.0-20230712191918-387b3d561d14 - github.com/nudgebee/logparser v0.0.0-20260406040008-c408d0a5c69c + github.com/nudgebee/logparser v0.0.0-20261008200612-6da0f06ad2f3 github.com/opencontainers/runtime-spec v1.1.0 github.com/prometheus/client_golang v1.20.5 github.com/prometheus/client_model v0.6.2 diff --git a/go.sum b/go.sum index d3064d40..a99293a4 100644 --- a/go.sum +++ b/go.sum @@ -295,8 +295,8 @@ github.com/morikuni/aec v1.0.0 h1:nP9CBfwrvYnBRgY6qfDQkygYDmYwOilePFkwzv4dU8A= github.com/morikuni/aec v1.0.0/go.mod h1:BbKIizmSmc5MMPqRYbxO4ZU0S0+P200+tUnFx7PXmsc= github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822 h1:C3w9PqII01/Oq1c1nUAm88MOHcQC9l5mIlSMApZMrHA= github.com/munnerz/goautoneg v0.0.0-20191010083416-a7dc8b61c822/go.mod h1:+n7T8mK8HuQTcFwEeznm/DIxMOiR9yIdICNftLE1DvQ= -github.com/nudgebee/logparser v0.0.0-20260406040008-c408d0a5c69c h1:drsGFCiBXp5OZm92HE0iVx8U4XhpzAxL5vVKubBGl+I= -github.com/nudgebee/logparser v0.0.0-20260406040008-c408d0a5c69c/go.mod h1:oFnM9D6YEjZzb1jy0kJ/NUkTbkZA+BgNYjpLO4n4szA= +github.com/nudgebee/logparser v0.0.0-20261008200612-6da0f06ad2f3 h1:LsztuTRobqh4hdv5d1hCRFs++NUET+AJiKSzh1za2KA= +github.com/nudgebee/logparser v0.0.0-20261008200612-6da0f06ad2f3/go.mod h1:mUWBZ/SsHHevKDuVZ7MzFwIT5oLprAFRlp8oYq/Q/0U= github.com/oklog/ulid v1.3.1 h1:EGfNDEx6MqHz8B3uNV6QAib1UR2Lm97sHi3ocA6ESJ4= github.com/oklog/ulid v1.3.1/go.mod h1:CirwcVhetQ6Lv90oh/F+FBtV6XMibvdAFo93nm5qn4U= github.com/onsi/ginkgo v1.16.5 h1:8xi0RTUf59SOSfEtZMvwTvXYMzG4gV23XVHOZiXNtnE= diff --git a/logs/otel.go b/logs/otel.go index 9aeb2bec..5614ade3 100644 --- a/logs/otel.go +++ b/logs/otel.go @@ -2,6 +2,7 @@ package logs import ( "context" + "strings" "time" otel "github.com/agoda-com/opentelemetry-logs-go" @@ -15,11 +16,25 @@ import ( "go.opentelemetry.io/otel/attribute" "go.opentelemetry.io/otel/sdk/resource" semconv "go.opentelemetry.io/otel/semconv/v1.18.0" + "go.opentelemetry.io/otel/trace" + "golang.org/x/time/rate" "k8s.io/klog/v2" ) var otelLogger otelLogs.Logger +// PatternExtractionRateLimiter caps, per container, how many warning and +// error messages per second get a log pattern extracted (nil: unlimited). +func PatternExtractionRateLimiter() *rate.Limiter { + limit := *flags.LogPatternExtractionLimit + if limit <= 0 { + return nil + } + // At least 1: with a burst of 0 the limiter rejects every event, so a + // limit below 0.1/s would turn pattern extraction off entirely. + return rate.NewLimiter(rate.Limit(limit), max(1, int(limit*10))) +} + func Init(machineId, hostname, version string) { endpointUrl := *flags.LogsEndpoint if endpointUrl == nil { @@ -60,11 +75,72 @@ func Init(machineId, hostname, version string) { otelLogger = loggerProvider.Logger("nudgebee-node-agent", otelLogs.WithInstrumentationVersion(version)) } +const ( + traceIdKey = "traceid" + spanIdKey = "spanid" +) + +var ( + traceIdKeys = []string{traceIdKey, "trace_id", "trace-id", "trace.id"} + spanIdKeys = []string{spanIdKey, "span_id", "span-id", "span.id"} +) + +func normalizeTraceContextKey(k string) string { + k = strings.ToLower(k) + switch k { + case "@tr": + return traceIdKey + case "@sp": + return spanIdKey + } + if strings.Contains(k, "parent") { // e.g. parent_span_id is not the record's span id + return "" + } + for _, s := range traceIdKeys { + if k == s || strings.HasSuffix(k, "."+s) { + return traceIdKey + } + } + for _, s := range spanIdKeys { + if k == s || strings.HasSuffix(k, "."+s) { + return spanIdKey + } + } + return "" +} + +func logRecordAttrs(patternHash string, attributes map[string]string) ([]attribute.KeyValue, *trace.TraceID, *trace.SpanID) { + var traceId *trace.TraceID + var spanId *trace.SpanID + attrs := make([]attribute.KeyValue, 0, len(attributes)+1) + attrs = append(attrs, attribute.Key("pattern.hash").String(patternHash)) + for k, v := range attributes { + switch normalizeTraceContextKey(k) { + case traceIdKey: + if traceId == nil { + if id, err := trace.TraceIDFromHex(strings.ToLower(v)); err == nil { + traceId = &id + continue + } + } + case spanIdKey: + if spanId == nil { + if id, err := trace.SpanIDFromHex(strings.ToLower(v)); err == nil { + spanId = &id + continue + } + } + } + attrs = append(attrs, attribute.Key(k).String(v)) + } + return attrs, traceId, spanId +} + func OtelLogEmitter(containerId string) logparser.OnMsgCallbackF { if otelLogger == nil { return nil } - return func(ts time.Time, level logparser.Level, patternHash string, msg string) { + return func(ts time.Time, level logparser.Level, patternHash string, msg string, attributes map[string]string) { severityText := level.String() severityNumber := otelLogs.UNSPECIFIED switch level { @@ -80,9 +156,13 @@ func OtelLogEmitter(containerId string) logparser.OnMsgCallbackF { severityNumber = otelLogs.DEBUG } + attrs, traceId, spanId := logRecordAttrs(patternHash, attributes) + otelLogger.Emit( otelLogs.NewLogRecord(otelLogs.LogRecordConfig{ ObservedTimestamp: ts, + TraceId: traceId, + SpanId: spanId, SeverityText: &severityText, SeverityNumber: &severityNumber, Body: &msg, @@ -90,9 +170,7 @@ func OtelLogEmitter(containerId string) logparser.OnMsgCallbackF { semconv.ServiceName(common.ContainerIdToOtelServiceName(containerId)), semconv.ContainerID(containerId), ), - Attributes: &[]attribute.KeyValue{ - attribute.Key("pattern.hash").String(patternHash), - }, + Attributes: &attrs, }), ) } diff --git a/logs/tail_reader.go b/logs/tail_reader.go index 57f6eb8e..ac0fb389 100644 --- a/logs/tail_reader.go +++ b/logs/tail_reader.go @@ -37,18 +37,20 @@ func NewTailReader(fileName string, ch chan<- logparser.LogEntry) (*TailReader, stopped: make(chan struct{}), } var err error - if r.file, err = os.Open(fileName); err != nil { + if r.file, err = os.Open(fileName); err != nil && !os.IsNotExist(err) { return nil, err } - if r.info, err = r.file.Stat(); err != nil { - _ = r.file.Close() - return nil, err - } - if _, err = r.file.Seek(0, io.SeekEnd); err != nil { - _ = r.file.Close() - return nil, err + if r.file != nil { // the file may not exist yet, poll() waits for it to appear + if r.info, err = r.file.Stat(); err != nil { + _ = r.file.Close() + return nil, err + } + if _, err = r.file.Seek(0, io.SeekEnd); err != nil { + _ = r.file.Close() + return nil, err + } + r.reader = bufio.NewReader(r.file) } - r.reader = bufio.NewReader(r.file) go func() { const maxPrefixLen = 64 * 1024 // 64KB cap on partial line buffer @@ -59,6 +61,10 @@ func NewTailReader(fileName string, ch chan<- logparser.LogEntry) (*TailReader, r.stopped <- struct{}{} return default: + if r.reader == nil { + r.poll(ctx) + continue + } line, err := r.reader.ReadString('\n') if err != nil { if len(prefix)+len(line) > maxPrefixLen { diff --git a/logs/tail_reader_test.go b/logs/tail_reader_test.go index 2097c2d7..49d7712f 100644 --- a/logs/tail_reader_test.go +++ b/logs/tail_reader_test.go @@ -15,6 +15,10 @@ func TestTailReader(t *testing.T) { assert.NoError(t, err) defer os.Remove(f.Name()) + // the file doesn't exist yet + err = os.Remove(f.Name()) + assert.NoError(t, err) + tailPollInterval = time.Millisecond * 100 ch := make(chan logparser.LogEntry, 10) tr, err := NewTailReader(f.Name(), ch) @@ -35,6 +39,11 @@ func TestTailReader(t *testing.T) { assert.Equal(t, expected, entry.Content) } + // the file is created + wait() + f, err = os.Create(f.Name()) + assert.NoError(t, err) + write("foo 1\n") get("foo 1")