From bca2e038d028d0b782467ac6bee3331204b5f438 Mon Sep 17 00:00:00 2001 From: mayankpande88 Date: Wed, 7 Oct 2026 17:57:42 +0530 Subject: [PATCH 1/4] cap unique DNS domain labels per container Port of coroot/coroot-node-agent@33c46ec ("cap unique FQDN labels per container in container_dns_requests_total"). After the first --max-fqdns-per-container (default 50) distinct domains a container resolves, new domains are counted under domain="~other". Upstream keeps DNS metrics in a separate per-container vector. Here DNS is recorded through L7Stats.observe with destination and workload labels, so the cap lives in L7Stats and applies to the counter and the histogram alike. L7Stats never deletes series, so without a cap the domain label grows for the container's lifetime. --- containers/l7.go | 42 ++++++++++++++++++++++++++++++++++++++---- flags/flags.go | 1 + 2 files changed, 39 insertions(+), 4 deletions(-) diff --git a/containers/l7.go b/containers/l7.go index 0821dc3b..53b788d7 100644 --- a/containers/l7.go +++ b/containers/l7.go @@ -7,6 +7,7 @@ import ( "github.com/coroot/coroot-node-agent/common" "github.com/coroot/coroot-node-agent/ebpftracer/l7" + "github.com/coroot/coroot-node-agent/flags" "github.com/prometheus/client_golang/prometheus" "k8s.io/klog/v2" ) @@ -64,17 +65,45 @@ type L7Stats struct { latency map[l7.Protocol]*prometheus.HistogramVec initialized map[l7.Protocol]bool promConstLabels prometheus.Labels // container_id, app_id, machine_id, system_uuid, az, region + + // DNS domain label values in use; at most --max-fqdns-per-container. + seenFQDNs map[string]struct{} } +// fqdnOverflowLabel is the domain label of DNS requests for domains beyond +// the first --max-fqdns-per-container seen in a container. +const fqdnOverflowLabel = "~other" + func NewL7Stats(constLabels prometheus.Labels) L7Stats { return L7Stats{ requests: make(map[l7.Protocol]*prometheus.CounterVec), latency: make(map[l7.Protocol]*prometheus.HistogramVec), initialized: make(map[l7.Protocol]bool), promConstLabels: constLabels, + seenFQDNs: make(map[string]struct{}), } } +// dnsDomainLabel returns fqdn while the container has seen fewer than +// --max-fqdns-per-container distinct domains, and fqdnOverflowLabel for any +// new domain after that: a process resolving many unique names (crawlers, +// per-tenant hostnames) would otherwise create series without bound. +func (s *L7Stats) dnsDomainLabel(fqdn string) string { + if fqdn == "" { + return fqdn + } + s.mu.Lock() + defer s.mu.Unlock() + if _, ok := s.seenFQDNs[fqdn]; ok { + return fqdn + } + if len(s.seenFQDNs) >= *flags.MaxFQDNsPerContainer { + return fqdnOverflowLabel + } + s.seenFQDNs[fqdn] = struct{}{} + return fqdn +} + func (s *L7Stats) observe(protocol l7.Protocol, status, method, path string, duration time.Duration, key common.DestinationKey, srcWorkload common.Workload, r *l7.RequestData, traceId string) { // HTTP/1 and HTTP/2 are one set of metrics, so they must share one vector: // two vectors under the same name emit identical series for a destination @@ -110,6 +139,13 @@ func (s *L7Stats) observe(protocol l7.Protocol, status, method, path string, dur labelInterner.intern(actualDestWorkload.Namespace), } + var dnsRequestType, dnsDomain string + if metricsProtocol == l7.ProtocolDNS { + var domain string + dnsRequestType, domain, _ = l7.ParseDns(r.Payload) + dnsDomain = s.dnsDomainLabel(common.NormalizeFQDN(domain, dnsRequestType)) + } + // Protocol-specific labels for counters (keep all labels including path for HTTP) counterLabelValues := make([]string, len(labelValues)) copy(counterLabelValues, labelValues) @@ -125,8 +161,7 @@ func (s *L7Stats) observe(protocol l7.Protocol, status, method, path string, dur } counterLabelValues = append(counterLabelValues, labelInterner.intern(method)) case l7.ProtocolDNS: - requestType, domain, _ := l7.ParseDns(r.Payload) - counterLabelValues = append(counterLabelValues, labelInterner.intern(requestType), labelInterner.intern(common.NormalizeFQDN(domain, requestType))) + counterLabelValues = append(counterLabelValues, labelInterner.intern(dnsRequestType), labelInterner.intern(dnsDomain)) } // Protocol-specific labels for histograms (exclude path and method for HTTP, use grouped status to reduce cardinality) @@ -137,8 +172,7 @@ func (s *L7Stats) observe(protocol l7.Protocol, status, method, path string, dur switch metricsProtocol { case l7.ProtocolDNS: - requestType, domain, _ := l7.ParseDns(r.Payload) - histogramLabelValues = append(histogramLabelValues, labelInterner.intern(requestType), labelInterner.intern(common.NormalizeFQDN(domain, requestType))) + histogramLabelValues = append(histogramLabelValues, labelInterner.intern(dnsRequestType), labelInterner.intern(dnsDomain)) } // Map reads are safe after ensureInitialized — the protocol entry exists. diff --git a/flags/flags.go b/flags/flags.go index 3edb2d5c..9d9fdcf1 100644 --- a/flags/flags.go +++ b/flags/flags.go @@ -51,6 +51,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() + 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 aea8202bf115738d287bbcfbb0b0b07f3222e533 Mon Sep 17 00:00:00 2001 From: Nikolay Sivko Date: Fri, 8 May 2026 13:15:49 -0300 Subject: [PATCH 2/4] suppress metrics for short-lived containers (cherry picked from commit 34ea61f99a4e5acce57f582665adc87cafcb1083) Includes the createdAt half of coroot/coroot-node-agent@75d6656 ("fix: systemd services disappearing after a restart"), which only applies to this flag: age counts from the earlier of the first process start and container discovery, so a restarted unit is not hidden for the minimum age after each restart. Fork adaptations: Collect here does not hold c.lock, so the dead-pid sweep goes through onProcessExit and the age check reads its fields under RLock (youngerThan). A pid counts as dead only when /proc/ is gone; upstream treats any taskstats error as an exit, which here would close a live process's uprobes. The default is upstream's 30s: containers that live less than 30s, such as short jobs, produce no series. --- containers/container.go | 56 +++++++++++++++++++++++++++++++++++++---- flags/flags.go | 1 + 2 files changed, 52 insertions(+), 5 deletions(-) diff --git a/containers/container.go b/containers/container.go index 29d2f3fb..8edbd069 100644 --- a/containers/container.go +++ b/containers/container.go @@ -132,6 +132,7 @@ type Container struct { processes map[uint32]*Process + createdAt time.Time // when the agent found the container startedAt time.Time zombieAt time.Time @@ -246,6 +247,8 @@ func NewContainer(id ContainerID, cg *cgroup.Cgroup, md *ContainerMetadata, pid cgroup: cg, metadata: md, + createdAt: time.Now(), + processes: map[uint32]*Process{}, delaysByPid: map[uint32]Delays{}, @@ -322,6 +325,18 @@ func (c *Container) Collect(ch chan<- prometheus.Metric) { defer c.collectMu.Unlock() collectStart := time.Now() + if taskstatsClient != nil { + // Processes whose exit event was missed: without this they keep + // the container out of zombie state and its age wrong. + for _, pid := range c.updateDelays() { + c.onProcessExit(pid, false) + } + } + + if minAge := *flags.MinContainerAge; minAge > 0 && c.youngerThan(minAge) { + return + } + if c.metadata.image != "" || !c.metadata.systemd.IsEmpty() { ch <- c.gauge(metrics.ContainerInfo, 1, c.metadata.image, c.metadata.systemd.TriggeredBy, c.metadata.systemd.Type) } @@ -337,7 +352,6 @@ func (c *Container) Collect(ch chan<- prometheus.Metric) { } if taskstatsClient != nil { - c.updateDelays() ch <- c.counter(metrics.CPUDelay, float64(c.delays.cpu)/float64(time.Second)) ch <- c.counter(metrics.DiskDelay, float64(c.delays.disk)/float64(time.Second)) } @@ -590,9 +604,7 @@ func (c *Container) onProcessStart(pid uint32) *Process { return p } -func (c *Container) onProcessExit(pid uint32, oomKill bool) { - c.lock.Lock() - defer c.lock.Unlock() +func (c *Container) onProcessExitLocked(pid uint32, oomKill bool) { if p := c.processes[pid]; p != nil { c.closeProcess(pid, p) } @@ -640,6 +652,12 @@ func (c *Container) closeProcess(pid uint32, p *Process) { } } +func (c *Container) onProcessExit(pid uint32, oomKill bool) { + c.lock.Lock() + defer c.lock.Unlock() + c.onProcessExitLocked(pid, oomKill) +} + func (c *Container) onFileOpen(pid uint32, fd uint64, mnt uint64, log bool) { if mnt > 0 && !log { c.lock.Lock() @@ -1668,7 +1686,9 @@ func (c *Container) onRetransmission(src netaddr.IPPort, dst netaddr.IPPort) boo return true } -func (c *Container) updateDelays() { +// updateDelays refreshes the per-process CPU and disk delay counters and +// returns the pids that no longer exist. +func (c *Container) updateDelays() []uint32 { // Get a snapshot of PIDs under read lock to avoid concurrent map access c.lock.RLock() pids := make([]uint32, 0, len(c.processes)) @@ -1684,9 +1704,15 @@ func (c *Container) updateDelays() { diskDelay time.Duration } pidStats := make([]pidDelayStats, 0, len(pids)) + var deadPids []uint32 for _, pid := range pids { stats, err := TaskstatsTGID(pid) if err != nil { + // Only a pid that is gone is dead: a failed netlink call for a + // live process would otherwise close its uprobes. + if _, statErr := os.Stat(proc.Path(pid)); os.IsNotExist(statErr) { + deadPids = append(deadPids, pid) + } continue } pidStats = append(pidStats, pidDelayStats{ @@ -1707,6 +1733,26 @@ func (c *Container) updateDelays() { c.delaysByPid[ps.pid] = d } c.lock.Unlock() + return deadPids +} + +// youngerThan reports whether the container has existed for less than d, +// counted from its earliest process start or from when the agent found it, +// whichever is earlier, until now or until it became a zombie. Counting from +// discovery keeps a restarted unit, whose startedAt moves to the new +// process, from disappearing for d after every restart. +func (c *Container) youngerThan(d time.Duration) bool { + c.lock.RLock() + defer c.lock.RUnlock() + since := c.startedAt + if since.IsZero() || c.createdAt.Before(since) { + since = c.createdAt + } + end := time.Now() + if !c.zombieAt.IsZero() && c.zombieAt.Before(end) { + end = c.zombieAt + } + return end.Sub(since) < d } func (c *Container) updateNodejsStats(s NodejsStatsUpdate) { diff --git a/flags/flags.go b/flags/flags.go index 9d9fdcf1..5ab50e54 100644 --- a/flags/flags.go +++ b/flags/flags.go @@ -30,6 +30,7 @@ var ( ContainerAllowlist = kingpin.Flag("container-allowlist", "List of allowed containers (regex patterns)").Envar("CONTAINER_ALLOWLIST").Strings() ContainerDenylist = kingpin.Flag("container-denylist", "List of denied containers (regex patterns)").Envar("CONTAINER_DENYLIST").Strings() + MinContainerAge = kingpin.Flag("min-container-age", "Don't report metrics for containers younger than this. Suppresses short-lived job/cronjob pods that produce high-cardinality series. 0 disables.").Default("30s").Envar("MIN_CONTAINER_AGE").Duration() SkipSystemdSystemServices = kingpin.Flag("skip-systemd-system-services", "Skip well-known systemd system services (apt, motd, udev, etc.)").Default("true").Envar("SKIP_SYSTEMD_SYSTEM_SERVICES").Bool() From 9b45e846a62429d53307a1eab46b7a1cab6b3bf1 Mon Sep 17 00:00:00 2001 From: mayankpande88 Date: Wed, 7 Oct 2026 18:19:43 +0530 Subject: [PATCH 3/4] fix: exit only the process that was found gone The sweep in Collect finds exited processes without holding c.lock. If the pid was reused and the new process registered before the exit is handled, onProcessExit would untrack the new process. updateDelays now returns the *Process it found gone, and onProcessExitIf handles the exit only while that same process is registered. --- containers/container.go | 32 ++++++++++++++++++++++++-------- 1 file changed, 24 insertions(+), 8 deletions(-) diff --git a/containers/container.go b/containers/container.go index 8edbd069..29b61ef3 100644 --- a/containers/container.go +++ b/containers/container.go @@ -328,8 +328,8 @@ func (c *Container) Collect(ch chan<- prometheus.Metric) { if taskstatsClient != nil { // Processes whose exit event was missed: without this they keep // the container out of zombie state and its age wrong. - for _, pid := range c.updateDelays() { - c.onProcessExit(pid, false) + for pid, p := range c.updateDelays() { + c.onProcessExitIf(pid, p) } } @@ -658,6 +658,17 @@ func (c *Container) onProcessExit(pid uint32, oomKill bool) { c.onProcessExitLocked(pid, oomKill) } +// onProcessExitIf handles the exit of pid only if it is still registered as +// p: between finding p gone and getting here, the pid may have been reused +// by a new process that must not be untracked. +func (c *Container) onProcessExitIf(pid uint32, p *Process) { + c.lock.Lock() + defer c.lock.Unlock() + if c.processes[pid] == p { + c.onProcessExitLocked(pid, false) + } +} + func (c *Container) onFileOpen(pid uint32, fd uint64, mnt uint64, log bool) { if mnt > 0 && !log { c.lock.Lock() @@ -1687,13 +1698,15 @@ func (c *Container) onRetransmission(src netaddr.IPPort, dst netaddr.IPPort) boo } // updateDelays refreshes the per-process CPU and disk delay counters and -// returns the pids that no longer exist. -func (c *Container) updateDelays() []uint32 { +// returns the processes whose pid no longer exists. +func (c *Container) updateDelays() map[uint32]*Process { // Get a snapshot of PIDs under read lock to avoid concurrent map access c.lock.RLock() pids := make([]uint32, 0, len(c.processes)) - for pid := range c.processes { + procs := make(map[uint32]*Process, len(c.processes)) + for pid, p := range c.processes { pids = append(pids, pid) + procs[pid] = p } c.lock.RUnlock() @@ -1704,14 +1717,17 @@ func (c *Container) updateDelays() []uint32 { diskDelay time.Duration } pidStats := make([]pidDelayStats, 0, len(pids)) - var deadPids []uint32 + var dead map[uint32]*Process for _, pid := range pids { stats, err := TaskstatsTGID(pid) if err != nil { // Only a pid that is gone is dead: a failed netlink call for a // live process would otherwise close its uprobes. if _, statErr := os.Stat(proc.Path(pid)); os.IsNotExist(statErr) { - deadPids = append(deadPids, pid) + if dead == nil { + dead = map[uint32]*Process{} + } + dead[pid] = procs[pid] } continue } @@ -1733,7 +1749,7 @@ func (c *Container) updateDelays() []uint32 { c.delaysByPid[ps.pid] = d } c.lock.Unlock() - return deadPids + return dead } // youngerThan reports whether the container has existed for less than d, From d43f186b594b8c202c77d4f0dab8f92f90a8163e Mon Sep 17 00:00:00 2001 From: mayankpande88 Date: Thu, 8 Oct 2026 09:25:33 +0530 Subject: [PATCH 4/4] default --min-container-age to 0 With 30s, a pod in CrashLoopBackOff at the 5-minute maximum backoff stays a zombie past gcInterval, is removed between restarts and comes back as a new container each time. If it dies within 30s it never reports: no OOM kills, memory or CPU for the pods most worth looking at. Jobs that finish in under 30s never appear either. Keep the flag, and enable it per install where short-lived series are a measured problem. --- flags/flags.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/flags/flags.go b/flags/flags.go index 5ab50e54..10fca964 100644 --- a/flags/flags.go +++ b/flags/flags.go @@ -30,7 +30,7 @@ var ( ContainerAllowlist = kingpin.Flag("container-allowlist", "List of allowed containers (regex patterns)").Envar("CONTAINER_ALLOWLIST").Strings() ContainerDenylist = kingpin.Flag("container-denylist", "List of denied containers (regex patterns)").Envar("CONTAINER_DENYLIST").Strings() - MinContainerAge = kingpin.Flag("min-container-age", "Don't report metrics for containers younger than this. Suppresses short-lived job/cronjob pods that produce high-cardinality series. 0 disables.").Default("30s").Envar("MIN_CONTAINER_AGE").Duration() + MinContainerAge = kingpin.Flag("min-container-age", "Don't report metrics for containers younger than this. Suppresses short-lived job/cronjob pods that produce high-cardinality series. 0 disables.").Default("0s").Envar("MIN_CONTAINER_AGE").Duration() SkipSystemdSystemServices = kingpin.Flag("skip-systemd-system-services", "Skip well-known systemd system services (apt, motd, udev, etc.)").Default("true").Envar("SKIP_SYSTEMD_SYSTEM_SERVICES").Bool()