Created
July 26, 2026 01:01
-
-
Save mohashari/45107c26c40aacc6daaf7fb81355f2ab to your computer and use it in GitHub Desktop.
Debugging Socket Buffer Bloat and TCP Queueing Delay in High-Throughput Go Services with eBPF — code snippets
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| #include <vmlinux.h> | |
| #include <bpf/bpf_helpers.h> | |
| #include <bpf/bpf_tracing.h> | |
| #include <bpf/bpf_endian.h> | |
| char LICENSE[] SEC("license") = "GPL"; | |
| struct event { | |
| __u32 pid; | |
| __u64 duration_ns; | |
| __u32 saddr; | |
| __u32 daddr; | |
| __u16 sport; | |
| __u16 dport; | |
| } __attribute__((packed)); | |
| struct { | |
| __uint(type, BPF_MAP_TYPE_HASH); | |
| __uint(max_entries, 10240); | |
| __type(key, struct sock *); | |
| __type(value, __u64); | |
| } write_start SEC(".maps"); | |
| struct { | |
| __uint(type, BPF_MAP_TYPE_RINGBUF); | |
| __uint(max_entries, 256 * 1024); | |
| } rb SEC(".maps"); | |
| SEC("kprobe/tcp_sendmsg") | |
| int BPF_KPROBE(kprobe_tcp_sendmsg, struct sock *sk, struct msghdr *msg, size_t size) { | |
| __u64 ts = bpf_ktime_get_ns(); | |
| bpf_map_update_elem(&write_start, &sk, &ts, BPF_ANY); | |
| return 0; | |
| } | |
| SEC("kprobe/tcp_write_xmit") | |
| int BPF_KPROBE(kprobe_tcp_write_xmit, struct sock *sk, unsigned int mss_now, int nonblock, int push_one, gfp_t gfp) { | |
| __u64 *tsp = bpf_map_lookup_elem(&write_start, &sk); | |
| if (!tsp) { | |
| return 0; | |
| } | |
| __u64 duration = bpf_ktime_get_ns() - *tsp; | |
| bpf_map_delete_elem(&write_start, &sk); | |
| // Limit logging to queues taking longer than 1ms to prevent logging noise | |
| if (duration < 1000000) { | |
| return 0; | |
| } | |
| struct event *ev = bpf_ringbuf_reserve(&rb, sizeof(*ev), 0); | |
| if (!ev) { | |
| return 0; | |
| } | |
| ev->pid = bpf_get_current_pid_tgid() >> 32; | |
| ev->duration_ns = duration; | |
| // Extract socket metadata using BPF CO-RE (Compile Once, Run Everywhere) | |
| ev->saddr = BPF_CORE_READ(sk, __sk_common.skc_rcv_saddr); | |
| ev->daddr = BPF_CORE_READ(sk, __sk_common.skc_daddr); | |
| ev->sport = BPF_CORE_READ(sk, __sk_common.skc_num); | |
| ev->dport = bpf_ntohs(BPF_CORE_READ(sk, __sk_common.skc_dport)); | |
| bpf_ringbuf_submit(ev, 0); | |
| return 0; | |
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| package main | |
| import ( | |
| "bytes" | |
| "encoding/binary" | |
| "errors" | |
| "log" | |
| "net" | |
| "os" | |
| "os/signal" | |
| "syscall" | |
| "time" | |
| "github.com/cilium/ebpf/link" | |
| "github.com/cilium/ebpf/ringbuf" | |
| "github.com/cilium/ebpf/rlimit" | |
| ) | |
| //go:generate go run github.com/cilium/ebpf/cmd/bpf2go -target bpf main main.c -- -I./headers | |
| type TCPQueueEvent struct { | |
| PID uint32 | |
| DurationNS uint64 | |
| SAddr uint32 | |
| DAddr uint32 | |
| SPort uint16 | |
| DPort uint16 | |
| } | |
| func main() { | |
| // Remove memory limits for locking eBPF maps in kernel memory space | |
| if err := rlimit.RemoveMemlock(); err != nil { | |
| log.Fatalf("failed to remove memlock limit: %v", err) | |
| } | |
| // Load the compiled BPF structures generated by bpf2go | |
| var objs mainObjects | |
| if err := loadMainObjects(&objs, nil); err != nil { | |
| log.Fatalf("loading eBPF objects failed: %v", err) | |
| } | |
| defer objs.Close() | |
| // Attach kprobe to tcp_sendmsg | |
| kpSend, err := link.Kprobe("tcp_sendmsg", objs.KprobeTcpSendmsg, nil) | |
| if err != nil { | |
| log.Fatalf("attaching tcp_sendmsg kprobe failed: %v", err) | |
| } | |
| defer kpSend.Close() | |
| // Attach kprobe to tcp_write_xmit | |
| kpXmit, err := link.Kprobe("tcp_write_xmit", objs.KprobeTcpWriteXmit, nil) | |
| if err != nil { | |
| log.Fatalf("attaching tcp_write_xmit kprobe failed: %v", err) | |
| } | |
| defer kpXmit.Close() | |
| log.Println("eBPF latency tracker initialized. Polling ring buffer...") | |
| // Listen to the eBPF ring buffer events | |
| rd, err := ringbuf.NewReader(objs.Rb) | |
| if err != nil { | |
| log.Fatalf("failed to create ringbuf reader: %v", err) | |
| } | |
| defer rd.Close() | |
| // Graceful shutdown handling | |
| stopChan := make(chan os.Signal, 1) | |
| signal.Notify(stopChan, os.Interrupt, syscall.SIGTERM) | |
| go func() { | |
| <-stopChan | |
| _ = rd.Close() | |
| }() | |
| var event TCPQueueEvent | |
| for { | |
| record, err := rd.Read() | |
| if err != nil { | |
| if errors.Is(err, ringbuf.ErrClosed) { | |
| log.Println("Ring buffer reader closed, exiting daemon.") | |
| return | |
| } | |
| log.Printf("error reading from ring buffer: %v", err) | |
| continue | |
| } | |
| // Read the packed binary event directly into the Go struct | |
| err = binary.Read(bytes.NewBuffer(record.RawSample), binary.LittleEndian, &event) | |
| if err != nil { | |
| log.Printf("failed to parse event structure: %v", err) | |
| continue | |
| } | |
| srcIP := parseIP(event.SAddr) | |
| dstIP := parseIP(event.DAddr) | |
| log.Printf("[TCP Bloat Alert] PID: %d | Delay: %v | Source: %s:%d -> Destination: %s:%d", | |
| event.PID, | |
| time.Duration(event.DurationNS), | |
| srcIP.String(), event.SPort, | |
| dstIP.String(), event.DPort, | |
| ) | |
| } | |
| } | |
| func parseIP(val uint32) net.IP { | |
| ip := make(net.IP, 4) | |
| binary.LittleEndian.PutUint32(ip, val) | |
| return ip | |
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| package main | |
| import ( | |
| "context" | |
| "log" | |
| "net" | |
| "syscall" | |
| "golang.org/x/sys/unix" | |
| ) | |
| // NewWatermarkedListener returns a listener that configures TCP_NOTSENT_LOWAT on connections | |
| func NewWatermarkedListener(network, address string, lowatBytes int) (net.Listener, error) { | |
| config := net.ListenConfig{ | |
| Control: func(network, address string, c syscall.RawConn) error { | |
| var socketErr error | |
| controlErr := c.Control(func(fd uintptr) { | |
| // Apply TCP_NOTSENT_LOWAT to limit the unsent byte queue size in kernel space | |
| socketErr = unix.SetsockoptInt(int(fd), unix.IPPROTO_TCP, unix.TCP_NOTSENT_LOWAT, lowatBytes) | |
| if socketErr != nil { | |
| log.Printf("warning: failed to set TCP_NOTSENT_LOWAT: %v", socketErr) | |
| } | |
| // Optimize socket buffer size to prevent memory waste under high connection concurrency | |
| _ = unix.SetsockoptInt(int(fd), unix.SOL_SOCKET, unix.SO_SNDBUF, 65536) | |
| }) | |
| if controlErr != nil { | |
| return controlErr | |
| } | |
| return socketErr | |
| }, | |
| } | |
| return config.Listen(context.Background(), network, address) | |
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| package main | |
| import ( | |
| "io" | |
| "net" | |
| "sync" | |
| "testing" | |
| "time" | |
| ) | |
| func BenchmarkSlowReaderBackpressure(b *testing.B) { | |
| serverAddr := "127.0.0.1:48080" | |
| // Start server listening with a tight 16KB TCP_NOTSENT_LOWAT limit | |
| ln, err := NewWatermarkedListener("tcp", serverAddr, 16384) | |
| if err != nil { | |
| b.Fatalf("failed to start listener: %v", err) | |
| } | |
| defer ln.Close() | |
| var wg sync.WaitGroup | |
| wg.Add(1) | |
| // Server routine reads incoming data slowly to create backpressure | |
| go func() { | |
| defer wg.Done() | |
| conn, err := ln.Accept() | |
| if err != nil { | |
| return | |
| } | |
| defer conn.Close() | |
| buf := make([]byte, 4096) | |
| for { | |
| _, err := conn.Read(buf) | |
| if err != nil { | |
| if err == io.EOF { | |
| break | |
| } | |
| return | |
| } | |
| // Artificially restrict read throughput to force buffering on the write side | |
| time.Sleep(10 * time.Millisecond) | |
| } | |
| }() | |
| clientConn, err := net.Dial("tcp", serverAddr) | |
| if err != nil { | |
| b.Fatalf("failed to connect: %v", err) | |
| } | |
| defer clientConn.Close() | |
| payload := make([]byte, 1024*1024) // 1MB payload writes | |
| b.ResetTimer() | |
| for i := 0; i < b.N; i++ { | |
| start := time.Now() | |
| _, err := clientConn.Write(payload) | |
| if err != nil { | |
| b.Fatalf("write failed: %v", err) | |
| } | |
| writeDuration := time.Since(start) | |
| // Without TCP_NOTSENT_LOWAT, a 1MB write completes instantly (<1ms) because | |
| // the kernel buffers the entire payload. | |
| // With a 16KB TCP_NOTSENT_LOWAT, the write blocks in Go until the reader drains the buffer. | |
| if writeDuration > 5*time.Millisecond { | |
| b.Logf("Iteration %d: Backpressure applied successfully. Write blocked for %v", i, writeDuration) | |
| } | |
| } | |
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| package main | |
| import ( | |
| "log" | |
| "net/http" | |
| "strconv" | |
| "time" | |
| "github.com/prometheus/client_golang/prometheus" | |
| "github.com/prometheus/client_golang/prometheus/promhttp" | |
| ) | |
| var ( | |
| tcpQueueDuration = prometheus.NewHistogramVec( | |
| prometheus.HistogramOpts{ | |
| Name: "tcp_kernel_queue_latency_seconds", | |
| Help: "TCP socket write queue latency in seconds, measured via eBPF probes.", | |
| Buckets: []float64{0.001, 0.005, 0.010, 0.025, 0.050, 0.100, 0.250, 0.500, 1.0, 2.5}, | |
| }, | |
| []string{"src_port", "dst_ip"}, | |
| ) | |
| ) | |
| func init() { | |
| prometheus.MustRegister(tcpQueueDuration) | |
| } | |
| func startMetricsServer(addr string) { | |
| http.Handle("/metrics", promhttp.Handler()) | |
| log.Printf("Starting prometheus server on %s", addr) | |
| if err := http.ListenAndServe(addr, nil); err != nil { | |
| log.Fatalf("metrics server failed to run: %v", err) | |
| } | |
| } | |
| func trackEvent(ev TCPQueueEvent, dstIP string) { | |
| durationSec := float64(ev.DurationNS) / float64(time.Second) | |
| tcpQueueDuration.With(prometheus.Labels{ | |
| "src_port": strconv.Itoa(int(ev.SPort)), | |
| "dst_ip": dstIP, | |
| }).Observe(durationSec) | |
| } |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| #!/usr/bin/env bash | |
| # build-ebpf.sh | |
| set -euo pipefail | |
| # 1. Install modern Clang, LLVM, and kernel headers | |
| sudo apt-get update && sudo apt-get install -y clang llvm libelf-dev libbpf-dev | |
| # 2. Extract vmlinux.h from the running kernel for CO-RE compatibility | |
| bpftool btf dump file /sys/kernel/btf/vmlinux format c > headers/vmlinux.h | |
| # 3. Trigger go generate to invoke bpf2go wrapper compiler | |
| go generate ./... | |
| # 4. Compile the final binary incorporating the eBPF asset | |
| go build -o tcp-latency-daemon main.go main_bpfel.go |
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
| # /etc/sysctl.d/99-tcp-tuning.conf | |
| # Set default queueing discipline to Fair Queueing (FQ) | |
| net.core.default_qdisc = fq | |
| # Enable BBR congestion control algorithm | |
| net.ipv4.tcp_congestion_control = bbr | |
| # Limit system-wide TCP write (send) buffers (Min, Default, Max in bytes) | |
| net.ipv4.tcp_wmem = 4096 16384 4194304 | |
| # Limit system-wide TCP read (receive) buffers (Min, Default, Max in bytes) | |
| net.ipv4.tcp_rmem = 4096 87380 6291456 | |
| # Enable dynamic TCP buffer autotuning | |
| net.ipv4.tcp_moderate_rcvbuf = 1 |
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment