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
102 changes: 62 additions & 40 deletions containers/container.go
Original file line number Diff line number Diff line change
Expand Up @@ -1161,7 +1161,7 @@ func (c *Container) onL7RequestWithResult(pid uint32, fd uint64, timestamp uint6
// HTTP/2 frames cannot be parsed out of order, which is what a
// retry would deliver them as.
if r.Protocol == l7.ProtocolHTTP2 {
countL7Drop("unknown_connection", r)
dropL7Event(c.id, unknownConnectionReason(socketInfo), pid, fd, r, socketInfo)
return nil, L7RequestProcessed
}
klog.V(3).Infof("L7_EVENT_CONN_NOT_FOUND: pid=%d fd=%d container=%s num_connections=%d",
Expand All @@ -1178,7 +1178,7 @@ func (c *Container) onL7RequestWithResult(pid uint32, fd uint64, timestamp uint6
if prev := conn; socketInfo != nil && socketInfo.Valid {
fresh, filtered := c.createConnectionFromSocketInfo(pid, fd, timestamp, socketInfo)
if fresh == nil && !filtered {
countL7Drop("unknown_connection", r)
dropL7Event(c.id, "unknown_connection", pid, fd, r, socketInfo)
return nil, L7RequestProcessed
}
prev.Closed = time.Now()
Expand All @@ -1197,7 +1197,7 @@ func (c *Container) onL7RequestWithResult(pid uint32, fd uint64, timestamp uint6
if timestamp != 0 && conn.Timestamp != timestamp {
klog.V(5).Infof("L7_EVENT_TIMESTAMP_MISMATCH: pid=%d fd=%d event_ts=%d conn_ts=%d protocol=%d",
pid, fd, timestamp, conn.Timestamp, r.Protocol)
countL7Drop("stale_connection", r)
dropL7Event(c.id, "stale_connection", pid, fd, r, socketInfo)
return nil, L7RequestTimestampMismatch
}

Expand Down Expand Up @@ -2158,18 +2158,16 @@ const (
openSslRecheckInterval = 10 * time.Second
)

// tlsExeRecheckInterval bounds how often the periodic sweep re-reads a
// process's executable to notice an exec (see recheckTlsAfterExec). New
// sockets always re-read it.
const tlsExeRecheckInterval = time.Second

// processForSocket returns the process that opened a socket, registering it
// if its start event has not been handled yet. Socket events often come
// first: they arrive on separate buffers, and connect events are read more
// often. Waiting for the start event dropped the connection (onConnectionOpen
// ignores unknown pids) and left a short-lived process's first TLS calls,
// often all of them, unprobed.
func (c *Container) processForSocket(pid uint32) *Process {
// tlsExeRecheckInterval bounds how often a process's executable is re-read to
// catch an exec whose event was lost (see recheckTlsAfterExec).
const tlsExeRecheckInterval = 10 * time.Second

// ensureProcess returns the process, registering it if its start event has
// not been handled yet. Socket and exec events often come first: they arrive
// on other buffers or CPUs. Waiting for the start event dropped the
// 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 {
return p
}
Expand All @@ -2182,12 +2180,12 @@ func (c *Container) attachTlsUprobes(tracer *ebpftracer.Tracer, pid uint32, newS
p := c.processes[pid]
if p == nil {
if newSocket {
// processForSocket could not register it: the process is gone.
// ensureProcess could not register it: the process is gone.
countTLSAttach("-", "not_registered")
}
return
}
c.recheckTlsAfterExec(p, newSocket)
c.recheckTlsAfterExec(p)
if !p.openSslUprobesChecked && time.Since(p.openSslLastCheck) >= openSslRecheckInterval {
p.openSslLastCheck = time.Now()
p.openSslChecks++
Expand All @@ -2197,41 +2195,65 @@ func (c *Container) attachTlsUprobes(tracer *ebpftracer.Tracer, pid uint32, newS
p.tlsAttached = p.tlsAttached || len(links) > 0
p.openSslUprobesChecked = len(links) > 0 || p.openSslChecks >= openSslMaxChecks
}
if !p.goTlsUprobesChecked {
p.tlsExe, _ = exeIdentityOf(pid)
p.tlsExeName, _ = os.Readlink(proc.Path(pid, "exe"))
p.tlsExeCheckedAt = time.Now()
links, isGolangApp, result := tracer.AttachGoTlsUprobes(pid)
countTLSAttach("go", result)
p.setGolang(isGolangApp)
p.addUprobes(links)
p.tlsAttached = p.tlsAttached || len(links) > 0
p.goTlsUprobesChecked = true
c.attachGoTls(tracer, p)
}

// attachGoTls probes the process's Go TLS functions, once per program image.
func (c *Container) attachGoTls(tracer *ebpftracer.Tracer, p *Process) {
if p.goTlsUprobesChecked {
return
}
p.tlsExe, _ = exeIdentityOf(p.Pid)
p.tlsExeName, _ = os.Readlink(proc.Path(p.Pid, "exe"))
p.tlsExeCheckedAt = time.Now()
links, isGolangApp, result := tracer.AttachGoTlsUprobes(p.Pid)
countTLSAttach("go", result)
p.setGolang(isGolangApp)
p.addUprobes(links)
p.tlsAttached = p.tlsAttached || len(links) > 0
p.goTlsUprobesChecked = true
}

// recheckTlsAfterExec lets a process be probed again after it execs a
// different binary. Attach runs on the first connection, so a wrapper that
// connects before exec'ing the real program (a secrets launcher, for example)
// would otherwise leave that program unprobed for its whole life, or probed
// at addresses of the wrapper's binary.
//
// Every new socket checks: the program a wrapper execs usually connects
// within a second, and a throttle shared with the periodic sweep skipped that
// first connection. The cost is one stat of /proc/<pid>/exe per connection.
func (c *Container) recheckTlsAfterExec(p *Process, newSocket bool) {
if !p.goTlsUprobesChecked {
// onProcessExec handles a process replacing its image. Its Go TLS probes are
// attached now rather than at its first connection, which a short-lived
// program often makes before its connect event is handled. OpenSSL is left
// to the first connection: the dynamic loader maps the libraries after the
// exec, so they are not there yet.
func (c *Container) onProcessExec(tracer *ebpftracer.Tracer, pid uint32) {
p := c.ensureProcess(pid)
if p == nil {
return
}
if !newSocket && time.Since(p.tlsExeCheckedAt) < tlsExeRecheckInterval {
if p.goTlsUprobesChecked {
exe, err := exeIdentityOf(pid)
if err != nil || exe == p.tlsExe {
// Gone, or a connection was handled first and already probed
// this image.
return
}
c.resetTlsForNewImage(p)
}
c.attachGoTls(tracer, p)
}

// recheckTlsAfterExec catches an exec whose event was lost (onProcessExec
// handles the rest): a process that runs a different binary than the one it
// was probed for is probed again.
func (c *Container) recheckTlsAfterExec(p *Process) {
if !p.goTlsUprobesChecked || time.Since(p.tlsExeCheckedAt) < tlsExeRecheckInterval {
return
}
p.tlsExeCheckedAt = time.Now()
exe, err := exeIdentityOf(p.Pid)
if err != nil || exe == p.tlsExe {
return
}
// The old image's probes point into a binary the process no longer runs.
c.resetTlsForNewImage(p)
}

// resetTlsForNewImage forgets what was probed for the process's previous
// image: those probes point into a binary the process no longer runs.
func (c *Container) resetTlsForNewImage(p *Process) {
p.dropUprobes()
if p.golang() {
c.registry.tracer.ReleaseGoTLSOffsets(p.Pid)
Expand Down
43 changes: 40 additions & 3 deletions containers/l7_self_metrics.go
Original file line number Diff line number Diff line change
@@ -1,10 +1,14 @@
package containers

import (
"net"
"strconv"

"github.com/coroot/coroot-node-agent/ebpftracer"
"github.com/coroot/coroot-node-agent/ebpftracer/l7"
lru "github.com/hashicorp/golang-lru/v2"
"github.com/prometheus/client_golang/prometheus"
"k8s.io/klog/v2"
)

var (
Expand Down Expand Up @@ -199,7 +203,8 @@ var TLSAttachTotal = prometheus.NewCounterVec(

// L7EventsDroppedTotal counts L7 events discarded in the agent before any
// protocol parsing: the connection or process they belong to was never
// found, or the retry queue for such events was full. Events lost in the
// found, or the retry queue for such events was full. no_ip_socket marks
// events on sockets the agent does not track (Unix sockets), not a loss. Events lost in the
// kernel are counted separately (node_agent_l7_ringbuf_drops_total,
// node_agent_tls_plaintext_dropped_total).
var L7EventsDroppedTotal = prometheus.NewCounterVec(
Expand All @@ -210,8 +215,40 @@ var L7EventsDroppedTotal = prometheus.NewCounterVec(
[]string{"reason", "protocol", "tls"},
)

func countL7Drop(reason string, r *l7.RequestData) {
L7EventsDroppedTotal.WithLabelValues(reason, protocolLabel(r.Protocol), strconv.FormatBool(r.TLS)).Inc()
type l7DropLogKey struct {
container ContainerID
reason string
protocol l7.Protocol
}
Comment thread
mayankpande88 marked this conversation as resolved.

// l7DropLogged keeps one log line per container, reason and protocol: the
// counter says how many events are dropped, the line says where.
var l7DropLogged, _ = lru.New[l7DropLogKey, struct{}](4096)

// unknownConnectionReason names why an event's connection could not be found:
// no_ip_socket when the event carries no IP socket tuple, as for gRPC over a
// Unix socket, which the agent does not track; unknown_connection otherwise.
func unknownConnectionReason(si *ebpftracer.SocketInfo) string {
if si == nil || !si.Valid {
return "no_ip_socket"
}
return "unknown_connection"
}

// dropL7Event counts an L7 event dropped before parsing and, the first time
// for its container, reason and protocol, logs the process and destination.
// container is empty when the event's process belongs to no known container.
func dropL7Event(container ContainerID, reason string, pid uint32, fd uint64, req *l7.RequestData, si *ebpftracer.SocketInfo) {
L7EventsDroppedTotal.WithLabelValues(reason, protocolLabel(req.Protocol), strconv.FormatBool(req.TLS)).Inc()
if ok, _ := l7DropLogged.ContainsOrAdd(l7DropLogKey{container: container, reason: reason, protocol: req.Protocol}, struct{}{}); ok {
return
}
Comment thread
mayankpande88 marked this conversation as resolved.
dst := "unknown"
if si != nil && si.Valid {
dst = net.JoinHostPort(si.DstIP, strconv.Itoa(int(si.DstPort)))
}
klog.Infof("L7 events dropped before parsing: reason=%s protocol=%s tls=%t container=%s pid=%d fd=%d dst=%s (logged once per container, reason and protocol)",
reason, protocolLabel(req.Protocol), req.TLS, container, pid, fd, dst)
}

// RegisterL7SelfMetrics registers the agent's L7 self-observability counters
Expand Down
29 changes: 20 additions & 9 deletions containers/registry.go
Original file line number Diff line number Diff line change
Expand Up @@ -336,6 +336,10 @@ func (r *Registry) handleEvents(ch <-chan ebpftracer.Event) {
r.processInfoCh <- ProcessInfo{Pid: p.Pid, ContainerId: c.id, StartedAt: p.StartedAt, Flags: p.Flags}
}
}
case ebpftracer.EventTypeProcessExec:
if c := r.getOrCreateContainer(e.Pid); c != nil {
c.onProcessExec(r.tracer, e.Pid)
}
case ebpftracer.EventTypeProcessExit:
r.containerLock.RLock()
c := r.containersByPid[e.Pid]
Expand All @@ -354,7 +358,7 @@ func (r *Registry) handleEvents(ch <-chan ebpftracer.Event) {

case ebpftracer.EventTypeListenOpen:
if c := r.getOrCreateContainer(e.Pid); c != nil {
c.processForSocket(e.Pid)
c.ensureProcess(e.Pid)
c.onListenOpen(e.Pid, e.SrcAddr, false)
c.attachTlsUprobes(r.tracer, e.Pid, true)
}
Expand All @@ -368,13 +372,13 @@ func (r *Registry) handleEvents(ch <-chan ebpftracer.Event) {

case ebpftracer.EventTypeConnectionOpen:
if c := r.getOrCreateContainer(e.Pid); c != nil {
c.processForSocket(e.Pid)
c.ensureProcess(e.Pid)
c.onConnectionOpen(e.Pid, e.Fd, e.SrcAddr, e.DstAddr, e.ActualDstAddr, e.Timestamp, false, e.Duration)
c.attachTlsUprobes(r.tracer, e.Pid, true)
}
case ebpftracer.EventTypeConnectionError:
if c := r.getOrCreateContainer(e.Pid); c != nil {
c.processForSocket(e.Pid)
c.ensureProcess(e.Pid)
c.onConnectionOpen(e.Pid, e.Fd, e.SrcAddr, e.DstAddr, e.ActualDstAddr, 0, true, e.Duration)
}
case ebpftracer.EventTypeConnectionClose:
Expand Down Expand Up @@ -417,12 +421,19 @@ func (r *Registry) handleEvents(ch <-chan ebpftracer.Event) {
}
}

func containerIdOf(c *Container) ContainerID {
if c == nil {
return ""
}
return c.id
}

// processL7Event handles an L7 event, queueing it for retry if the connection isn't found yet
func (r *Registry) processL7Event(e ebpftracer.Event) {
defer func() {
if p := recover(); p != nil {
klog.Errorf("recovered from panic in L7 event handler: pid=%d fd=%d protocol=%d: %v", e.Pid, e.Fd, e.L7Request.Protocol, p)
countL7Drop("panic", e.L7Request)
dropL7Event(containerIdOf(r.containersByPid[e.Pid]), "panic", e.Pid, e.Fd, e.L7Request, e.SocketInfo)
}
}()
if c := r.containersByPid[e.Pid]; c != nil {
Expand Down Expand Up @@ -466,7 +477,7 @@ func (r *Registry) queueL7EventForRetry(e ebpftracer.Event) {
const maxPendingEvents = 500
if len(r.pendingL7Events) >= maxPendingEvents {
klog.V(3).Infof("L7_EVENT_QUEUE_FULL: dropping event pid=%d fd=%d", e.Pid, e.Fd)
countL7Drop("retry_queue_full", e.L7Request)
dropL7Event(containerIdOf(r.containersByPid[e.Pid]), "retry_queue_full", e.Pid, e.Fd, e.L7Request, e.SocketInfo)
return
}

Expand Down Expand Up @@ -501,7 +512,7 @@ func (r *Registry) processPendingL7Events() {
// Expire old events
if now.Sub(p.addedAt) > maxAge {
klog.V(3).Infof("L7_EVENT_EXPIRED: pid=%d fd=%d age=%v", p.event.Pid, p.event.Fd, now.Sub(p.addedAt))
countL7Drop("retry_expired", p.event.L7Request)
dropL7Event(containerIdOf(r.containersByPid[p.event.Pid]), "retry_expired", p.event.Pid, p.event.Fd, p.event.L7Request, p.event.SocketInfo)
continue
}

Expand All @@ -512,14 +523,14 @@ func (r *Registry) processPendingL7Events() {
p.retryCount++
stillPending = append(stillPending, p)
} else {
countL7Drop("unknown_process", p.event.L7Request)
dropL7Event("", "unknown_process", p.event.Pid, p.event.Fd, p.event.L7Request, p.event.SocketInfo)
}
continue
}

ip2fqdn, result, panicked := r.safeOnL7Request(c, p.event)
if panicked {
countL7Drop("panic", p.event.L7Request)
dropL7Event(c.id, "panic", p.event.Pid, p.event.Fd, p.event.L7Request, p.event.SocketInfo)
continue
}
if result == L7RequestConnNotFound {
Expand All @@ -529,7 +540,7 @@ func (r *Registry) processPendingL7Events() {
stillPending = append(stillPending, p)
} else {
klog.V(3).Infof("L7_EVENT_MAX_RETRIES: pid=%d fd=%d", p.event.Pid, p.event.Fd)
countL7Drop("unknown_connection", p.event.L7Request)
dropL7Event(c.id, unknownConnectionReason(p.event.SocketInfo), p.event.Pid, p.event.Fd, p.event.L7Request, p.event.SocketInfo)
}
continue
}
Expand Down
20 changes: 10 additions & 10 deletions ebpftracer/ebpf.go

Large diffs are not rendered by default.

1 change: 1 addition & 0 deletions ebpftracer/ebpf/ebpf.c
Original file line number Diff line number Diff line change
Expand Up @@ -15,6 +15,7 @@
#define EVENT_TYPE_FILE_OPEN 8
#define EVENT_TYPE_TCP_RETRANSMIT 9
#define EVENT_TYPE_PYTHON_THREAD_LOCK 11
#define EVENT_TYPE_PROCESS_EXEC 13

#define EVENT_REASON_OOM_KILL 1

Expand Down
15 changes: 15 additions & 0 deletions ebpftracer/ebpf/proc.c
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,21 @@ struct trace_event_raw_sched_process_template__stub {
__u32 pid;
};

// sched_process_exec reports that a process replaced its image, so userspace
// can probe the new program at once: waiting for its first connection let a
// short-lived program make its first requests unprobed, and left the program
// a wrapper execs probed at the wrapper's addresses.
SEC("tracepoint/sched/sched_process_exec")
int sched_process_exec(void *ctx)
{
struct proc_event e = {
.type = EVENT_TYPE_PROCESS_EXEC,
.pid = bpf_get_current_pid_tgid() >> 32,
};
bpf_perf_event_output(ctx, &proc_events, BPF_F_CURRENT_CPU, &e, sizeof(e));
return 0;
}

SEC("tracepoint/sched/sched_process_exit")
int sched_process_exit(struct trace_event_raw_sched_process_template__stub *args)
{
Expand Down
7 changes: 6 additions & 1 deletion ebpftracer/tracer.go
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@ const (
EventTypeL7Request EventType = 10
EventTypePythonThreadLock EventType = 11
EventTypeLLMData EventType = 12
EventTypeProcessExec EventType = 13

EventReasonNone EventReason = 0
EventReasonOOMKill EventReason = 1
Expand Down Expand Up @@ -609,7 +610,9 @@ func (t *Tracer) ebpf(ch chan<- Event) error {
}

perfMaps := []perfMap{
{name: "proc_events", typ: perfMapTypeProcEvents, perCPUBufferSizePages: 4},
// Read as often as connect events: an exec is acted on (TLS probes
// attached) before the new program makes its first connection.
{name: "proc_events", typ: perfMapTypeProcEvents, perCPUBufferSizePages: 4, readTimeout: 10 * time.Millisecond},
{name: "tcp_listen_events", typ: perfMapTypeTCPEvents, perCPUBufferSizePages: 4},
{name: "tcp_connect_events", typ: perfMapTypeTCPEvents, perCPUBufferSizePages: 8, readTimeout: 10 * time.Millisecond},
{name: "tcp_retransmit_events", typ: perfMapTypeTCPEvents, perCPUBufferSizePages: 4},
Expand Down Expand Up @@ -704,6 +707,8 @@ func (t EventType) String() string {
return "process-start"
case EventTypeProcessExit:
return "process-exit"
case EventTypeProcessExec:
return "process-exec"
case EventTypeConnectionOpen:
return "connection-open"
case EventTypeConnectionClose:
Expand Down
Loading