Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
74 changes: 68 additions & 6 deletions containers/container.go
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down Expand Up @@ -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{},
Expand Down Expand Up @@ -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)
}
}
Comment thread
mayankpande88 marked this conversation as resolved.

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)
}
Expand All @@ -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))
}
Expand Down Expand Up @@ -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)
}
Expand Down Expand Up @@ -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)
}
Comment thread
mayankpande88 marked this conversation as resolved.

// 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()
Expand Down Expand Up @@ -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()

Expand All @@ -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{
Expand All @@ -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) {
Expand Down
42 changes: 38 additions & 4 deletions containers/l7.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
)
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand All @@ -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)
Expand All @@ -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.
Expand Down
2 changes: 2 additions & 0 deletions flags/flags.go
Original file line number Diff line number Diff line change
Expand Up @@ -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()

Expand All @@ -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()

Expand Down
Loading