#!/usr/bin/env python3 """Emit bridge ingress, egress, and drop telemetry events via eBPF.""" from __future__ import annotations import argparse import ctypes as ct import hashlib import ipaddress import json import os import signal import socket import sys from typing import Iterable try: from bcc import BPF # type: ignore except Exception as exc: # pragma: no cover - depends on host runtime print(f"Failed to import python3-bpfcc: {exc}", file=sys.stderr, flush=True) raise from src.utilities.packet_mark import packet_id_from_mark, verdict_from_mark IDENTITY_FIELDS = ( "src_mac", "dst_mac", "eth_type_raw", "vlan_id", "src_ip", "dst_ip", "protocol_raw", "src_port", "dst_port", "length", "ip_id", "icmp_type", "icmp_code", "arp_op", "tcp_seq", "tcp_ack", "tcp_flags", ) BPF_SOURCE = r""" #include #include #include #include #include #include #include #include #include #include #include #include #include #include #define EVENT_INGRESS 1 #define EVENT_EGRESS 2 #define EVENT_DROP 3 struct vlan_hdr_t { __be16 h_vlan_TCI; __be16 h_vlan_encapsulated_proto; }; struct arp_eth_ipv4_t { __u8 sha[6]; __u8 spa[4]; __u8 tha[6]; __u8 tpa[4]; }; struct event_t { __u64 ts_ns; __u32 skb_mark; __u32 length; __u32 reason; __u16 eth_type_raw; __u16 vlan_id; __u16 src_port; __u16 dst_port; __u16 ip_id; __u16 arp_op; __u32 protocol_raw; __u32 tcp_seq; __u32 tcp_ack; __u8 event_type; __u8 ip_version; __u8 icmp_type; __u8 icmp_code; __u8 tcp_flags; char ifname[IFNAMSIZ]; unsigned char src_mac[6]; unsigned char dst_mac[6]; unsigned char src_ip[16]; unsigned char dst_ip[16]; }; BPF_PERF_OUTPUT(events); static __always_inline int fill_ifname(struct sk_buff *skb, struct event_t *event) { if (!skb) { return 0; } struct net_device *dev = NULL; bpf_probe_read_kernel(&dev, sizeof(dev), &skb->dev); if (!dev) { return 0; } bpf_probe_read_kernel(event->ifname, sizeof(event->ifname), dev->name); return 1; } static __always_inline int parse_skb(struct sk_buff *skb, struct event_t *event) { unsigned char *head = NULL; __u16 mac_header = 0; __u16 network_header = 0; __u16 transport_header = 0; if (!skb) { return 0; } bpf_probe_read_kernel(&head, sizeof(head), &skb->head); bpf_probe_read_kernel(&mac_header, sizeof(mac_header), &skb->mac_header); bpf_probe_read_kernel(&network_header, sizeof(network_header), &skb->network_header); bpf_probe_read_kernel(&transport_header, sizeof(transport_header), &skb->transport_header); bpf_probe_read_kernel(&event->length, sizeof(event->length), &skb->len); bpf_probe_read_kernel(&event->skb_mark, sizeof(event->skb_mark), &skb->mark); if (!head) { return 0; } struct ethhdr eth = {}; unsigned char *eth_ptr = head + mac_header; bpf_probe_read_kernel(ð, sizeof(eth), eth_ptr); __builtin_memcpy(event->src_mac, eth.h_source, 6); __builtin_memcpy(event->dst_mac, eth.h_dest, 6); __be16 eth_proto = eth.h_proto; unsigned char *l3_ptr = head + network_header; if (eth_proto == htons(ETH_P_8021Q) || eth_proto == htons(ETH_P_8021AD)) { struct vlan_hdr_t vlan = {}; bpf_probe_read_kernel(&vlan, sizeof(vlan), eth_ptr + sizeof(struct ethhdr)); event->vlan_id = ntohs(vlan.h_vlan_TCI) & 0x0fff; eth_proto = vlan.h_vlan_encapsulated_proto; event->eth_type_raw = ntohs(eth_proto); } else { event->eth_type_raw = ntohs(eth_proto); } if (eth_proto == htons(ETH_P_ARP)) { struct arphdr arph = {}; struct arp_eth_ipv4_t arp_body = {}; bpf_probe_read_kernel(&arph, sizeof(arph), l3_ptr); event->arp_op = ntohs(arph.ar_op); if (arph.ar_hrd == htons(ARPHRD_ETHER) && arph.ar_pro == htons(ETH_P_IP) && arph.ar_hln == ETH_ALEN && arph.ar_pln == 4) { bpf_probe_read_kernel(&arp_body, sizeof(arp_body), l3_ptr + sizeof(struct arphdr)); __builtin_memcpy(event->src_ip, arp_body.spa, 4); __builtin_memcpy(event->dst_ip, arp_body.tpa, 4); event->ip_version = 4; } return 1; } if (eth_proto == htons(ETH_P_IP)) { struct iphdr iph = {}; bpf_probe_read_kernel(&iph, sizeof(iph), l3_ptr); event->ip_version = 4; event->protocol_raw = iph.protocol; event->ip_id = ntohs(iph.id); bpf_probe_read_kernel(event->src_ip, 4, &iph.saddr); bpf_probe_read_kernel(event->dst_ip, 4, &iph.daddr); if (iph.protocol == IPPROTO_TCP) { struct tcphdr tcph = {}; unsigned char flags = 0; bpf_probe_read_kernel(&tcph, sizeof(tcph), head + transport_header); event->src_port = ntohs(tcph.source); event->dst_port = ntohs(tcph.dest); event->tcp_seq = ntohl(tcph.seq); event->tcp_ack = ntohl(tcph.ack_seq); bpf_probe_read_kernel(&flags, sizeof(flags), (void *)(head + transport_header + 13)); event->tcp_flags = flags; } else if (iph.protocol == IPPROTO_UDP) { struct udphdr udph = {}; bpf_probe_read_kernel(&udph, sizeof(udph), head + transport_header); event->src_port = ntohs(udph.source); event->dst_port = ntohs(udph.dest); } else if (iph.protocol == IPPROTO_ICMP) { struct icmphdr icmph = {}; bpf_probe_read_kernel(&icmph, sizeof(icmph), head + transport_header); event->icmp_type = icmph.type; event->icmp_code = icmph.code; } return 1; } if (eth_proto == htons(ETH_P_IPV6)) { struct ipv6hdr ip6h = {}; bpf_probe_read_kernel(&ip6h, sizeof(ip6h), l3_ptr); event->ip_version = 6; event->protocol_raw = ip6h.nexthdr; __builtin_memcpy(event->src_ip, &ip6h.saddr, 16); __builtin_memcpy(event->dst_ip, &ip6h.daddr, 16); if (ip6h.nexthdr == IPPROTO_TCP) { struct tcphdr tcph = {}; unsigned char flags = 0; bpf_probe_read_kernel(&tcph, sizeof(tcph), head + transport_header); event->src_port = ntohs(tcph.source); event->dst_port = ntohs(tcph.dest); event->tcp_seq = ntohl(tcph.seq); event->tcp_ack = ntohl(tcph.ack_seq); bpf_probe_read_kernel(&flags, sizeof(flags), (void *)(head + transport_header + 13)); event->tcp_flags = flags; } else if (ip6h.nexthdr == IPPROTO_UDP) { struct udphdr udph = {}; bpf_probe_read_kernel(&udph, sizeof(udph), head + transport_header); event->src_port = ntohs(udph.source); event->dst_port = ntohs(udph.dest); } else if (ip6h.nexthdr == IPPROTO_ICMPV6) { struct icmp6hdr icmp6 = {}; bpf_probe_read_kernel(&icmp6, sizeof(icmp6), head + transport_header); event->icmp_type = icmp6.icmp6_type; event->icmp_code = icmp6.icmp6_code; } return 1; } return 1; } static __always_inline int emit_event(struct pt_regs *ctx, struct sk_buff *skb, __u8 event_type, __u32 reason) { struct event_t event = {}; event.ts_ns = bpf_ktime_get_ns(); event.event_type = event_type; event.reason = reason; if (!fill_ifname(skb, &event)) { return 0; } if (!parse_skb(skb, &event)) { return 0; } events.perf_submit(ctx, &event, sizeof(event)); return 0; } int trace_ingress(struct pt_regs *ctx, struct sk_buff *skb) { return emit_event(ctx, skb, EVENT_INGRESS, 0); } int trace_egress(struct pt_regs *ctx, struct sk_buff *skb) { return emit_event(ctx, skb, EVENT_EGRESS, 0); } TRACEPOINT_PROBE(skb, kfree_skb) { struct sk_buff *skb = (struct sk_buff *)args->skbaddr; struct event_t event = {}; event.ts_ns = bpf_ktime_get_ns(); event.event_type = EVENT_DROP; event.reason = args->reason; if (!fill_ifname(skb, &event)) { return 0; } if (!parse_skb(skb, &event)) { return 0; } events.perf_submit(args, &event, sizeof(event)); return 0; } TRACEPOINT_PROBE(net, net_dev_queue) { struct sk_buff *skb = (struct sk_buff *)args->skbaddr; struct event_t event = {}; event.ts_ns = bpf_ktime_get_ns(); event.event_type = EVENT_EGRESS; event.reason = 0; if (!fill_ifname(skb, &event)) { return 0; } if (!parse_skb(skb, &event)) { return 0; } events.perf_submit(args, &event, sizeof(event)); return 0; } """ class Event(ct.Structure): _fields_ = [ ("ts_ns", ct.c_ulonglong), ("skb_mark", ct.c_uint), ("length", ct.c_uint), ("reason", ct.c_uint), ("eth_type_raw", ct.c_ushort), ("vlan_id", ct.c_ushort), ("src_port", ct.c_ushort), ("dst_port", ct.c_ushort), ("ip_id", ct.c_ushort), ("arp_op", ct.c_ushort), ("protocol_raw", ct.c_uint), ("tcp_seq", ct.c_uint), ("tcp_ack", ct.c_uint), ("event_type", ct.c_ubyte), ("ip_version", ct.c_ubyte), ("icmp_type", ct.c_ubyte), ("icmp_code", ct.c_ubyte), ("tcp_flags", ct.c_ubyte), ("ifname", ct.c_char * 16), ("src_mac", ct.c_ubyte * 6), ("dst_mac", ct.c_ubyte * 6), ("src_ip", ct.c_ubyte * 16), ("dst_ip", ct.c_ubyte * 16), ] def _mac_to_str(value: Iterable[int]) -> str: return ":".join(f"{byte:02x}" for byte in value) def _ip_to_str(ip_version: int, raw: Iterable[int]) -> str | None: data = bytes(raw) if ip_version == 4: try: return str(ipaddress.IPv4Address(data[:4])) except ipaddress.AddressValueError: return None if ip_version == 6: try: return str(ipaddress.IPv6Address(data[:16])) except ipaddress.AddressValueError: return None return None def _build_packet_uid(payload: dict[str, object]) -> str: normalized = [] for field in IDENTITY_FIELDS: value = payload.get(field) normalized.append("" if value is None else str(value)) return hashlib.sha1("|".join(normalized).encode("utf-8")).hexdigest() def _event_name(value: int) -> str: return {1: "ingress", 2: "egress", 3: "drop"}.get(value, "unknown") def _reason_name(reason: int) -> str: return f"skb_drop_reason_{reason}" def _emit_event(cpu: int, data: int, size: int) -> None: del cpu, size event = ct.cast(data, ct.POINTER(Event)).contents iface = bytes(event.ifname).split(b"\x00", 1)[0].decode("utf-8", "replace") if iface not in TARGET_INTERFACES: return payload: dict[str, object] = { "event_type": _event_name(event.event_type), "iface": iface, "skb_mark": int(event.skb_mark) or None, "length": int(event.length), "src_mac": _mac_to_str(event.src_mac), "dst_mac": _mac_to_str(event.dst_mac), "eth_type_raw": int(event.eth_type_raw) or None, "vlan_id": int(event.vlan_id) or None, "src_ip": _ip_to_str(int(event.ip_version), event.src_ip), "dst_ip": _ip_to_str(int(event.ip_version), event.dst_ip), "protocol_raw": int(event.protocol_raw) or None, "src_port": int(event.src_port) or None, "dst_port": int(event.dst_port) or None, "ip_id": int(event.ip_id) or None, "arp_op": int(event.arp_op) or None, "icmp_type": int(event.icmp_type) or None, "icmp_code": int(event.icmp_code) or None, "tcp_seq": int(event.tcp_seq) or None, "tcp_ack": int(event.tcp_ack) or None, "tcp_flags": int(event.tcp_flags) or None, "reason": _reason_name(int(event.reason)) if event.event_type == 3 else None, "reason_code": int(event.reason) if event.event_type == 3 else None, } packet_id = packet_id_from_mark(payload.get("skb_mark")) if packet_id: payload["packet_id"] = packet_id payload["correlation_key"] = f"pid:{packet_id}" payload["correlation_source"] = "kernel_mark" verdict_hint = verdict_from_mark(payload.get("skb_mark")) if verdict_hint: payload["verdict_hint"] = verdict_hint else: payload["packet_uid"] = _build_packet_uid(payload) payload["correlation_key"] = f"uid:{payload['packet_uid']}" payload["correlation_source"] = "legacy_hash" print(json.dumps(payload, separators=(",", ":")), flush=True) def _attach_kprobe_first(bpf: BPF, symbols: list[str], fn_name: str) -> str: supported = _supported_kprobe_symbols(symbols) if not supported: raise RuntimeError(f"No supported kprobe symbols found for {fn_name}: {symbols}") for symbol in supported: try: bpf.attach_kprobe(event=symbol, fn_name=fn_name) return symbol except Exception: continue raise RuntimeError(f"Failed to attach {fn_name} to any of {supported}") def _supported_kprobe_symbols(symbols: list[str]) -> list[str]: try: available = set() for symbol in symbols: for candidate in BPF.get_kprobe_functions(symbol.encode()): decoded = candidate.decode("utf-8", "replace") if decoded == symbol: available.add(symbol) if available: return [symbol for symbol in symbols if symbol in available] except Exception: pass if os.path.exists("/proc/kallsyms"): try: with open("/proc/kallsyms", "r", encoding="utf-8", errors="replace") as handle: names = {line.rsplit(" ", 1)[-1].strip() for line in handle} return [symbol for symbol in symbols if symbol in names] except Exception: pass return symbols def _attach_egress_probe(bpf: BPF) -> str: symbols = ["__dev_queue_xmit", "dev_queue_xmit"] try: return _attach_kprobe_first(bpf, symbols, "trace_egress") except Exception: pass try: bpf.attach_tracepoint(tp="net:net_dev_queue", fn_name="tracepoint__net__net_dev_queue") return "tracepoint:net:net_dev_queue" except Exception as exc: raise RuntimeError( "Failed to attach egress telemetry to any of " f"{symbols} or tracepoint net:net_dev_queue" ) from exc def _parse_args() -> argparse.Namespace: parser = argparse.ArgumentParser(description="eBPF bridge telemetry collector") parser.add_argument("--ifaces", required=True, help="Comma-separated list of interfaces to keep") return parser.parse_args() def _sigterm(_signum: int, _frame: object) -> None: raise KeyboardInterrupt def main() -> int: args = _parse_args() global TARGET_INTERFACES TARGET_INTERFACES = {iface.strip() for iface in args.ifaces.split(",") if iface.strip()} if not TARGET_INTERFACES: print("No interfaces provided", file=sys.stderr) return 1 signal.signal(signal.SIGTERM, _sigterm) signal.signal(signal.SIGINT, _sigterm) bpf = BPF(text=BPF_SOURCE) ingress_symbol = _attach_kprobe_first( bpf, ["__netif_receive_skb_core", "netif_receive_skb", "__netif_receive_skb_one_core"], "trace_ingress", ) egress_symbol = _attach_egress_probe(bpf) print( json.dumps( { "status": "collector_started", "ifaces": sorted(TARGET_INTERFACES), "ingress_symbol": ingress_symbol, "egress_symbol": egress_symbol, }, separators=(",", ":"), ), flush=True, ) bpf["events"].open_perf_buffer(_emit_event, page_cnt=128) try: while True: bpf.perf_buffer_poll() except KeyboardInterrupt: return 0 TARGET_INTERFACES: set[str] = set() if __name__ == "__main__": sys.exit(main())