diff --git a/containers/container.go b/containers/container.go index 29d2f3fb..29b61ef3 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, p := range c.updateDelays() { + c.onProcessExitIf(pid, p) + } + } + + 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,23 @@ 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) +} + +// 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() @@ -1668,12 +1697,16 @@ 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 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() @@ -1684,9 +1717,18 @@ func (c *Container) updateDelays() { diskDelay time.Duration } pidStats := make([]pidDelayStats, 0, len(pids)) + 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) { + if dead == nil { + dead = map[uint32]*Process{} + } + dead[pid] = procs[pid] + } continue } pidStats = append(pidStats, pidDelayStats{ @@ -1707,6 +1749,26 @@ func (c *Container) updateDelays() { c.delaysByPid[ps.pid] = d } c.lock.Unlock() + return dead +} + +// 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/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..10fca964 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("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() @@ -51,6 +52,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()