Skip to content
Merged
3 changes: 3 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -39,3 +39,6 @@ coverage.html

# eBPF build outputs (regenerated via ebpftracer/make build)
ebpftracer/ebpf/*.o

# Local build helper (upstream)
/build.sh
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -9,6 +9,7 @@ gathers container and host metrics, logs, and L7 traffic using eBPF and
exposes them in Prometheus format.

Minimum Linux kernel: **5.8** (L7 events use a BPF ring buffer).
The kernel must also be built with `CONFIG_BPF_EVENTS=y` (kprobe and tracepoint BPF programs); some embedded and vendor kernels disable it.

> This project is a fork of
> [coroot/coroot-node-agent](https://github.com/coroot/coroot-node-agent)
Expand Down
4 changes: 2 additions & 2 deletions cgroup/cgroup.go
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ var (
crioIdRegexp = regexp.MustCompile(`crio-([a-z0-9]{64})`)
containerdIdRegexp = regexp.MustCompile(`cri-containerd[-:]([a-z0-9]{64})`)
lxcIdRegexp = regexp.MustCompile(`/lxc/([^/]+)`)
systemSliceIdRegexp = regexp.MustCompile(`(/(system|runtime|reserved|kube|azure)\.slice/([^/]+))`)
systemSliceIdRegexp = regexp.MustCompile(`(/(system|runtime|reserved|kube|azure|podruntime)\.slice/([^/]+))`)
talosIdRegexp = regexp.MustCompile(`/(system|podruntime)/([^/]+)`)
lxcPayloadRegexp = regexp.MustCompile(`/lxc\.payload\.([^/]+)`)
)
Expand Down Expand Up @@ -196,7 +196,7 @@ func containerByCgroup(cgroupPath string) (ContainerType, string, error) {
return ContainerTypeUnknown, "", fmt.Errorf("invalid talos runtime cgroup %s", cgroupPath)
}
return ContainerTypeTalosRuntime, path.Join("/talos/", matches[2]), nil
case prefix == "system.slice" || prefix == "runtime.slice" || prefix == "reserved.slice" || prefix == "kube.slice" || prefix == "azure.slice":
case prefix == "system.slice" || prefix == "runtime.slice" || prefix == "reserved.slice" || prefix == "kube.slice" || prefix == "azure.slice" || prefix == "podruntime.slice":
if strings.HasSuffix(cgroupPath, ".scope") {
return ContainerTypeStandaloneProcess, "", nil
}
Expand Down
10 changes: 10 additions & 0 deletions cgroup/cgroup_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -171,6 +171,16 @@ func TestContainerByCgroup(t *testing.T) {
as.Equal("/azure.slice/walinuxagent.service", id)
as.Nil(err)

typ, id, err = containerByCgroup("/podruntime.slice/containerd.service")
as.Equal(typ, ContainerTypeSystemdService)
as.Equal("/podruntime.slice/containerd.service", id)
as.Nil(err)

typ, id, err = containerByCgroup("/podruntime.slice/kubelet.service")
as.Equal(typ, ContainerTypeSystemdService)
as.Equal("/podruntime.slice/kubelet.service", id)
as.Nil(err)

typ, id, err = containerByCgroup("/system.slice/system-postgresql.slice/postgresql@9.4-main.service")
as.Equal(typ, ContainerTypeSystemdService)
as.Equal("/system.slice/system-postgresql.slice", id)
Expand Down
26 changes: 26 additions & 0 deletions common/api.go
Original file line number Diff line number Diff line change
@@ -1,7 +1,12 @@
package common

import (
"crypto/tls"
"crypto/x509"
"os"

"github.com/coroot/coroot-node-agent/flags"
"k8s.io/klog/v2"
)

func AuthHeaders() map[string]string {
Expand All @@ -11,3 +16,24 @@ func AuthHeaders() map[string]string {
}
return res
}

func TlsConfig() *tls.Config {
cfg := &tls.Config{InsecureSkipVerify: *flags.InsecureSkipVerify}
if *flags.CAFile != "" {
ca, err := os.ReadFile(*flags.CAFile)
if err != nil {
klog.Fatalln(err)
return cfg
Comment thread
mayankpande88 marked this conversation as resolved.
Dismissed
}
pool, err := x509.SystemCertPool()
if err != nil {
klog.Warningln("failed to load system cert pool, starting with empty pool:", err)
pool = x509.NewCertPool()
}
Comment thread
mayankpande88 marked this conversation as resolved.
if !pool.AppendCertsFromPEM(ca) {
klog.Fatalf("failed to parse CA from %s", *flags.CAFile)
}
cfg.RootCAs = pool
}
return cfg
}
67 changes: 47 additions & 20 deletions common/net.go
Original file line number Diff line number Diff line change
Expand Up @@ -36,24 +36,9 @@ func init() {
}
if r := flags.EphemeralPortRange; r != nil && *r != "" {
klog.Infoln("ephemeral-port-range:", *r)
parts := strings.Split(*r, "-")
if len(parts) != 2 {
klog.Exitf("invalid port range: %s", *r)
}
from, err := strconv.ParseUint(parts[0], 10, 16)
if err != nil {
klog.Exitf("invalid port range: %s", *r)
}
to, err := strconv.ParseUint(parts[1], 10, 16)
if err != nil {
klog.Exitf("invalid port range: %s", *r)
}
if from > to {
klog.Exitf("invalid port range: %s", *r)
}
PortFilter = &portFilter{
from: uint16(from),
to: uint16(to),
var err error
if PortFilter, err = newPortFilter(*r); err != nil {
klog.Exitln(err)
}
}
var err error
Expand Down Expand Up @@ -121,16 +106,58 @@ func (f connectionFilter) ShouldBeSkipped(dst, actualDst netaddr.IP) bool {
return true
}

type portFilter struct {
type portRange struct {
from uint16
to uint16
}

type portFilter struct {
ranges []portRange
}

func newPortFilter(s string) (*portFilter, error) {
f := &portFilter{}
for _, r := range strings.Fields(strings.ReplaceAll(s, ",", " ")) {
from, to, ok := strings.Cut(r, "-")
if !ok {
return nil, fmt.Errorf("invalid port range: %s", r)
}
f1, err := strconv.ParseUint(from, 10, 16)
if err != nil {
return nil, fmt.Errorf("invalid port range: %s", r)
}
t1, err := strconv.ParseUint(to, 10, 16)
if err != nil {
return nil, fmt.Errorf("invalid port range: %s", r)
}
if f1 > t1 {
return nil, fmt.Errorf("invalid port range: %s", r)
}
f.ranges = append(f.ranges, portRange{from: uint16(f1), to: uint16(t1)})
}
if len(f.ranges) == 0 {
return nil, fmt.Errorf("invalid port range: %s", s)
}
return f, nil
}

var wellKnownPorts = map[uint16]struct{}{
50051: {}, // gRPC
}

func (f *portFilter) ShouldBeSkipped(port uint16) bool {
if f == nil {
return false
}
return port >= f.from && port <= f.to
if _, ok := wellKnownPorts[port]; ok {
return false
}
for _, r := range f.ranges {
if port >= r.from && port <= r.to {
return true
}
}
return false
}

type HostPort struct {
Expand Down
36 changes: 36 additions & 0 deletions common/net_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -157,3 +157,39 @@ func BenchmarkNormalizeFQDN(b *testing.B) {
NormalizeFQDN("example.io.svc.default.cluster.local", "TypeA")
}
}

func TestPortFilter(t *testing.T) {
var nilFilter *portFilter
assert.False(t, nilFilter.ShouldBeSkipped(40000))

f, err := newPortFilter("32768-60999")
assert.NoError(t, err)
assert.False(t, f.ShouldBeSkipped(32767))
assert.True(t, f.ShouldBeSkipped(32768))
assert.True(t, f.ShouldBeSkipped(60999))
assert.False(t, f.ShouldBeSkipped(61000))
assert.False(t, f.ShouldBeSkipped(50051)) // well-known port within the range

for _, s := range []string{"1024-23768 30000-65535", "1024-23768,30000-65535", " 1024-23768 , 30000-65535 "} {
f, err = newPortFilter(s)
assert.NoError(t, err, s)
assert.False(t, f.ShouldBeSkipped(1023), s)
assert.True(t, f.ShouldBeSkipped(1024), s)
assert.True(t, f.ShouldBeSkipped(23768), s)
assert.False(t, f.ShouldBeSkipped(23769), s) // between the ranges
assert.False(t, f.ShouldBeSkipped(29999), s)
assert.True(t, f.ShouldBeSkipped(30000), s)
assert.True(t, f.ShouldBeSkipped(65535), s)
assert.False(t, f.ShouldBeSkipped(50051), s)
}

f, err = newPortFilter("8080-8080")
assert.NoError(t, err)
assert.True(t, f.ShouldBeSkipped(8080))
assert.False(t, f.ShouldBeSkipped(8081))

for _, s := range []string{"", " ", "32768", "-", "a-b", "32768-", "-60999", "60999-32768", "32768-65536", "-1-10"} {
_, err = newPortFilter(s)
assert.Error(t, err, s)
}
}
2 changes: 1 addition & 1 deletion containers/app.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,7 @@ import (
)

var (
phpCmd = regexp.MustCompile(`.*php\d*\.?\d*$`)
phpCmd = regexp.MustCompile(`.*php(-fpm)?\d*\.?\d*$`)
pythonCmd = regexp.MustCompile(`.*python\d*\.?\d*$`)
rubyCmd = regexp.MustCompile(`.*ruby\d*\.?\d*$`)
nodejsCmd = regexp.MustCompile(`.*node(js)?\d*\.?\d*$`)
Expand Down
42 changes: 38 additions & 4 deletions containers/container.go
Original file line number Diff line number Diff line change
Expand Up @@ -361,12 +361,17 @@ func (c *Container) Collect(ch chan<- prometheus.Metric) {

if disks, err := node.GetDisks(); err == nil {
ioStat := c.cgroup.IOStat()
seenVolumes := map[string]struct{}{}
for majorMinor, mounts := range c.getMounts() {
var device string
if dev := disks.GetParentBlockDevice(majorMinor); dev != nil {
device = dev.Name
}
for mountPoint, fsStat := range mounts {
if _, ok := seenVolumes[mountPoint+":"+device]; ok {
continue
}
seenVolumes[mountPoint+":"+device] = struct{}{}
dls := []string{mountPoint, device, c.metadata.volumes[mountPoint]}
ch <- c.gauge(metrics.DiskSize, float64(fsStat.CapacityBytes), dls...)
ch <- c.gauge(metrics.DiskUsed, float64(fsStat.UsedBytes), dls...)
Expand Down Expand Up @@ -1749,19 +1754,45 @@ func (c *Container) updatePythonStats(s PythonStatsUpdate) {
}

func (c *Container) getMounts() map[string]map[string]*proc.FSStat {
// c.processes and c.mounts are mutated by the handleEvents goroutine
// (onFileOpen), so both are read and updated under c.lock.
c.lock.RLock()
if len(c.mounts) == 0 {
c.lock.RUnlock()
return nil
}
// Copy pids under read lock — c.processes is mutated by handleEvents goroutine
c.lock.RLock()
pids := make([]uint32, 0, len(c.processes))
for pid := range c.processes {
pids = append(pids, pid)
}
c.lock.RUnlock()

res := map[string]map[string]*proc.FSStat{}
// Drop mounts that are gone from the container's mount namespace, so a
// remounted volume is not reported twice (stale + current entry).
var current map[string]proc.MountInfo
for _, pid := range pids {
if current = proc.GetMountInfo(pid); current != nil {
break
}
}
c.lock.Lock()
if len(current) > 0 {
for mntId := range c.mounts {
if mi, ok := current[mntId]; ok {
c.mounts[mntId] = mi
} else {
delete(c.mounts, mntId)
}
}
}
mounts := make([]proc.MountInfo, 0, len(c.mounts))
for _, mi := range c.mounts {
mounts = append(mounts, mi)
}
c.lock.Unlock()

res := map[string]map[string]*proc.FSStat{}
for _, mi := range mounts {
var stat *proc.FSStat
for _, pid := range pids {
s, err := proc.StatFS(proc.Path(pid, "root", mi.MountPoint))
Expand Down Expand Up @@ -2180,7 +2211,10 @@ const tlsExeRecheckInterval = 10 * time.Second
// connection (onConnectionOpen ignores unknown pids) and left a short-lived
// process's first TLS calls, often all of them, unprobed.
func (c *Container) ensureProcess(pid uint32) *Process {
if p := c.processes[pid]; p != nil {
c.lock.RLock()
p := c.processes[pid]
c.lock.RUnlock()
if p != nil {
return p
}
return c.onProcessStart(pid)
Expand Down
8 changes: 8 additions & 0 deletions containers/registry.go
Original file line number Diff line number Diff line change
Expand Up @@ -331,6 +331,8 @@ func (r *Registry) handleEvents(ch <-chan ebpftracer.Event) {
}
}
if c := r.getOrCreateContainer(e.Pid); c != nil {
// onProcessStart, not ensureProcess: it also replaces a
// registered process whose pid was reused.
p := c.onProcessStart(e.Pid)
if r.processInfoCh != nil && p != nil {
r.processInfoCh <- ProcessInfo{Pid: p.Pid, ContainerId: c.id, StartedAt: p.StartedAt, Flags: p.Flags}
Expand Down Expand Up @@ -613,6 +615,10 @@ func (r *Registry) getOrCreateContainer(pid uint32) *Container {
r.containerLock.Lock()
r.containersByPid[pid] = c
r.containerLock.Unlock()
// The pid may have been mapped without its start event (e.g. a unit's
// new main process after a restart); register it so the container
// does not go zombie while the process is alive.
c.ensureProcess(pid)
return c
}
r.containerLock.RUnlock()
Expand Down Expand Up @@ -688,6 +694,7 @@ func (r *Registry) getOrCreateContainer(pid uint32) *Container {
}
if c := r.containersByCgroupId[cg.Id]; c != nil {
r.containersByPid[pid] = c
c.ensureProcess(pid)
return c
}
if c := r.containersById[id]; c != nil {
Expand Down Expand Up @@ -718,6 +725,7 @@ func (r *Registry) getOrCreateContainer(pid uint32) *Container {
r.containersByPid[pid] = c
r.containersByCgroupId[cg.Id] = c
r.containersById[id] = c
c.ensureProcess(pid)
return c
}

Expand Down
6 changes: 6 additions & 0 deletions ebpftracer/tracer.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@ import (
"time"

"github.com/cilium/ebpf"
"github.com/cilium/ebpf/features"
"github.com/cilium/ebpf/link"
"github.com/cilium/ebpf/perf"
"github.com/cilium/ebpf/ringbuf"
Expand Down Expand Up @@ -589,6 +590,11 @@ func (t *Tracer) ebpf(ch chan<- Event) error {
}
t.programVariant = variant
_ = unix.Setrlimit(unix.RLIMIT_MEMLOCK, &unix.Rlimit{Cur: unix.RLIM_INFINITY, Max: unix.RLIM_INFINITY})
for _, pt := range []ebpf.ProgramType{ebpf.TracePoint, ebpf.Kprobe} {
if err := features.HaveProgramType(pt); errors.Is(err, ebpf.ErrNotSupported) {
return fmt.Errorf("kernel does not support BPF %s programs (CONFIG_BPF_EVENTS is not set?): %w", pt, ebpf.ErrNotSupported)
}
}
c, err := ebpf.NewCollectionWithOptions(collectionSpec, ebpf.CollectionOptions{
//Programs: ebpf.ProgramOptions{LogLevel: 2, LogSize: 20 * 1024 * 1024},
})
Expand Down
3 changes: 2 additions & 1 deletion flags/flags.go
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ var (
Envar("TRACK_PUBLIC_NETWORK").
Default("0.0.0.0/0").
Strings()
EphemeralPortRange = kingpin.Flag("ephemeral-port-range", "Destination and Listen TCP ports from this range will be skipped").Default("32768-60999").Envar("EPHEMERAL_PORT_RANGE").String()
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()
Expand All @@ -62,6 +62,7 @@ var (
LogsEndpoint = kingpin.Flag("logs-endpoint", "The URL of the endpoint to send logs to").Envar("LOGS_ENDPOINT").URL()
ProfilesEndpoint = kingpin.Flag("profiles-endpoint", "The URL of the endpoint to send profiles to").Envar("PROFILES_ENDPOINT").URL()
InsecureSkipVerify = kingpin.Flag("insecure-skip-verify", "whether to skip verifying the certificate or not").Envar("INSECURE_SKIP_VERIFY").Default("false").Bool()
CAFile = kingpin.Flag("ca-file", "Path to the custom CA certificate file").Envar("CA_FILE").String()

ScrapeInterval = kingpin.Flag("scrape-interval", "How often to gather metrics from the agent").Default("15s").Envar("SCRAPE_INTERVAL").Duration()
WalDir = kingpin.Flag("wal-dir", "Path to where the agent stores data (e.g. the metrics Write-Ahead Log)").Default("/tmp/nudgebee-node-agent").Envar("WAL_DIR").String()
Expand Down
3 changes: 1 addition & 2 deletions logs/otel.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,6 @@ package logs

import (
"context"
"crypto/tls"
"time"

otel "github.com/agoda-com/opentelemetry-logs-go"
Expand Down Expand Up @@ -41,7 +40,7 @@ func Init(machineId, hostname, version string) {
if endpointUrl.Scheme != "https" {
opts = append(opts, otlplogshttp.WithInsecure())
} else {
opts = append(opts, otlplogshttp.WithTLSClientConfig(&tls.Config{InsecureSkipVerify: *flags.InsecureSkipVerify}))
opts = append(opts, otlplogshttp.WithTLSClientConfig(common.TlsConfig()))
}
client := otlplogshttp.NewClient(opts...)
exporter, _ := otlplogs.NewExporter(context.Background(), otlplogs.WithClient(client))
Expand Down
Loading
Loading