From 92475a2a1e73fb12d02b3ff8401465bc0e198b31 Mon Sep 17 00:00:00 2001 From: Nikolay Sivko Date: Mon, 27 Jul 2026 11:19:35 -0300 Subject: [PATCH 1/8] extract attributes from JSON-formatted logs (cherry picked from commit e1cd714d00d10de91b4596835cdc8d962d83a215) Conflicts: this fork imports github.com/nudgebee/logparser, which is bumped to the commit that merges coroot/logparser v1.4.2 (JSON parser, pattern rate limit). NewParser keeps this fork's SensitiveConfig argument; the rate limiter is nil until the next port. windows/ is not part of this fork. --- containers/container.go | 6 +++--- flags/flags.go | 13 +++++++------ go.mod | 2 +- go.sum | 4 ++-- logs/otel.go | 12 ++++++++---- 5 files changed, 21 insertions(+), 16 deletions(-) diff --git a/containers/container.go b/containers/container.go index 3ad9ede8..67c35bce 100644 --- a/containers/container.go +++ b/containers/container.go @@ -2112,7 +2112,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, nil, sensitiveCfg) reader, err := logs.NewTailReader(proc.HostPath(logPath), ch) if err != nil { klog.Warningln(err) @@ -2132,7 +2132,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, nil, sensitiveCfg) stop := func() { JournaldUnsubscribe(c.metadata.systemd.Unit) } @@ -2149,7 +2149,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, nil, sensitiveCfg) reader, err := logs.NewTailReader(proc.HostPath(c.metadata.logPath), ch) if err != nil { klog.Warningln(err) diff --git a/flags/flags.go b/flags/flags.go index ad95f9a9..c436074a 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 diff --git a/go.mod b/go.mod index 590a6559..ba2d2dd2 100644 --- a/go.mod +++ b/go.mod @@ -21,7 +21,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-20261008193348-2d6d6c935397 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 c0322e78..6ce08160 100644 --- a/go.sum +++ b/go.sum @@ -301,8 +301,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-20261008193348-2d6d6c935397 h1:y+19/3dvhfbrFLeYl5v6KZNHLiD3soQ6GsNnVTFlO3E= +github.com/nudgebee/logparser v0.0.0-20261008193348-2d6d6c935397/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..6933173f 100644 --- a/logs/otel.go +++ b/logs/otel.go @@ -64,7 +64,7 @@ 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,6 +80,12 @@ func OtelLogEmitter(containerId string) logparser.OnMsgCallbackF { severityNumber = otelLogs.DEBUG } + attrs := make([]attribute.KeyValue, 0, len(attributes)+1) + attrs = append(attrs, attribute.Key("pattern.hash").String(patternHash)) + for k, v := range attributes { + attrs = append(attrs, attribute.Key(k).String(v)) + } + otelLogger.Emit( otelLogs.NewLogRecord(otelLogs.LogRecordConfig{ ObservedTimestamp: ts, @@ -90,9 +96,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, }), ) } From cf11c788b16456120bfaf9494884f2177e13c54f Mon Sep 17 00:00:00 2001 From: Nikolay Sivko Date: Tue, 4 Aug 2026 14:04:53 -0300 Subject: [PATCH 2/8] populate OTLP trace context from trace_id/span_id fields in JSON logs (cherry picked from commit 18cfe3ce6740e741a10b4d5a0a791ee114474549) --- logs/otel.go | 71 ++++++++++++++++++++++++++++++++++++++++++++++++---- 1 file changed, 66 insertions(+), 5 deletions(-) diff --git a/logs/otel.go b/logs/otel.go index 6933173f..218f7f41 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,6 +16,7 @@ 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" "k8s.io/klog/v2" ) @@ -60,6 +62,67 @@ 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 @@ -80,15 +143,13 @@ func OtelLogEmitter(containerId string) logparser.OnMsgCallbackF { severityNumber = otelLogs.DEBUG } - attrs := make([]attribute.KeyValue, 0, len(attributes)+1) - attrs = append(attrs, attribute.Key("pattern.hash").String(patternHash)) - for k, v := range attributes { - attrs = append(attrs, attribute.Key(k).String(v)) - } + attrs, traceId, spanId := logRecordAttrs(patternHash, attributes) otelLogger.Emit( otelLogs.NewLogRecord(otelLogs.LogRecordConfig{ ObservedTimestamp: ts, + TraceId: traceId, + SpanId: spanId, SeverityText: &severityText, SeverityNumber: &severityNumber, Body: &msg, From 85995fd2f74ef0980a62c979f6d2f646170df231 Mon Sep 17 00:00:00 2001 From: Nikolay Sivko Date: Thu, 6 Aug 2026 11:23:01 -0300 Subject: [PATCH 3/8] retry opening a container log file if it doesn't exist yet (cherry picked from commit cd65ce5581bc73e989a8941e3a86c47d710349de) Conflict: this fork closes the file when Stat or Seek fails; kept. --- logs/tail_reader.go | 24 +++++++++++++++--------- logs/tail_reader_test.go | 9 +++++++++ 2 files changed, 24 insertions(+), 9 deletions(-) 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") From b3dc81b495c3e79a31dbd562b70638bf00600a0f Mon Sep 17 00:00:00 2001 From: Nikolay Sivko Date: Mon, 10 Aug 2026 14:23:54 -0300 Subject: [PATCH 4/8] rate limit log pattern extraction to cap CPU usage under log storms (cherry picked from commit ab9b04971fbe16977cc224fc574c5b7622dbd3c8) Conflicts: kept this fork's logs.Init, flag style and SensitiveConfig argument. With the logparser sync, an over-limit message is still exported (under the sampled pattern); only extraction is skipped. --- containers/container.go | 6 +++--- flags/flags.go | 21 +++++++++++---------- logs/otel.go | 11 +++++++++++ 3 files changed, 25 insertions(+), 13 deletions(-) diff --git a/containers/container.go b/containers/container.go index 67c35bce..49326564 100644 --- a/containers/container.go +++ b/containers/container.go @@ -2112,7 +2112,7 @@ func (c *Container) runLogParser(logPath string) { return } ch := make(chan logparser.LogEntry) - parser := logparser.NewParser(ch, nil, logs.OtelLogEmitter(containerId), multilineCollectorTimeout, *flags.LogPatternsPerContainer, !*flags.DisableJsonLogParsing, nil, 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 +2132,7 @@ func (c *Container) runLogParser(logPath string) { klog.Warningln(err) return } - parser := logparser.NewParser(ch, nil, logs.OtelLogEmitter(containerId), multilineCollectorTimeout, *flags.LogPatternsPerContainer, !*flags.DisableJsonLogParsing, nil, 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 +2149,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, !*flags.DisableJsonLogParsing, nil, 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) diff --git a/flags/flags.go b/flags/flags.go index c436074a..a1a999a3 100644 --- a/flags/flags.go +++ b/flags/flags.go @@ -44,16 +44,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("100").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/logs/otel.go b/logs/otel.go index 218f7f41..8f83ffc9 100644 --- a/logs/otel.go +++ b/logs/otel.go @@ -17,11 +17,22 @@ import ( "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 + } + return rate.NewLimiter(rate.Limit(limit), int(limit*10)) +} + func Init(machineId, hostname, version string) { endpointUrl := *flags.LogsEndpoint if endpointUrl == nil { From d34db8bf597b7f33c7a1553a6530b7ff4feb91ca Mon Sep 17 00:00:00 2001 From: Nikolay Sivko Date: Wed, 30 Sep 2026 17:18:42 -0300 Subject: [PATCH 5/8] logs: find the log file when the open() fd is already gone (freopen), bump logparser to v1.4.2 (cherry picked from commit e4be15099c51ac92116a3a231b8775d54c1a2eba) In this fork, tailing log files a process opens stays behind --enable-dynamic-log-tailing; the freopen fallback is inside that gate. The logparser bump is the nudgebee/logparser sync from the first commit. --- containers/container.go | 36 +++++++++++++++++++++++++++++++----- 1 file changed, 31 insertions(+), 5 deletions(-) diff --git a/containers/container.go b/containers/container.go index 49326564..832b09e5 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() @@ -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 +} From f63389b22e170b120ec35b13a0e487f5011b1914 Mon Sep 17 00:00:00 2001 From: mayankpande88 Date: Fri, 9 Oct 2026 01:16:57 +0530 Subject: [PATCH 6/8] test: pass parseJson and limiter to logparser.NewParser --- common/log_parser_test.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) 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, From f5ebf171940912f951dc1bdad9aaa7549d16a285 Mon Sep 17 00:00:00 2001 From: mayankpande88 Date: Fri, 9 Oct 2026 01:22:45 +0530 Subject: [PATCH 7/8] default --log-pattern-extraction-limit to 0 (unlimited) In a storm of about 59,000 error lines in 30s, a limit of 100/s saved 139 ms of agent CPU (about 4.6 millicores), but counted 55,070 of the messages under the 'event was sampled' pattern: the storm's own pattern showed about 4,000. Pattern extraction is cheap here, and log pattern spikes are what incident detection looks at. Keep the flag for installs that need the cap. --- flags/flags.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/flags/flags.go b/flags/flags.go index a1a999a3..417bb8ab 100644 --- a/flags/flags.go +++ b/flags/flags.go @@ -53,7 +53,7 @@ var ( 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("100").Envar("LOG_PATTERN_EXTRACTION_LIMIT").Float64() + 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() From e3456454f8f81d1a4d705d8c7a00c67e9b5faacb Mon Sep 17 00:00:00 2001 From: mayankpande88 Date: Fri, 9 Oct 2026 01:38:49 +0530 Subject: [PATCH 8/8] fix: keep the pattern extraction limiter's burst at least 1; bump logparser With a limit below 0.1/s, int(limit*10) was 0 and a zero-burst limiter rejects every event, so pattern extraction stopped entirely. Also pin nudgebee/logparser to the commit that decodes a JSON line once and accepts CRLF line endings. --- go.mod | 2 +- go.sum | 4 ++-- logs/otel.go | 4 +++- 3 files changed, 6 insertions(+), 4 deletions(-) diff --git a/go.mod b/go.mod index ba2d2dd2..09d77b03 100644 --- a/go.mod +++ b/go.mod @@ -21,7 +21,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-20261008193348-2d6d6c935397 + 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 6ce08160..f40a2a9f 100644 --- a/go.sum +++ b/go.sum @@ -301,8 +301,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-20261008193348-2d6d6c935397 h1:y+19/3dvhfbrFLeYl5v6KZNHLiD3soQ6GsNnVTFlO3E= -github.com/nudgebee/logparser v0.0.0-20261008193348-2d6d6c935397/go.mod h1:mUWBZ/SsHHevKDuVZ7MzFwIT5oLprAFRlp8oYq/Q/0U= +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 8f83ffc9..5614ade3 100644 --- a/logs/otel.go +++ b/logs/otel.go @@ -30,7 +30,9 @@ func PatternExtractionRateLimiter() *rate.Limiter { if limit <= 0 { return nil } - return rate.NewLimiter(rate.Limit(limit), int(limit*10)) + // 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) {