Skip to content

Instantly share code, notes, and snippets.

@mohashari
Created July 26, 2026 01:01
Show Gist options
  • Select an option

  • Save mohashari/45107c26c40aacc6daaf7fb81355f2ab to your computer and use it in GitHub Desktop.

Select an option

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
#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;
}
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
}
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)
}
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)
}
}
}
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)
}
#!/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
# /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