From ab9b04971fbe16977cc224fc574c5b7622dbd3c8 Mon Sep 17 00:00:00 2001 From: Nikolay Sivko Date: Mon, 10 Aug 2026 14:23:54 -0300 Subject: [PATCH] rate limit log pattern extraction to cap CPU usage under log storms --- containers/container.go | 6 +++--- flags/flags.go | 5 +++-- go.mod | 2 +- go.sum | 4 ++-- logs/otel.go | 10 ++++++++++ windows/containers/container.go | 4 ++-- 6 files changed, 21 insertions(+), 10 deletions(-) diff --git a/containers/container.go b/containers/container.go index 8d93e69..0dbaa79 100644 --- a/containers/container.go +++ b/containers/container.go @@ -1249,7 +1249,7 @@ func (c *Container) runLogParser(logPath string) { return } ch := make(chan logparser.LogEntry) - parser := logparser.NewParser(ch, nil, logs.OtelLogEmitter(containerId), logs.MultilineCollectorTimeout, *flags.LogPatternsPerContainer, !*flags.DisableJsonLogParsing) + parser := logparser.NewParser(ch, nil, logs.OtelLogEmitter(containerId), logs.MultilineCollectorTimeout, *flags.LogPatternsPerContainer, !*flags.DisableJsonLogParsing, logs.PatternExtractionRateLimiter()) reader, err := logs.NewTailReader(proc.HostPath(logPath), ch) if err != nil { klog.Warningln(err) @@ -1268,7 +1268,7 @@ func (c *Container) runLogParser(logPath string) { klog.Warningln(err) return } - parser := logparser.NewParser(ch, nil, logs.OtelLogEmitter(containerId), logs.MultilineCollectorTimeout, *flags.LogPatternsPerContainer, !*flags.DisableJsonLogParsing) + parser := logparser.NewParser(ch, nil, logs.OtelLogEmitter(containerId), logs.MultilineCollectorTimeout, *flags.LogPatternsPerContainer, !*flags.DisableJsonLogParsing, logs.PatternExtractionRateLimiter()) stop := func() { JournaldUnsubscribe(c.metadata.systemd.Unit) } @@ -1284,7 +1284,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), logs.MultilineCollectorTimeout, *flags.LogPatternsPerContainer, !*flags.DisableJsonLogParsing) + parser := logparser.NewParser(ch, c.metadata.logDecoder, logs.OtelLogEmitter(containerId), logs.MultilineCollectorTimeout, *flags.LogPatternsPerContainer, !*flags.DisableJsonLogParsing, logs.PatternExtractionRateLimiter()) 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 257f30e..a97db03 100644 --- a/flags/flags.go +++ b/flags/flags.go @@ -35,8 +35,9 @@ var ( InstanceType = kingpin.Flag("instance-type", "`instance_type` label for `node_cloud_info` metric").Envar(envar("INSTANCE_TYPE")).String() InstanceLifeCycle = kingpin.Flag("instance-life-cycle", "`instance_life_cycle` label for `node_cloud_info` metric").Envar(envar("INSTANCE_LIFE_CYCLE")).String() - LogPatternsPerContainer = kingpin.Flag("log-patterns-per-container", "Max unique log patterns per container per level").Default("256").Envar(envar("LOG_PATTERNS_PER_CONTAINER")).Int() - MaxLabelLength = kingpin.Flag("max-label-length", "Maximum length of a metric label value").Default("4096").Envar(envar("MAX_LABEL_LENGTH")).Int() + LogPatternsPerContainer = kingpin.Flag("log-patterns-per-container", "Max unique log patterns per container per level").Default("256").Envar(envar("LOG_PATTERNS_PER_CONTAINER")).Int() + LogPatternExtractionLimit = kingpin.Flag("log-pattern-extraction-limit", "Max log messages per second per container for which patterns are extracted. Over-limit messages are counted under a dedicated 'event was sampled' pattern (0 - unlimited)").Default("100").Envar(envar("LOG_PATTERN_EXTRACTION_LIMIT")).Float64() + MaxLabelLength = kingpin.Flag("max-label-length", "Maximum length of a metric label value").Default("4096").Envar(envar("MAX_LABEL_LENGTH")).Int() CollectorEndpoint = kingpin.Flag("collector-endpoint", "A base endpoint URL for metrics, traces, logs, and profiles").Envar(envar("COLLECTOR_ENDPOINT")).URL() ApiKey = kingpin.Flag("api-key", "Coroot API key").Envar(envar("API_KEY")).String() diff --git a/go.mod b/go.mod index d085c1c..1ba070a 100644 --- a/go.mod +++ b/go.mod @@ -11,7 +11,7 @@ require ( github.com/containerd/cgroups v1.1.0 github.com/containerd/containerd v1.7.29 github.com/coreos/go-systemd/v22 v22.7.0 - github.com/coroot/logparser v1.3.2 + github.com/coroot/logparser v1.4.0 github.com/docker/docker v27.4.0+incompatible github.com/florianl/go-conntrack v0.3.0 github.com/go-kit/log v0.2.1 diff --git a/go.sum b/go.sum index f7ddce1..98c0710 100644 --- a/go.sum +++ b/go.sum @@ -92,8 +92,8 @@ github.com/coreos/go-systemd/v22 v22.7.0 h1:LAEzFkke61DFROc7zNLX/WA2i5J8gYqe0rSj github.com/coreos/go-systemd/v22 v22.7.0/go.mod h1:xNUYtjHu2EDXbsxz1i41wouACIwT7Ybq9o0BQhMwD0w= github.com/coroot/dotnetdiag v1.2.2 h1:PVP/By8o+xhPjfVolJYcjHLbFQInM7pkaD6/otPLc8Q= github.com/coroot/dotnetdiag v1.2.2/go.mod h1:veXCMlFzm1yNl7wwJb/ZLxO4WbzhDBoy1VG1XtkH2ls= -github.com/coroot/logparser v1.3.2 h1:5osbFfws9/AYUL2mdX3EU0Gjh4ZJ2NVnPN3M+sqbuQE= -github.com/coroot/logparser v1.3.2/go.mod h1:/7qHU4/I4zWRYIzRchQPehlTzbcMv5HV6cwBqg2zl6I= +github.com/coroot/logparser v1.4.0 h1:/b+XnAh7kuKR1mVSekPCSPpKloWx6wJveHa0HDcXueY= +github.com/coroot/logparser v1.4.0/go.mod h1:5mgr/LIFAEuhxLgnFTtEXVWKthT1l2ZTsPHF6oOGc/E= github.com/coroot/pyroscope/ebpf v0.0.0-20260804213318-758a0e72af3a h1:whKF1ZhIGcaAn9ZV5hbBSkrw7C5OBgWasZ0bYktOW+U= github.com/coroot/pyroscope/ebpf v0.0.0-20260804213318-758a0e72af3a/go.mod h1:IepHM9FJ0n3n3k+ZV23Y7vNAfvWI7LDuLqWPO4rB6sQ= github.com/cpuguy83/go-md2man/v2 v2.0.4/go.mod h1:tgQtvFlXSQOSOSIRvRPT7W67SCa46tRHOmNcaadrF8o= diff --git a/logs/otel.go b/logs/otel.go index 826808f..c059d6e 100644 --- a/logs/otel.go +++ b/logs/otel.go @@ -13,11 +13,13 @@ import ( otelLogs "github.com/agoda-com/opentelemetry-logs-go/logs" sdk "github.com/agoda-com/opentelemetry-logs-go/sdk/logs" "github.com/coroot/coroot-node-agent/common" + "github.com/coroot/coroot-node-agent/flags" "github.com/coroot/logparser" "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" ) @@ -25,6 +27,14 @@ import ( // multi-line log message is complete. const MultilineCollectorTimeout = time.Second +func PatternExtractionRateLimiter() *rate.Limiter { + limit := *flags.LogPatternExtractionLimit + if limit <= 0 { + return nil + } + return rate.NewLimiter(rate.Limit(limit), int(limit*10)) +} + type Config struct { Endpoint *url.URL AuthHeaders map[string]string diff --git a/windows/containers/container.go b/windows/containers/container.go index 0b4d8f8..33db75b 100644 --- a/windows/containers/container.go +++ b/windows/containers/container.go @@ -265,7 +265,7 @@ func (c *Container) eventLogInput() chan logparser.LogEntry { return c.eventLogCh } ch := make(chan logparser.LogEntry, 100) - parser := logparser.NewParser(ch, nil, logs.OtelLogEmitter(c.ID), logs.MultilineCollectorTimeout, *flags.LogPatternsPerContainer, !*flags.DisableJsonLogParsing) + parser := logparser.NewParser(ch, nil, logs.OtelLogEmitter(c.ID), logs.MultilineCollectorTimeout, *flags.LogPatternsPerContainer, !*flags.DisableJsonLogParsing, logs.PatternExtractionRateLimiter()) c.addLogParserLocked(logSourceEventLog, logs.NewPipeline(parser, nil)) c.eventLogCh = ch return ch @@ -279,7 +279,7 @@ func (c *Container) startLogTailer() { defer c.lock.Unlock() c.stopLogParserLocked(logSourceStdout) ch := make(chan logparser.LogEntry, 100) - parser := logparser.NewParser(ch, logparser.DockerJsonDecoder{}, logs.OtelLogEmitter(c.ID), logs.MultilineCollectorTimeout, *flags.LogPatternsPerContainer, !*flags.DisableJsonLogParsing) + parser := logparser.NewParser(ch, logparser.DockerJsonDecoder{}, logs.OtelLogEmitter(c.ID), logs.MultilineCollectorTimeout, *flags.LogPatternsPerContainer, !*flags.DisableJsonLogParsing, logs.PatternExtractionRateLimiter()) reader, err := logs.NewTailReader(c.logPath, ch) if err != nil { parser.Stop()