From baede43d4a8fa2619c06a245d4164f9670bd3488 Mon Sep 17 00:00:00 2001 From: Rafael Roquetto Date: Fri, 21 Nov 2025 14:54:11 -0700 Subject: [PATCH 01/14] Introduce TCP context propagation --- bpf/common/common.h | 2 + bpf/common/connection_info.h | 27 +- bpf/generictracer/k_tracer.c | 21 +- bpf/generictracer/protocol_http.h | 23 +- bpf/generictracer/protocol_tcp.h | 12 +- bpf/logger/bpf_dbg.h | 5 + bpf/rdns/rdns_xdp.c | 40 +- bpf/tctracer/tc_ip.h | 6 +- bpf/tpinjector/maps/egress_key_mem.h | 16 - bpf/tpinjector/maps/sk_tp_info_pid_map.h | 16 + bpf/tpinjector/tpinjector.c | 547 +++++++++++++----- .../test/integration/multiprocess_test.go | 33 ++ pkg/appolly/discover/finder.go | 20 +- pkg/config/ebpf_tracer.go | 104 +++- pkg/config/ebpf_tracer_test.go | 345 +++++++++++ .../ebpf/generictracer/generictracer.go | 7 +- pkg/internal/ebpf/tpinjector/tpinjector.go | 12 +- .../ebpf/tpinjector/tpinjector_test.go | 135 +++++ pkg/obi/config.go | 3 +- pkg/obi/os.go | 3 +- 20 files changed, 1107 insertions(+), 270 deletions(-) delete mode 100644 bpf/tpinjector/maps/egress_key_mem.h create mode 100644 bpf/tpinjector/maps/sk_tp_info_pid_map.h create mode 100644 pkg/config/ebpf_tracer_test.go create mode 100644 pkg/internal/ebpf/tpinjector/tpinjector_test.go diff --git a/bpf/common/common.h b/bpf/common/common.h index 8ce8389ad3..2893e10098 100644 --- a/bpf/common/common.h +++ b/bpf/common/common.h @@ -15,6 +15,8 @@ #pragma once +#include + #include #include diff --git a/bpf/common/connection_info.h b/bpf/common/connection_info.h index 4c947219c3..b8aa99c4c2 100644 --- a/bpf/common/connection_info.h +++ b/bpf/common/connection_info.h @@ -67,8 +67,11 @@ typedef struct connection_info_part { u8 __pad; } connection_info_part_t; -#ifdef BPF_DEBUG static __always_inline void dbg_print_http_connection_info(connection_info_t *info) { + if (!k_bpf_debug) { + return; + } + bpf_dbg_printk("[conn] s_h = %llx, s_l = %llx, s_port=%d", *(u64 *)(&info->s_addr), *(u64 *)(&info->s_addr[8]), @@ -79,18 +82,30 @@ static __always_inline void dbg_print_http_connection_info(connection_info_t *in info->d_port); } static __always_inline void dbg_print_http_connection_info_part(connection_info_part_t *info) { + if (!k_bpf_debug) { + return; + } + bpf_dbg_printk("[conn part] s_h = %llx, s_l = %llx, s_port=%d", *(u64 *)(&info->addr), *(u64 *)(&info->addr[8]), info->port); } static __always_inline void d_print_http_connection_info_part(connection_info_part_t *info) { + if (!k_bpf_debug) { + return; + } + bpf_d_printk("[conn part] s_h = %llx, s_l = %llx, s_port=%d", *(u64 *)(&info->addr), *(u64 *)(&info->addr[8]), info->port); } static __always_inline void d_print_http_connection_info(connection_info_t *info) { + if (!k_bpf_debug) { + return; + } + bpf_d_printk("[conn] s_h = %llx, s_l = %llx, s_port=%d", *(u64 *)(&info->s_addr), *(u64 *)(&info->s_addr[8]), @@ -100,16 +115,6 @@ static __always_inline void d_print_http_connection_info(connection_info_t *info *(u64 *)(&info->d_addr[8]), info->d_port); } -#else -static __always_inline void dbg_print_http_connection_info(connection_info_t *info) { -} -static __always_inline void dbg_print_http_connection_info_part(connection_info_part_t *info) { -} -static __always_inline void d_print_http_connection_info_part(connection_info_part_t *info) { -} -static __always_inline void d_print_http_connection_info(connection_info_t *info) { -} -#endif const u8 ip4ip6_prefix[] = {0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0xff, 0xff}; diff --git a/bpf/generictracer/k_tracer.c b/bpf/generictracer/k_tracer.c index 3942f3e8ba..cb163af631 100644 --- a/bpf/generictracer/k_tracer.c +++ b/bpf/generictracer/k_tracer.c @@ -7,6 +7,8 @@ #include #include +#include +#include #include #include #include @@ -647,12 +649,12 @@ int BPF_KPROBE(obi_kprobe_tcp_close, struct sock *sk, long timeout) { if (is_tcp_socket_never_connected(sk)) { cp_support_data_t *ct = bpf_map_lookup_elem(&cp_support_connect_info, &info); bpf_dbg_printk("=== possibly never connected sock %d %llx ct=%llx ===", id, sk, ct); -#ifdef BPF_DEBUG - if (ct) { + + if (k_bpf_debug && ct) { bpf_dbg_printk( "=== established %d, already failed %d ===", ct->established, ct->failed); } -#endif + if (ct && !ct->established && !ct->failed) { dbg_print_http_connection_info(&info.conn); failed_to_connect_event(&info, orig_dport, ct->ts); @@ -922,7 +924,6 @@ static __always_inline int return_recvmsg(void *ctx, struct sock *in_sock, u64 i if (parse_sock_info((struct sock *)sock_ptr, &info.conn)) { const u16 orig_dport = info.conn.d_port; - //dbg_print_http_connection_info(&info.conn); sort_connection_info(&info.conn); info.pid = pid_from_pid_tgid(id); @@ -983,14 +984,14 @@ int BPF_KPROBE(obi_kprobe_tcp_cleanup_rbuf, struct sock *sk, int copied) { bpf_dbg_printk("=== tcp_cleanup_rbuf id=%d copied_len %d ===", id, copied); -#ifdef BPF_DEBUG - connection_info_t conn = {}; + if (k_bpf_debug) { + connection_info_t conn = {}; - if (parse_sock_info(sk, &conn)) { - sort_connection_info(&conn); - dbg_print_http_connection_info(&conn); + if (parse_sock_info(sk, &conn)) { + sort_connection_info(&conn); + dbg_print_http_connection_info(&conn); + } } -#endif return return_recvmsg(ctx, sk, id, copied); } diff --git a/bpf/generictracer/protocol_http.h b/bpf/generictracer/protocol_http.h index f5970461e3..a3fdac6cb9 100644 --- a/bpf/generictracer/protocol_http.h +++ b/bpf/generictracer/protocol_http.h @@ -20,6 +20,8 @@ #include #include +#include + #include #include #include @@ -78,6 +80,7 @@ http_get_or_create_trace_info(http_connection_metadata_t *meta, set_trace_info_for_connection(conn, TRACE_TYPE_CLIENT, tp_p); // clean up so that TC does not pick it up + // FIXME do we really? bpf_map_delete_elem(&outgoing_trace_map, &e_key); return; } @@ -121,11 +124,11 @@ http_get_or_create_trace_info(http_connection_metadata_t *meta, bpf_dbg_printk("Using old traceparent id"); } -#ifdef BPF_DEBUG - unsigned char tp_buf[TP_MAX_VAL_LENGTH]; - make_tp_string(tp_buf, &tp_p->tp); - bpf_dbg_printk("tp: %s", tp_buf); -#endif + if (k_bpf_debug) { + unsigned char tp_buf[TP_MAX_VAL_LENGTH]; + make_tp_string(tp_buf, &tp_p->tp); + bpf_dbg_printk("tp: %s", tp_buf); + } u8 skip_tp_parsing = 0; @@ -175,10 +178,12 @@ http_get_or_create_trace_info(http_connection_metadata_t *meta, if (meta && meta->type != EVENT_HTTP_CLIENT) { decode_hex(tp_p->tp.parent_id, s_id, SPAN_ID_CHAR_LEN); } -#ifdef BPF_DEBUG - make_tp_string(tp_buf, &tp_p->tp); - bpf_dbg_printk("new tp: %s", tp_buf); -#endif + + if (k_bpf_debug) { + unsigned char tp_buf[TP_MAX_VAL_LENGTH]; + make_tp_string(tp_buf, &tp_p->tp); + bpf_dbg_printk("new tp: %s", tp_buf); + } } else { bpf_dbg_printk("No additional traceparent in headers, using what was made before"); } diff --git a/bpf/generictracer/protocol_tcp.h b/bpf/generictracer/protocol_tcp.h index 99273e4f7b..aa08677ea7 100644 --- a/bpf/generictracer/protocol_tcp.h +++ b/bpf/generictracer/protocol_tcp.h @@ -21,6 +21,8 @@ #include #include +#include + static __always_inline tcp_req_t *empty_tcp_req() { int zero = 0; tcp_req_t *value = bpf_map_lookup_elem(&tcp_req_mem, &zero); @@ -36,11 +38,11 @@ static __always_inline void init_new_trace(tp_info_t *tp) { urand_bytes(tp->span_id, SPAN_ID_SIZE_BYTES); __builtin_memset(tp->parent_id, 0, sizeof(tp->span_id)); -#ifdef BPF_DEBUG - unsigned char tp_buf[TP_MAX_VAL_LENGTH]; - make_tp_string(tp_buf, tp); - bpf_dbg_printk("tp: %s", tp_buf); -#endif + if (k_bpf_debug) { + unsigned char tp_buf[TP_MAX_VAL_LENGTH]; + make_tp_string(tp_buf, tp); + bpf_dbg_printk("tp: %s", tp_buf); + } } static __always_inline u8 already_tracked_tcp(const pid_connection_info_t *p_conn) { diff --git a/bpf/logger/bpf_dbg.h b/bpf/logger/bpf_dbg.h index 1cabcafda4..9ab73b06ad 100644 --- a/bpf/logger/bpf_dbg.h +++ b/bpf/logger/bpf_dbg.h @@ -22,6 +22,8 @@ #ifdef BPF_DEBUG +enum { k_bpf_debug = 1 }; + typedef struct log_info { u64 pid; unsigned char log[80]; @@ -64,6 +66,9 @@ enum bpf_func_id___x { bpf_printk(fmt, ##args); \ } #else + +enum { k_bpf_debug = 0 }; + #define bpf_dbg_printk(fmt, args...) #define bpf_d_printk(fmt, args...) #endif diff --git a/bpf/rdns/rdns_xdp.c b/bpf/rdns/rdns_xdp.c index 487c1f01f2..429367108e 100644 --- a/bpf/rdns/rdns_xdp.c +++ b/bpf/rdns/rdns_xdp.c @@ -180,25 +180,27 @@ static __always_inline void parse_dns_response(struct xdp_md *ctx, return; } -#ifdef BPF_DEBUG - const __u16 id = bpf_ntohs(*(const __be16 *)(data)); - const __u8 ra = get_bit(flags1, RA_OFFSET); - const __u8 aa = get_bit(flags0, AA_OFFSET); - const __u8 tc = get_bit(flags0, TC_OFFSET); - const __u8 rd = get_bit(flags0, RD_OFFSET); - bpf_dbg_printk("Found possible DNS response: %x!\n", id); - bpf_dbg_printk("flags[0] = %x\n", flags0); - bpf_dbg_printk("id: %x, qr: %u, opcode: %u, aa: %u, tc: %u, rd: %u, ra: %u\n", - id, - qr, - opcode, - aa, - tc, - rd, - ra); - bpf_dbg_printk("flags[1] = %x\n", flags1); - bpf_dbg_printk("z: %u, rcode: %u, qdcount = %u, ancount = %u\n", z, rcode, ancount, qdcount); -#endif //BPF_DEBUG + if (k_bpf_debug) { + [[maybe_unused]] const __u16 id = bpf_ntohs(*(const __be16 *)(data)); + [[maybe_unused]] const __u8 ra = get_bit(flags1, RA_OFFSET); + [[maybe_unused]] const __u8 aa = get_bit(flags0, AA_OFFSET); + [[maybe_unused]] const __u8 tc = get_bit(flags0, TC_OFFSET); + [[maybe_unused]] const __u8 rd = get_bit(flags0, RD_OFFSET); + + bpf_dbg_printk("Found possible DNS response: %x!\n", id); + bpf_dbg_printk("flags[0] = %x\n", flags0); + bpf_dbg_printk("id: %x, qr: %u, opcode: %u, aa: %u, tc: %u, rd: %u, ra: %u\n", + id, + qr, + opcode, + aa, + tc, + rd, + ra); + bpf_dbg_printk("flags[1] = %x\n", flags1); + bpf_dbg_printk( + "z: %u, rcode: %u, qdcount = %u, ancount = %u\n", z, rcode, ancount, qdcount); + } // Parse question sections __u32 __attribute__((unused)) dns_packet_size = 0; diff --git a/bpf/tctracer/tc_ip.h b/bpf/tctracer/tc_ip.h index 19c69966a0..11a5db41dc 100644 --- a/bpf/tctracer/tc_ip.h +++ b/bpf/tctracer/tc_ip.h @@ -20,12 +20,14 @@ static __always_inline void populate_span_id_from_tcp_info(tp_info_t *tp, protoc } static __always_inline void print_tp(tp_info_pid_t *new_tp) { -#ifdef BPF_DEBUG + if (!k_bpf_debug) { + return; + } + unsigned char tp_buf[TP_MAX_VAL_LENGTH]; make_tp_string(tp_buf, &new_tp->tp); bpf_dbg_printk("tp: %s", tp_buf); -#endif } static __always_inline void diff --git a/bpf/tpinjector/maps/egress_key_mem.h b/bpf/tpinjector/maps/egress_key_mem.h deleted file mode 100644 index 121d10445e..0000000000 --- a/bpf/tpinjector/maps/egress_key_mem.h +++ /dev/null @@ -1,16 +0,0 @@ -// Copyright The OpenTelemetry Authors -// SPDX-License-Identifier: Apache-2.0 - -#pragma once - -#include -#include - -#include - -struct { - __uint(type, BPF_MAP_TYPE_PERCPU_ARRAY); - __type(key, int); - __type(value, egress_key_t); - __uint(max_entries, 1); -} egress_key_mem SEC(".maps"); diff --git a/bpf/tpinjector/maps/sk_tp_info_pid_map.h b/bpf/tpinjector/maps/sk_tp_info_pid_map.h new file mode 100644 index 0000000000..4dba8bc0b0 --- /dev/null +++ b/bpf/tpinjector/maps/sk_tp_info_pid_map.h @@ -0,0 +1,16 @@ +// Copyright The OpenTelemetry Authors +// SPDX-License-Identifier: Apache-2.0 + +#pragma once + +#include +#include + +#include + +struct { + __uint(type, BPF_MAP_TYPE_SK_STORAGE); + __uint(map_flags, BPF_F_NO_PREALLOC); + __type(key, u32); + __type(value, tp_info_pid_t); +} sk_tp_info_pid_map SEC(".maps"); diff --git a/bpf/tpinjector/tpinjector.c b/bpf/tpinjector/tpinjector.c index 9e99f445d8..c0a63567e5 100644 --- a/bpf/tpinjector/tpinjector.c +++ b/bpf/tpinjector/tpinjector.c @@ -6,12 +6,16 @@ #include #include +#include +#include #include #include #include +#include #include #include #include +#include #include #include #include @@ -22,39 +26,145 @@ #include +#include #include #include #include +#include -#include #include #include +#include char __license[] SEC("license") = "Dual MIT/GPL"; +// Flags to control what tpinjector should inject +enum { + k_inject_http_headers = 1 << 0, // Bit 0: inject HTTP headers + k_inject_tcp_options = 1 << 1, // Bit 1: inject TCP options +}; + +volatile const u32 inject_flags = + k_inject_http_headers | k_inject_tcp_options; // default: both enabled + enum { k_tail_write_msg_traceparent = 0 }; +SCRATCH_MEM_SIZED(tp_str_buf, 64) + +#ifndef ENOMSG +#define ENOMSG 42 +#endif + +struct tp_option { + u8 kind; + u8 len; + unsigned char trace_id[TRACE_ID_SIZE_BYTES]; + unsigned char span_id[SPAN_ID_SIZE_BYTES]; +}; + +static __always_inline void +encode_hex_skb(unsigned char *dst, const unsigned char *src, u32 src_len) { + +#pragma clang loop unroll(full) + for (u32 i = 0, j = 0; i < src_len; i++) { + unsigned char p = src[i]; + + dst[j++] = hex[(p >> 4) & 0xff]; + dst[j++] = hex[p & 0x0f]; + } +} + +static __always_inline const char *tp_string_from_opt(const struct tp_option *opt) { + unsigned char *buf = tp_str_buf_mem(); + + if (!buf) { + return NULL; + } + + unsigned char *ptr = buf; + + // Version + *ptr++ = '0'; + *ptr++ = '0'; + *ptr++ = '-'; + + // Trace ID + encode_hex_skb(ptr, opt->trace_id, TRACE_ID_SIZE_BYTES); + ptr += TRACE_ID_CHAR_LEN; + + *ptr++ = '-'; + + // SpanID + encode_hex_skb(ptr, opt->span_id, SPAN_ID_SIZE_BYTES); + ptr += SPAN_ID_CHAR_LEN; + + *ptr++ = '-'; + + *ptr++ = '0'; + *ptr++ = '\0'; + + return (const char *)buf; +} + static __always_inline pid_connection_info_t *pid_conn_info_buf() { const int zero = 0; return bpf_map_lookup_elem(&pid_connection_info_mem, &zero); } -static __always_inline egress_key_t *egress_key_buf() { - const int zero = 0; - return bpf_map_lookup_elem(&egress_key_mem, &zero); +static __always_inline egress_key_t make_key(const connection_info_t *conn) { + egress_key_t e_key = { + .d_port = conn->d_port, + .s_port = conn->s_port, + }; + + sort_egress_key(&e_key); + + return e_key; +} + +// This is setup here for Go and SSL tracking. +// Essentially, when the Go or the OpenSSL userspace +// probes activate for an outgoing HTTP request they setup this +// outgoing_trace_map for us. We then know this is a connection we should +// be injecting the Traceparent in. Another place which sets up this map is +// the kprobe on tcp_sendmsg, however that happens after the sock_msg runs, +// so we have a different detection for that - protocol_detector. +static __always_inline tp_info_pid_t *get_tp_info_pid(const egress_key_t *e_key) { + return bpf_map_lookup_elem(&outgoing_trace_map, e_key); +} + +static __always_inline void set_tp_info_pid(const egress_key_t *e_key, const tp_info_pid_t *tp_p) { + bpf_map_update_elem(&outgoing_trace_map, e_key, tp_p, BPF_ANY); +} + +static __always_inline void clear_tp_info_pid(const egress_key_t *e_key) { + bpf_map_delete_elem(&outgoing_trace_map, e_key); +} + +static __always_inline u8 already_tracked(const pid_connection_info_t *p_conn) { + return already_tracked_http(p_conn) || already_tracked_tcp(p_conn) || + already_tracked_http2(p_conn); } // Extracts what we need for connection_info_t from bpf_sock_ops if the // communication is IPv4 -static __always_inline void sk_ops_extract_key_ip4(struct bpf_sock_ops *ops, - connection_info_t *conn) { - __builtin_memcpy(conn->s_addr, ip4ip6_prefix, sizeof(ip4ip6_prefix)); - conn->s_ip[3] = ops->local_ip4; - __builtin_memcpy(conn->d_addr, ip4ip6_prefix, sizeof(ip4ip6_prefix)); - conn->d_ip[3] = ops->remote_ip4; +static __always_inline connection_info_t sk_ops_extract_key_ip4(struct bpf_sock_ops *ops) { + connection_info_t conn = {}; + + const u32 local_ip4 = ops->local_ip4; + const u32 remote_ip4 = ops->remote_ip4; + const u32 local_port = ops->local_port; + const u32 remote_port = bpf_ntohl(ops->remote_port); - conn->s_port = ops->local_port; - conn->d_port = bpf_ntohl(ops->remote_port); + __builtin_memcpy(conn.s_addr, ip4ip6_prefix, sizeof(ip4ip6_prefix)); + conn.s_ip[3] = local_ip4; + __builtin_memcpy(conn.d_addr, ip4ip6_prefix, sizeof(ip4ip6_prefix)); + conn.d_ip[3] = remote_ip4; + + conn.s_port = local_port; + conn.d_port = remote_port; + + return conn; } // Extracts what we need for connection_info_t from bpf_sock_ops if the @@ -62,19 +172,29 @@ static __always_inline void sk_ops_extract_key_ip4(struct bpf_sock_ops *ops, // The order of copying the data from bpf_sock_ops matters and must match how // the struct is laid in vmlinux.h, otherwise the verifier thinks we are modifying // the context twice. -static __always_inline void sk_ops_extract_key_ip6(struct bpf_sock_ops *ops, - connection_info_t *conn) { - conn->d_ip[0] = ops->remote_ip6[0]; - conn->d_ip[1] = ops->remote_ip6[1]; - conn->d_ip[2] = ops->remote_ip6[2]; - conn->d_ip[3] = ops->remote_ip6[3]; - conn->s_ip[0] = ops->local_ip6[0]; - conn->s_ip[1] = ops->local_ip6[1]; - conn->s_ip[2] = ops->local_ip6[2]; - conn->s_ip[3] = ops->local_ip6[3]; +static __always_inline connection_info_t sk_ops_extract_key_ip6(struct bpf_sock_ops *ops) { + connection_info_t conn = {}; - conn->d_port = bpf_ntohl(ops->remote_port); - conn->s_port = ops->local_port; + conn.d_ip[0] = ops->remote_ip6[0]; + conn.d_ip[1] = ops->remote_ip6[1]; + conn.d_ip[2] = ops->remote_ip6[2]; + conn.d_ip[3] = ops->remote_ip6[3]; + conn.s_ip[0] = ops->local_ip6[0]; + conn.s_ip[1] = ops->local_ip6[1]; + conn.s_ip[2] = ops->local_ip6[2]; + conn.s_ip[3] = ops->local_ip6[3]; + + const u32 local_port = ops->local_port; + const u32 remote_port = bpf_ntohl(ops->remote_port); + + conn.d_port = remote_port; + conn.s_port = local_port; + + return conn; +} + +static __always_inline connection_info_t get_connection_info_ops(struct bpf_sock_ops *ops) { + return ops->family == AF_INET6 ? sk_ops_extract_key_ip6(ops) : sk_ops_extract_key_ip4(ops); } // Extracts what we need for connection_info_t from sk_msg_md if the @@ -110,69 +230,181 @@ static __always_inline connection_info_t sk_msg_extract_key_ip6(struct sk_msg_md return conn; } -// Helper that writes in the sock map for a sock_ops program -static __always_inline void bpf_sock_ops_establish_cb(struct bpf_sock_ops *skops) { - connection_info_t conn = {}; +static __always_inline bool +create_trace_info(u64 id, const connection_info_t *conn, tp_info_pid_t *tp_p) { + bpf_dbg_printk("=== %s ===", __FUNCTION__); - if (skops->family == AF_INET6) { - sk_ops_extract_key_ip6(skops, &conn); - } else { - sk_ops_extract_key_ip4(skops, &conn); + pid_connection_info_t *p_conn = pid_conn_info_buf(); + + if (!p_conn) { + return false; } - bpf_sock_hash_update(skops, &sock_dir, &conn, BPF_ANY); -} + const u32 pid = pid_from_pid_tgid(id); -// Tracks all outgoing sockets (BPF_SOCK_OPS_ACTIVE_ESTABLISHED_CB) -// We don't track incoming, those would be BPF_SOCK_OPS_PASSIVE_ESTABLISHED_CB -SEC("sockops") -int obi_sockmap_tracker(struct bpf_sock_ops *skops) { - switch (skops->op) { - case BPF_SOCK_OPS_ACTIVE_ESTABLISHED_CB: - bpf_sock_ops_establish_cb(skops); - break; - default: - break; + p_conn->conn = *conn; + p_conn->pid = pid; + + tp_p->tp.ts = bpf_ktime_get_ns(); + tp_p->tp.flags = 1; + tp_p->valid = 1; + tp_p->written = 0; + tp_p->pid = pid; + tp_p->req_type = EVENT_HTTP_CLIENT; //XXX double check + + urand_bytes(tp_p->tp.span_id, SPAN_ID_SIZE_BYTES); + + if (find_trace_for_client_request(p_conn, p_conn->conn.d_port, &tp_p->tp)) { + bpf_dbg_printk("found existing tp info"); + return true; } - return 0; + + bpf_dbg_printk("generating tp info"); + + new_trace_id(&tp_p->tp); + __builtin_memset(tp_p->tp.parent_id, 0, sizeof(tp_p->tp.parent_id)); + + return true; } -static __always_inline egress_key_t make_key(const connection_info_t *conn) { - egress_key_t e_key = { - .d_port = conn->d_port, - .s_port = conn->s_port, - }; +static __always_inline void bpf_sock_ops_set_flags(struct bpf_sock_ops *skops, u8 flags) { + bpf_sock_ops_cb_flags_set(skops, skops->bpf_sock_ops_cb_flags | flags); +} - sort_egress_key(&e_key); +// Helper that writes in the sock map for a sock_ops program +static __always_inline void bpf_sock_ops_active_est_cb(struct bpf_sock_ops *skops) { + connection_info_t conn = get_connection_info_ops(skops); - return e_key; + bpf_sock_hash_update(skops, &sock_dir, &conn, BPF_ANY); + bpf_sock_ops_set_flags(skops, BPF_SOCK_OPS_WRITE_HDR_OPT_CB_FLAG); } -// This is setup here for Go tracking. Essentially, when the Go userspace -// probes activate for an outgoing HTTP request they setup this -// outgoing_trace_map for us. We then know this is a connection we should -// be injecting the Traceparent in. Another place which sets up this map is -// the kprobe on tcp_sendmsg, however that happens after the sock_msg runs, -// so we have a different detection for that - protocol_detector. -static __always_inline tp_info_pid_t *get_tp_info_pid(const egress_key_t *e_key) { - return bpf_map_lookup_elem(&outgoing_trace_map, e_key); +static __always_inline void bpf_sock_ops_passive_est_cb(struct bpf_sock_ops *skops) { + bpf_sock_ops_set_flags(skops, BPF_SOCK_OPS_PARSE_ALL_HDR_OPT_CB_FLAG); } -static __always_inline void set_tp_info_pid(const egress_key_t *e_key, const tp_info_pid_t *tp_p) { - bpf_map_update_elem(&outgoing_trace_map, e_key, tp_p, BPF_ANY); +static __always_inline void bpf_sock_ops_opt_len_cb(struct bpf_sock_ops *skops) { + struct bpf_sock *sk = skops->sk; + + if (!sk) { + return; + } + + tp_info_pid_t *tp_pid = bpf_sk_storage_get(&sk_tp_info_pid_map, sk, NULL, 0); + + if (!tp_pid) { + return; + } + + const long ret = bpf_reserve_hdr_opt(skops, sizeof(struct tp_option), 0); + + if (ret != 0) { + bpf_dbg_printk("failed to reserve TCP option: %d", ret); + return; + } } -static __always_inline void clear_tp_info_pid(const egress_key_t *e_key) { - bpf_map_delete_elem(&outgoing_trace_map, e_key); +static __always_inline void bpf_sock_ops_write_hdr_cb(struct bpf_sock_ops *skops) { + struct bpf_sock *sk = skops->sk; + + if (!sk) { + return; + } + + const tp_info_pid_t *tp_pid = bpf_sk_storage_get(&sk_tp_info_pid_map, sk, NULL, 0); + + if (!tp_pid) { + bpf_dbg_printk("tp info not found"); + return; + } + + struct tp_option opt = {.kind = 25, .len = sizeof(struct tp_option)}; + + __builtin_memcpy(opt.trace_id, tp_pid->tp.trace_id, sizeof(opt.trace_id)); + __builtin_memcpy(opt.span_id, tp_pid->tp.span_id, sizeof(opt.span_id)); + + const long ret = bpf_store_hdr_opt(skops, &opt, sizeof(opt), 0); + + if (ret != 0) { + bpf_dbg_printk("failed to store option: %d", ret); + } + + if (k_bpf_debug) { + const char *tp_str = tp_string_from_opt(&opt); + + if (tp_str) { + bpf_dbg_printk("written TP to TCP options: %s", tp_str); + } + } } -static __always_inline u8 is_tracked_go_request(const tp_info_pid_t *tp) { - return tp != NULL && tp->valid; +static __always_inline void bpf_sock_ops_parse_hdr_cb(struct bpf_sock_ops *skops) { + struct tp_option opt = {}; + opt.kind = 25; + + const long ret = bpf_load_hdr_opt(skops, &opt, sizeof(opt), 0); + + if (ret == -ENOMSG) { + return; + } + + if (ret < 0) { + bpf_dbg_printk("error parsing TCP option = %d", ret); + return; + } + + if (k_bpf_debug) { + const char *tp_str = tp_string_from_opt(&opt); + + if (tp_str) { + bpf_dbg_printk("found TP in TCP options: %s", tp_str); + } + } + + tp_info_pid_t tp = {}; + tp.valid = 1; + + __builtin_memcpy(tp.tp.trace_id, opt.trace_id, sizeof(tp.tp.trace_id)); + __builtin_memcpy(tp.tp.span_id, opt.span_id, sizeof(tp.tp.span_id)); + + connection_info_t conn = get_connection_info_ops(skops); + sort_connection_info(&conn); + + dbg_print_http_connection_info(&conn); + bpf_map_update_elem(&incoming_trace_map, &conn, &tp, BPF_ANY); } -static __always_inline u8 already_tracked(const pid_connection_info_t *p_conn) { - return already_tracked_http(p_conn) || already_tracked_tcp(p_conn) || - already_tracked_http2(p_conn); +// Tracks all outgoing sockets (BPF_SOCK_OPS_ACTIVE_ESTABLISHED_CB) +// We don't track incoming, those would be BPF_SOCK_OPS_PASSIVE_ESTABLISHED_CB +SEC("sockops") +int obi_sockmap_tracker(struct bpf_sock_ops *skops) { + struct bpf_sock *sk = skops->sk; + + if (!sk) { + return 1; + } + + switch (skops->op) { + case BPF_SOCK_OPS_ACTIVE_ESTABLISHED_CB: + bpf_sock_ops_active_est_cb(skops); + break; + case BPF_SOCK_OPS_PASSIVE_ESTABLISHED_CB: + bpf_sock_ops_passive_est_cb(skops); + break; + case BPF_SOCK_OPS_HDR_OPT_LEN_CB: + bpf_sock_ops_opt_len_cb(skops); + break; + case BPF_SOCK_OPS_WRITE_HDR_OPT_CB: + bpf_sock_ops_write_hdr_cb(skops); + break; + case BPF_SOCK_OPS_PARSE_HDR_OPT_CB: + bpf_sock_ops_parse_hdr_cb(skops); + break; + default: + break; + } + + return 1; } // This code is copied from the kprobe on tcp_sendmsg and it's called from @@ -207,13 +439,17 @@ static __always_inline u8 protocol_detector(struct sk_msg_md *msg, }; bpf_probe_read_kernel(msg_buf.fallback_buf, k_kprobes_http2_buf_size, msg->data); - u16 copy_bytes = + + const u16 copy_bytes = msg_buf.real_size > k_kprobes_http2_buf_size ? msg_buf.real_size : k_kprobes_http2_buf_size; + unsigned char **msg_ptr = bpf_map_lookup_elem(&msg_buffer_mem, &(u32){0}); + if (!msg_ptr) { bpf_d_printk("protocol_detector: failed to reserve msg_buffer space"); return 0; } + bpf_probe_read_kernel(msg_ptr, copy_bytes & k_msg_buffer_size_max_mask, msg->data); bpf_map_update_elem(&msg_buffer_mem, &(u32){0}, msg_ptr, BPF_ANY); @@ -271,18 +507,6 @@ check_pkt_access(unsigned char *buf, //NOLINT(readability-non-const-parameter) return NULL; } -static __always_inline void -encode_hex_skb(unsigned char *dst, const unsigned char *src, u32 src_len) { - -#pragma clang loop unroll(full) - for (u32 i = 0, j = 0; i < src_len; i++) { - unsigned char p = src[i]; - - dst[j++] = hex[(p >> 4) & 0xff]; - dst[j++] = hex[p & 0x0f]; - } -} - static __always_inline void make_tp_string_skb(unsigned char *buf, const tp_info_t *tp, const unsigned char *end) { buf = check_pkt_access(buf, EXTEND_SIZE, end); @@ -380,77 +604,70 @@ static __always_inline bool write_msg_traceparent(struct sk_msg_md *msg, const t return extend_and_write_tp(msg, write_offset, tp); } -static __always_inline bool -create_trace_info(u64 id, const connection_info_t *conn, tp_info_pid_t *tp_p) { - bpf_dbg_printk("=== %s ===", __FUNCTION__); - - pid_connection_info_t *p_conn = pid_conn_info_buf(); +static __always_inline void schedule_write_tcp_option(struct sk_msg_md *msg, tp_info_pid_t *tp_p) { + struct bpf_sock *sk = msg->sk; - if (!p_conn) { - return false; + if (!sk) { + return; } - const u32 pid = pid_from_pid_tgid(id); - - p_conn->conn = *conn; - p_conn->pid = pid; - - tp_p->tp.ts = bpf_ktime_get_ns(); - tp_p->tp.flags = 1; - tp_p->valid = 1; - tp_p->written = 0; - tp_p->pid = pid; - tp_p->req_type = EVENT_HTTP_CLIENT; //XXX double check - - urand_bytes(tp_p->tp.span_id, SPAN_ID_SIZE_BYTES); + tp_info_pid_t *stp = + bpf_sk_storage_get(&sk_tp_info_pid_map, sk, NULL, BPF_SK_STORAGE_GET_F_CREATE); - if (find_trace_for_client_request(p_conn, p_conn->conn.d_port, &tp_p->tp)) { - bpf_dbg_printk("found existing tp info"); - return true; + if (!stp) { + return; } - bpf_dbg_printk("generating tp info"); + // associate it also with this socket for the tcp options program + *stp = *tp_p; - new_trace_id(&tp_p->tp); - __builtin_memset(tp_p->tp.parent_id, 0, sizeof(tp_p->tp.parent_id)); - - return true; + tp_p->written = 1; } -static __always_inline void -write_go_traceparent(struct sk_msg_md *msg, const egress_key_t *e_key, tp_info_pid_t *tp_pid) { - bpf_dbg_printk("writing go traceparent"); +static __always_inline void write_http_traceparent(struct sk_msg_md *msg, tp_info_pid_t *tp_pid) { + // used for the upcoming tailcall + tp_info_pid_t *tp_p = tp_buf(); - bpf_msg_pull_data(msg, 0, msg->size, 0); + if (!tp_p) { + return; + } - tp_pid->written = write_msg_traceparent(msg, &tp_pid->tp); + tp_pid->written = 1; + *tp_p = *tp_pid; - if (tp_pid->written) { - clear_tp_info_pid(e_key); - } else { - bpf_d_printk("failed to write go traceparent"); - } + bpf_tail_call(msg, &extender_jump_table, k_tail_write_msg_traceparent); + + bpf_d_printk("tailcall failed"); } -static __always_inline bool handle_go_request(struct sk_msg_md *msg, - u64 id, - const connection_info_t *conn, - const egress_key_t *e_key, - tp_info_pid_t *tp_pid) { - if (!is_tracked_go_request(tp_pid)) { - return false; +static __always_inline void handle_existing_tp_pid(struct sk_msg_md *msg, + u64 id, + const connection_info_t *conn, + const egress_key_t *e_key, + tp_info_pid_t *tp_pid) { + if (inject_flags & k_inject_tcp_options) { + schedule_write_tcp_option(msg, tp_pid); } - // We have metadata setup by the Go uprobes telling us we should extend - // this packet - if (!protocol_detector(msg, id, conn, e_key)) { - bpf_dbg_printk("found TLS or non HTTP go request, ignoring..."); - return false; + // shortcut: if valid == 0, this is not a HTTP request (likely SSL, but + // could be anything really - don't bother with protocol_detector) + if (tp_pid->valid == 0) { + clear_tp_info_pid(e_key); + return; } - write_go_traceparent(msg, e_key, tp_pid); + // check if this really is a HTTP request whose headers we can also extend + // (it could be an SSL packet instead, or just rubbish, for instance) + const bool is_http = protocol_detector(msg, id, conn, e_key); - return true; + if (is_http) { + // here we'll leave it for protocol_http clean it up + if (inject_flags & k_inject_http_headers) { + write_http_traceparent(msg, tp_pid); + } + } else { + clear_tp_info_pid(e_key); + } } // Sock_msg program which detects packets where it should add space for @@ -458,16 +675,28 @@ static __always_inline bool handle_go_request(struct sk_msg_md *msg, // Traceparent string. SEC("sk_msg") int obi_packet_extender(struct sk_msg_md *msg) { + // If neither injection method is enabled, nothing to do + if (!(inject_flags & (k_inject_http_headers | k_inject_tcp_options))) { + return SK_PASS; + } + const u64 id = bpf_get_current_pid_tgid(); const connection_info_t conn = get_connection_info(msg); const egress_key_t e_key = make_key(&conn); + bpf_dbg_printk("%s pid = %u", __FUNCTION__, id >> 32); + tp_info_pid_t *tp_pid = get_tp_info_pid(&e_key); - if (handle_go_request(msg, id, &conn, &e_key, tp_pid)) { + // Higher-level uprobes have already set the tp_pid for us (either Go, or SSL) + if (tp_pid) { + handle_existing_tp_pid(msg, id, &conn, &e_key, tp_pid); return SK_PASS; } + // At this stage, there were no previously TP information setup - it's the first + // time we are seeing this packet - so we need to detect whether this is the start + // of a new request and perform any injection if so. // Valid PID only works for kprobes since Go programs don't add their // PIDs to the PID map (we instrument the binaries), handled in the // previous check @@ -479,16 +708,21 @@ int obi_packet_extender(struct sk_msg_md *msg) { bpf_dbg_printk("MSG TO %llx:%d", conn.d_ip[3], conn.d_port); bpf_dbg_printk("MSG SIZE: %u", msg->size); + if (msg->size <= MIN_HTTP_SIZE) { + // not enough data to detect anything, bail + return SK_PASS; + } + bpf_msg_pull_data(msg, 0, msg->size, 0); // TODO: execute the protocol handlers here with tail calls, don't // rely on tcp_sendmsg to do it and record these message buffers. - // We must run the protocol detector always, the outgoing trace map - // might be setup for TCP traffic for L4 propagation. - const u8 tracked = protocol_detector(msg, id, &conn, &e_key); + const u8 is_http = protocol_detector(msg, id, &conn, &e_key); - if (!tracked || msg->size <= MIN_HTTP_SIZE) { + // at this point, we can't handle anything other than HTTP, as we need to be able + // to tell whether this is the start of a new request + if (!is_http) { return SK_PASS; } @@ -496,26 +730,31 @@ int obi_packet_extender(struct sk_msg_md *msg) { bpf_dbg_printk("ptr = %llx, end = %llx", ctx_msg_data(msg), ctx_msg_data_end(msg)); bpf_dbg_printk("BUF: '%s'", ctx_msg_data(msg)); - // used for the upcoming tailcall + // we've found the start of a new HTTP request, let's generate new TP info for it tp_info_pid_t *tp_p = tp_buf(); - egress_key_t *e_k = egress_key_buf(); - if (!tp_p || !e_k) { + if (!tp_p) { return SK_PASS; } - if (tp_pid) { - __builtin_memcpy(tp_p, tp_pid, sizeof(*tp_p)); - } else if (!create_trace_info(id, &conn, tp_p)) { - bpf_dbg_printk("no tp info found, bailing"); + if (!create_trace_info(id, &conn, tp_p)) { return SK_PASS; } - *e_k = e_key; + tp_p->written = 1; - bpf_tail_call(msg, &extender_jump_table, k_tail_write_msg_traceparent); + // associate this tp_info to this request + set_tp_info_pid(&e_key, tp_p); - bpf_d_printk("tailcall failed"); + if (inject_flags & k_inject_tcp_options) { + schedule_write_tcp_option(msg, tp_p); + } + + if (inject_flags & k_inject_http_headers) { + // write the HTTP headers + bpf_tail_call(msg, &extender_jump_table, k_tail_write_msg_traceparent); + bpf_d_printk("tailcall failed"); + } return SK_PASS; } @@ -527,20 +766,14 @@ int obi_packet_extender_write_msg_tp(struct sk_msg_md *msg) { tp_info_pid_t *tp_p = tp_buf(); - const egress_key_t *e_key = egress_key_buf(); - - if (!tp_p || !e_key) { - bpf_dbg_printk("empty tp_buf or e_key"); + if (!tp_p) { + bpf_dbg_printk("empty tp_buf"); return SK_PASS; } bpf_msg_pull_data(msg, 0, msg->size, 0); - tp_p->written = write_msg_traceparent(msg, &tp_p->tp); - - if (tp_p->written) { - set_tp_info_pid(e_key, tp_p); - } else { + if (!write_msg_traceparent(msg, &tp_p->tp)) { bpf_d_printk("failed to write traceparent"); } diff --git a/internal/test/integration/multiprocess_test.go b/internal/test/integration/multiprocess_test.go index fea653461d..6288b1f934 100644 --- a/internal/test/integration/multiprocess_test.go +++ b/internal/test/integration/multiprocess_test.go @@ -144,6 +144,39 @@ func TestMultiProcessAppCPNoIP(t *testing.T) { require.NoError(t, compose.Close()) } +func TestMultiProcessAppCPTCPOnly(t *testing.T) { + compose, err := docker.ComposeSuite("docker-compose-multiexec-host.yml", path.Join(pathOutput, "test-suite-multiexec-app-cp-tcp-only.log")) + require.NoError(t, err) + + // Test TCP-only context propagation (no HTTP headers, only TCP options) + // Explicitly disable request header tracking since we're not injecting HTTP headers + compose.Env = append(compose.Env, `OTEL_EBPF_BPF_DISABLE_BLACK_BOX_CP=1`, `OTEL_EBPF_BPF_CONTEXT_PROPAGATION=tcp`, `OTEL_EBPF_BPF_TRACK_REQUEST_HEADERS=false`) + + require.NoError(t, compose.Up()) + + t.Run("Nested traces with TCP-only propagation", func(t *testing.T) { + testNestedHTTPTracesKProbes(t) + }) + + require.NoError(t, compose.Close()) +} + +func TestMultiProcessAppCPHeadersAndTCP(t *testing.T) { + compose, err := docker.ComposeSuite("docker-compose-multiexec-host.yml", path.Join(pathOutput, "test-suite-multiexec-app-cp-headers-tcp.log")) + require.NoError(t, err) + + // Test combined headers and TCP context propagation + compose.Env = append(compose.Env, `OTEL_EBPF_BPF_DISABLE_BLACK_BOX_CP=1`, `OTEL_EBPF_BPF_CONTEXT_PROPAGATION=headers,tcp`, `OTEL_EBPF_BPF_TRACK_REQUEST_HEADERS=1`) + + require.NoError(t, compose.Up()) + + t.Run("Nested traces with headers and TCP propagation", func(t *testing.T) { + testNestedHTTPTracesKProbes(t) + }) + + require.NoError(t, compose.Close()) +} + // Addresses bug https://github.com/grafana/beyla/issues/370 for Go executables // Prevents that two instances of the same process report traces or metrics by duplicate func checkReportedOnlyOnce(t *testing.T, baseURL, serviceName string) { diff --git a/pkg/appolly/discover/finder.go b/pkg/appolly/discover/finder.go index f1b91d0357..d69e5ca4a9 100644 --- a/pkg/appolly/discover/finder.go +++ b/pkg/appolly/discover/finder.go @@ -8,7 +8,6 @@ import ( "fmt" "go.opentelemetry.io/obi/pkg/appolly/app/request" - "go.opentelemetry.io/obi/pkg/config" "go.opentelemetry.io/obi/pkg/ebpf" ebpfcommon "go.opentelemetry.io/obi/pkg/ebpf/common" "go.opentelemetry.io/obi/pkg/export/imetrics" @@ -130,16 +129,19 @@ func (pf *ProcessFinder) Done() <-chan error { // the common tracer group should get loaded for any tracer group, only once func newCommonTracersGroup(cfg *obi.Config) []ebpf.Tracer { - switch cfg.EBPF.ContextPropagation { - case config.ContextPropagationAll: - return []ebpf.Tracer{tctracer.New(cfg), tpinjector.New(cfg)} - case config.ContextPropagationHeadersOnly: - return []ebpf.Tracer{tpinjector.New(cfg)} - case config.ContextPropagationIPOptionsOnly: - return []ebpf.Tracer{tctracer.New(cfg)} + var tracers []ebpf.Tracer + + // Add tracers based on enabled propagation modes + // tpinjector handles both HTTP headers (sk_msg) and TCP options (BPF_SOCK_OPS) + if cfg.EBPF.ContextPropagation.HasHeaders() || cfg.EBPF.ContextPropagation.HasTCP() { + tracers = append(tracers, tpinjector.New(cfg)) + } + // tctracer handles IP options only (TC egress/ingress) + if cfg.EBPF.ContextPropagation.HasIPOptions() { + tracers = append(tracers, tctracer.New(cfg)) } - return []ebpf.Tracer{} + return tracers } func newGoTracersGroup(pidFilter ebpfcommon.ServiceFilter, cfg *obi.Config, metrics imetrics.Reporter) []ebpf.Tracer { diff --git a/pkg/config/ebpf_tracer.go b/pkg/config/ebpf_tracer.go index 32801f53dc..2e90619212 100644 --- a/pkg/config/ebpf_tracer.go +++ b/pkg/config/ebpf_tracer.go @@ -19,10 +19,19 @@ type RedisDBCacheConfig struct { } const ( - ContextPropagationAll = ContextPropagationMode(iota) - ContextPropagationHeadersOnly - ContextPropagationIPOptionsOnly - ContextPropagationDisabled + ContextPropagationDisabled ContextPropagationMode = 0 + ContextPropagationHeaders ContextPropagationMode = 1 << 0 // HTTP headers + ContextPropagationTCP ContextPropagationMode = 1 << 1 // TCP options + ContextPropagationIPOptions ContextPropagationMode = 1 << 2 // IP options + + // Convenience aliases + ContextPropagationAll = ContextPropagationHeaders | ContextPropagationTCP | ContextPropagationIPOptions +) + +// Deprecated aliases for backwards compatibility +const ( + ContextPropagationHeadersOnly = ContextPropagationHeaders + ContextPropagationIPOptionsOnly = ContextPropagationIPOptions ) // EBPFTracer configuration for eBPF programs @@ -58,7 +67,8 @@ type EBPFTracer struct { ContextPropagationEnabled bool `yaml:"enable_context_propagation" env:"OTEL_EBPF_BPF_ENABLE_CONTEXT_PROPAGATION" validate:"boolean"` // Enables distributed context propagation. - ContextPropagation ContextPropagationMode `yaml:"context_propagation" env:"OTEL_EBPF_BPF_CONTEXT_PROPAGATION" validate:"oneof=0 1 2 3"` + // Can be a combination of: headers, tcp, ip (e.g., "headers,tcp" or "all") + ContextPropagation ContextPropagationMode `yaml:"context_propagation" env:"OTEL_EBPF_BPF_CONTEXT_PROPAGATION"` // Skips checking the kernel version for bpf_loop functionality. Some modified kernels have this // backported prior to version 5.17. @@ -144,36 +154,84 @@ func (c *EBPFTracer) IsContextPropagationEnabled() { } } +// HasHeaders returns true if HTTP headers context propagation is enabled +func (m ContextPropagationMode) HasHeaders() bool { + return m&ContextPropagationHeaders != 0 +} + +// HasTCP returns true if TCP options context propagation is enabled +func (m ContextPropagationMode) HasTCP() bool { + return m&ContextPropagationTCP != 0 +} + +// HasIPOptions returns true if IP options context propagation is enabled +func (m ContextPropagationMode) HasIPOptions() bool { + return m&ContextPropagationIPOptions != 0 +} + +// IsEnabled returns true if any context propagation is enabled +func (m ContextPropagationMode) IsEnabled() bool { + return m != ContextPropagationDisabled +} + func (m *ContextPropagationMode) UnmarshalText(text []byte) error { - switch strings.TrimSpace(string(text)) { + str := strings.TrimSpace(string(text)) + + // Handle simple cases first + switch str { case "all": *m = ContextPropagationAll return nil - case "headers": - *m = ContextPropagationHeadersOnly - return nil - case "ip": - *m = ContextPropagationIPOptionsOnly - return nil - case "disabled": + case "disabled", "": *m = ContextPropagationDisabled return nil } - return fmt.Errorf("invalid value for context_propagation: '%s'", text) + // Parse comma-separated list + parts := strings.Split(str, ",") + var result ContextPropagationMode + + for _, part := range parts { + part = strings.TrimSpace(part) + switch part { + case "headers", "http": + result |= ContextPropagationHeaders + case "tcp": + result |= ContextPropagationTCP + case "ip": + result |= ContextPropagationIPOptions + default: + return fmt.Errorf("invalid value for context_propagation: '%s' (valid: all, disabled, headers, tcp, ip)", part) + } + } + + *m = result + return nil } func (m ContextPropagationMode) MarshalText() ([]byte, error) { - switch m { - case ContextPropagationAll: - return []byte("all"), nil - case ContextPropagationHeadersOnly: - return []byte("headers"), nil - case ContextPropagationIPOptionsOnly: - return []byte("ip"), nil - case ContextPropagationDisabled: + if m == ContextPropagationDisabled { return []byte("disabled"), nil } - return nil, fmt.Errorf("invalid context propagation mode: %d", m) + if m == ContextPropagationAll { + return []byte("all"), nil + } + + var parts []string + if m.HasHeaders() { + parts = append(parts, "headers") + } + if m.HasTCP() { + parts = append(parts, "tcp") + } + if m.HasIPOptions() { + parts = append(parts, "ip") + } + + if len(parts) == 0 { + return nil, fmt.Errorf("invalid context propagation mode: %d", m) + } + + return []byte(strings.Join(parts, ",")), nil } diff --git a/pkg/config/ebpf_tracer_test.go b/pkg/config/ebpf_tracer_test.go new file mode 100644 index 0000000000..f89d94a7c1 --- /dev/null +++ b/pkg/config/ebpf_tracer_test.go @@ -0,0 +1,345 @@ +// Copyright The OpenTelemetry Authors +// SPDX-License-Identifier: Apache-2.0 + +package config + +import ( + "testing" +) + +func TestContextPropagationMode_UnmarshalText(t *testing.T) { + tests := []struct { + name string + input string + want ContextPropagationMode + wantErr bool + }{ + { + name: "all", + input: "all", + want: ContextPropagationAll, + }, + { + name: "disabled", + input: "disabled", + want: ContextPropagationDisabled, + }, + { + name: "headers only", + input: "headers", + want: ContextPropagationHeaders, + }, + { + name: "http alias", + input: "http", + want: ContextPropagationHeaders, + }, + { + name: "tcp only", + input: "tcp", + want: ContextPropagationTCP, + }, + { + name: "ip only", + input: "ip", + want: ContextPropagationIPOptions, + }, + { + name: "headers and tcp", + input: "headers,tcp", + want: ContextPropagationHeaders | ContextPropagationTCP, + }, + { + name: "tcp and ip", + input: "tcp,ip", + want: ContextPropagationTCP | ContextPropagationIPOptions, + }, + { + name: "headers and ip", + input: "headers,ip", + want: ContextPropagationHeaders | ContextPropagationIPOptions, + }, + { + name: "all three", + input: "headers,tcp,ip", + want: ContextPropagationAll, + }, + { + name: "with spaces", + input: " headers , tcp , ip ", + want: ContextPropagationAll, + }, + { + name: "invalid value", + input: "invalid", + wantErr: true, + }, + { + name: "mixed valid and invalid", + input: "headers,invalid", + wantErr: true, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + var got ContextPropagationMode + err := got.UnmarshalText([]byte(tt.input)) + + if (err != nil) != tt.wantErr { + t.Errorf("UnmarshalText() error = %v, wantErr %v", err, tt.wantErr) + return + } + + if !tt.wantErr && got != tt.want { + t.Errorf("UnmarshalText() got = %v, want %v", got, tt.want) + } + }) + } +} + +func TestContextPropagationMode_MarshalText(t *testing.T) { + tests := []struct { + name string + mode ContextPropagationMode + want string + wantErr bool + }{ + { + name: "all", + mode: ContextPropagationAll, + want: "all", + }, + { + name: "disabled", + mode: ContextPropagationDisabled, + want: "disabled", + }, + { + name: "headers only", + mode: ContextPropagationHeaders, + want: "headers", + }, + { + name: "tcp only", + mode: ContextPropagationTCP, + want: "tcp", + }, + { + name: "ip only", + mode: ContextPropagationIPOptions, + want: "ip", + }, + { + name: "headers and tcp", + mode: ContextPropagationHeaders | ContextPropagationTCP, + want: "headers,tcp", + }, + { + name: "tcp and ip", + mode: ContextPropagationTCP | ContextPropagationIPOptions, + want: "tcp,ip", + }, + { + name: "headers and ip", + mode: ContextPropagationHeaders | ContextPropagationIPOptions, + want: "headers,ip", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got, err := tt.mode.MarshalText() + + if (err != nil) != tt.wantErr { + t.Errorf("MarshalText() error = %v, wantErr %v", err, tt.wantErr) + return + } + + if !tt.wantErr && string(got) != tt.want { + t.Errorf("MarshalText() got = %v, want %v", string(got), tt.want) + } + }) + } +} + +func TestContextPropagationMode_HasMethods(t *testing.T) { + tests := []struct { + name string + mode ContextPropagationMode + wantHeaders bool + wantTCP bool + wantIPOptions bool + wantIsEnabled bool + }{ + { + name: "all", + mode: ContextPropagationAll, + wantHeaders: true, + wantTCP: true, + wantIPOptions: true, + wantIsEnabled: true, + }, + { + name: "disabled", + mode: ContextPropagationDisabled, + wantHeaders: false, + wantTCP: false, + wantIPOptions: false, + wantIsEnabled: false, + }, + { + name: "headers only", + mode: ContextPropagationHeaders, + wantHeaders: true, + wantTCP: false, + wantIPOptions: false, + wantIsEnabled: true, + }, + { + name: "tcp only", + mode: ContextPropagationTCP, + wantHeaders: false, + wantTCP: true, + wantIPOptions: false, + wantIsEnabled: true, + }, + { + name: "ip only", + mode: ContextPropagationIPOptions, + wantHeaders: false, + wantTCP: false, + wantIPOptions: true, + wantIsEnabled: true, + }, + { + name: "headers and tcp", + mode: ContextPropagationHeaders | ContextPropagationTCP, + wantHeaders: true, + wantTCP: true, + wantIPOptions: false, + wantIsEnabled: true, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + if got := tt.mode.HasHeaders(); got != tt.wantHeaders { + t.Errorf("HasHeaders() = %v, want %v", got, tt.wantHeaders) + } + if got := tt.mode.HasTCP(); got != tt.wantTCP { + t.Errorf("HasTCP() = %v, want %v", got, tt.wantTCP) + } + if got := tt.mode.HasIPOptions(); got != tt.wantIPOptions { + t.Errorf("HasIPOptions() = %v, want %v", got, tt.wantIPOptions) + } + if got := tt.mode.IsEnabled(); got != tt.wantIsEnabled { + t.Errorf("IsEnabled() = %v, want %v", got, tt.wantIsEnabled) + } + }) + } +} + +func TestContextPropagationMode_BackwardsCompatibility(t *testing.T) { + // Test that old constants still work + tests := []struct { + name string + mode ContextPropagationMode + want ContextPropagationMode + }{ + { + name: "HeadersOnly equals Headers", + mode: ContextPropagationHeadersOnly, + want: ContextPropagationHeaders, + }, + { + name: "IPOptionsOnly equals IPOptions", + mode: ContextPropagationIPOptionsOnly, + want: ContextPropagationIPOptions, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + if tt.mode != tt.want { + t.Errorf("Backwards compatibility broken: %v != %v", tt.mode, tt.want) + } + }) + } +} + +func TestContextPropagationMode_TracerLoading(t *testing.T) { + // Test which tracers should be loaded for each configuration + // tpinjector handles: HTTP headers (sk_msg) and TCP options (BPF_SOCK_OPS) + // tctracer handles: IP options only (TC egress/ingress) + tests := []struct { + name string + mode ContextPropagationMode + wantTPInject bool // should load tpinjector + wantTCTracer bool // should load tctracer + }{ + { + name: "tcp only", + mode: ContextPropagationTCP, + wantTPInject: true, + wantTCTracer: false, + }, + { + name: "headers only", + mode: ContextPropagationHeaders, + wantTPInject: true, + wantTCTracer: false, + }, + { + name: "ip only", + mode: ContextPropagationIPOptions, + wantTPInject: false, + wantTCTracer: true, + }, + { + name: "headers and tcp", + mode: ContextPropagationHeaders | ContextPropagationTCP, + wantTPInject: true, + wantTCTracer: false, + }, + { + name: "tcp and ip", + mode: ContextPropagationTCP | ContextPropagationIPOptions, + wantTPInject: true, + wantTCTracer: true, + }, + { + name: "headers and ip", + mode: ContextPropagationHeaders | ContextPropagationIPOptions, + wantTPInject: true, + wantTCTracer: true, + }, + { + name: "all", + mode: ContextPropagationAll, + wantTPInject: true, + wantTCTracer: true, + }, + { + name: "disabled", + mode: ContextPropagationDisabled, + wantTPInject: false, + wantTCTracer: false, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + // Determine what should load based on the logic in finder.go + shouldLoadTPInject := tt.mode.HasHeaders() || tt.mode.HasTCP() + shouldLoadTCTracer := tt.mode.HasIPOptions() + + if shouldLoadTPInject != tt.wantTPInject { + t.Errorf("tpinjector loading = %v, want %v", shouldLoadTPInject, tt.wantTPInject) + } + if shouldLoadTCTracer != tt.wantTCTracer { + t.Errorf("tctracer loading = %v, want %v", shouldLoadTCTracer, tt.wantTCTracer) + } + }) + } +} diff --git a/pkg/internal/ebpf/generictracer/generictracer.go b/pkg/internal/ebpf/generictracer/generictracer.go index d01e9e800a..f26204d80c 100644 --- a/pkg/internal/ebpf/generictracer/generictracer.go +++ b/pkg/internal/ebpf/generictracer/generictracer.go @@ -24,7 +24,6 @@ import ( "go.opentelemetry.io/obi/pkg/appolly/app/request" "go.opentelemetry.io/obi/pkg/appolly/app/svc" "go.opentelemetry.io/obi/pkg/appolly/discover/exec" - "go.opentelemetry.io/obi/pkg/config" ebpfcommon "go.opentelemetry.io/obi/pkg/ebpf/common" "go.opentelemetry.io/obi/pkg/export/imetrics" "go.opentelemetry.io/obi/pkg/internal/goexec" @@ -140,7 +139,7 @@ func (p *Tracer) Load() (*ebpf.CollectionSpec, error) { } if p.cfg.EBPF.TrackRequestHeaders || - p.cfg.EBPF.ContextPropagation != config.ContextPropagationDisabled { + p.cfg.EBPF.ContextPropagation.IsEnabled() { loader = LoadBpfTP if p.cfg.EBPF.BpfDebug { loader = LoadBpfTPDebug @@ -194,7 +193,7 @@ func GenericTracerConstants(cfg *obi.Config) map[string]any { } if cfg.EBPF.TrackRequestHeaders || - cfg.EBPF.ContextPropagation != config.ContextPropagationDisabled { + cfg.EBPF.ContextPropagation.IsEnabled() { m["capture_header_buffer"] = int32(1) } else { m["capture_header_buffer"] = int32(0) @@ -330,7 +329,7 @@ func (p *Tracer) KProbes() map[string]ebpfcommon.ProbeDesc { }, } - if p.cfg.EBPF.ContextPropagation != config.ContextPropagationDisabled { + if p.cfg.EBPF.ContextPropagation.IsEnabled() { // tcp_rate_check_app_limited and tcp_sendmsg_fastopen are backup // for tcp_sendmsg_locked which doesn't fire on certain kernels // if sk_msg is attached. diff --git a/pkg/internal/ebpf/tpinjector/tpinjector.go b/pkg/internal/ebpf/tpinjector/tpinjector.go index c0ea168a24..c1e36d1f88 100644 --- a/pkg/internal/ebpf/tpinjector/tpinjector.go +++ b/pkg/internal/ebpf/tpinjector/tpinjector.go @@ -70,7 +70,7 @@ func (p *Tracer) SetupTailCalls() { } func (p *Tracer) Constants() map[string]any { - m := make(map[string]any, 2) + m := make(map[string]any, 3) // The eBPF side does some basic filtering of events that do not belong to // processes which we monitor. We filter more accurately in the userspace, but @@ -84,6 +84,16 @@ func (p *Tracer) Constants() map[string]any { m["max_transaction_time"] = uint64(p.cfg.EBPF.MaxTransactionTime.Nanoseconds()) + // Set injection flags based on context propagation configuration + flags := uint32(0) + if p.cfg.EBPF.ContextPropagation.HasHeaders() { + flags |= 1 // k_inject_http_headers + } + if p.cfg.EBPF.ContextPropagation.HasTCP() { + flags |= 2 // k_inject_tcp_options + } + m["inject_flags"] = flags + return m } diff --git a/pkg/internal/ebpf/tpinjector/tpinjector_test.go b/pkg/internal/ebpf/tpinjector/tpinjector_test.go new file mode 100644 index 0000000000..6e2a175990 --- /dev/null +++ b/pkg/internal/ebpf/tpinjector/tpinjector_test.go @@ -0,0 +1,135 @@ +// Copyright The OpenTelemetry Authors +// SPDX-License-Identifier: Apache-2.0 + +//go:build linux + +package tpinjector + +import ( + "testing" + "time" + + "github.com/stretchr/testify/assert" + + "go.opentelemetry.io/obi/pkg/appolly/services" + "go.opentelemetry.io/obi/pkg/config" + "go.opentelemetry.io/obi/pkg/obi" +) + +func TestTracer_Constants_InjectFlags(t *testing.T) { + tests := []struct { + name string + contextPropagation string + expectedInjectFlags uint32 + }{ + { + name: "disabled", + contextPropagation: "disabled", + expectedInjectFlags: 0, // neither HTTP headers nor TCP options + }, + { + name: "headers only", + contextPropagation: "headers", + expectedInjectFlags: 1, // k_inject_http_headers + }, + { + name: "tcp only", + contextPropagation: "tcp", + expectedInjectFlags: 2, // k_inject_tcp_options + }, + { + name: "headers and tcp", + contextPropagation: "headers,tcp", + expectedInjectFlags: 3, // k_inject_http_headers | k_inject_tcp_options + }, + { + name: "ip only", + contextPropagation: "ip", + expectedInjectFlags: 0, // tpinjector doesn't handle IP options + }, + { + name: "all", + contextPropagation: "all", + expectedInjectFlags: 3, // k_inject_http_headers | k_inject_tcp_options (IP handled by tctracer) + }, + { + name: "tcp and ip", + contextPropagation: "tcp,ip", + expectedInjectFlags: 2, // k_inject_tcp_options only (IP handled by tctracer) + }, + { + name: "headers and ip", + contextPropagation: "headers,ip", + expectedInjectFlags: 1, // k_inject_http_headers only (IP handled by tctracer) + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + cfg := &obi.Config{ + EBPF: config.EBPFTracer{ + MaxTransactionTime: 10 * time.Second, + }, + } + err := cfg.EBPF.ContextPropagation.UnmarshalText([]byte(tt.contextPropagation)) + assert.NoError(t, err) + + tracer := New(cfg) + constants := tracer.Constants() + + // Check that inject_flags is set correctly + injectFlags, ok := constants["inject_flags"] + assert.True(t, ok, "inject_flags should be present in constants") + assert.Equal(t, tt.expectedInjectFlags, injectFlags, "inject_flags value mismatch") + + // Verify the logic + expectedFlags := uint32(0) + if cfg.EBPF.ContextPropagation.HasHeaders() { + expectedFlags |= 1 + } + if cfg.EBPF.ContextPropagation.HasTCP() { + expectedFlags |= 2 + } + assert.Equal(t, expectedFlags, injectFlags, "inject_flags should match expected calculation") + }) + } +} + +func TestTracer_Constants_FilterPids(t *testing.T) { + tests := []struct { + name string + bpfPidFilterOff bool + expectedFilterVal int32 + }{ + { + name: "filter enabled", + bpfPidFilterOff: false, + expectedFilterVal: 1, + }, + { + name: "filter disabled", + bpfPidFilterOff: true, + expectedFilterVal: 0, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + cfg := &obi.Config{ + Discovery: services.DiscoveryConfig{ + BPFPidFilterOff: tt.bpfPidFilterOff, + }, + EBPF: config.EBPFTracer{ + MaxTransactionTime: 10 * time.Second, + }, + } + + tracer := New(cfg) + constants := tracer.Constants() + + filterPids, ok := constants["filter_pids"] + assert.True(t, ok, "filter_pids should be present in constants") + assert.Equal(t, tt.expectedFilterVal, filterPids, "filter_pids value mismatch") + }) + } +} diff --git a/pkg/obi/config.go b/pkg/obi/config.go index a03cdc24be..17cf7f9d0c 100644 --- a/pkg/obi/config.go +++ b/pkg/obi/config.go @@ -423,8 +423,7 @@ func (c *Config) otelNetO11yEnabled() bool { func (c *Config) willUseTC() bool { // remove after deleting ContextPropagationEnabled - return c.EBPF.ContextPropagation == config.ContextPropagationAll || - c.EBPF.ContextPropagation == config.ContextPropagationIPOptionsOnly || + return c.EBPF.ContextPropagation.HasIPOptions() || c.EBPF.ContextPropagationEnabled || (c.Enabled(FeatureNetO11y) && c.NetworkFlows.Source == EbpfSourceTC) } diff --git a/pkg/obi/os.go b/pkg/obi/os.go index bd9855889b..7d17db276b 100644 --- a/pkg/obi/os.go +++ b/pkg/obi/os.go @@ -11,7 +11,6 @@ import ( "golang.org/x/sys/unix" - ebpfcfg "go.opentelemetry.io/obi/pkg/config" ebpfcommon "go.opentelemetry.io/obi/pkg/ebpf/common" "go.opentelemetry.io/obi/pkg/internal/helpers" ) @@ -91,7 +90,7 @@ func checkCapabilitiesForSetOptions(config *Config, caps *helpers.OSCapabilities testAndSet(caps, capError, unix.CAP_PERFMON) testAndSet(caps, capError, unix.CAP_NET_RAW) - if config.EBPF.ContextPropagation != ebpfcfg.ContextPropagationDisabled { + if config.EBPF.ContextPropagation.IsEnabled() { testAndSet(caps, capError, unix.CAP_NET_ADMIN) } } From ba6c6caf5759ea195aefb52dcaf13ffae14a4904 Mon Sep 17 00:00:00 2001 From: Rafael Roquetto Date: Fri, 21 Nov 2025 17:59:14 -0700 Subject: [PATCH 02/14] Update test matrix --- .github/workflows/workflow_integration_tests_vm.yml | 11 ++++++++--- Makefile | 2 +- 2 files changed, 9 insertions(+), 4 deletions(-) diff --git a/.github/workflows/workflow_integration_tests_vm.yml b/.github/workflows/workflow_integration_tests_vm.yml index 98294e3357..56be37720e 100644 --- a/.github/workflows/workflow_integration_tests_vm.yml +++ b/.github/workflows/workflow_integration_tests_vm.yml @@ -35,9 +35,14 @@ jobs: - name: build matrix id: build-matrix env: - # 3 partitions, one for each of: - # TestMultiProcess, TestMultiProcessAppCP, TestMultiProcessAppCPNoIP - PARTITIONS: 3 + # 5 partitions, one for each of: + # 1. TestMultiProcess + # 2. TestMultiProcessAppCP + # 3. TestMultiProcessAppCPNoIP + # 4. TestMultiProcessAppCPTCPOnly + # 5. TestMultiProcessAppCPHeadersAndTCP + + PARTITIONS: 5 TEST_TAGS: integration run: | echo -n "matrix=" >> $GITHUB_OUTPUT diff --git a/Makefile b/Makefile index 8617db0766..2a5d49dad3 100644 --- a/Makefile +++ b/Makefile @@ -344,7 +344,7 @@ integration-test-matrix-json: .PHONY: vm-integration-test-matrix-json vm-integration-test-matrix-json: - @./scripts/generate-integration-matrix.sh "$${TEST_TAGS:-integration}" internal/test/integration "$${PARTITIONS:-3}" "TestMultiProcess" + @./scripts/generate-integration-matrix.sh "$${TEST_TAGS:-integration}" internal/test/integration "$${PARTITIONS:-5}" "TestMultiProcess" .PHONY: k8s-integration-test-matrix-json k8s-integration-test-matrix-json: From 1246974f224387e6507e00c4abec97683cbb43d3 Mon Sep 17 00:00:00 2001 From: Rafael Roquetto Date: Fri, 21 Nov 2025 18:01:38 -0700 Subject: [PATCH 03/14] Test format --- pkg/config/ebpf_tracer_test.go | 84 +++++++++++++++++----------------- 1 file changed, 42 insertions(+), 42 deletions(-) diff --git a/pkg/config/ebpf_tracer_test.go b/pkg/config/ebpf_tracer_test.go index f89d94a7c1..5b9707dfff 100644 --- a/pkg/config/ebpf_tracer_test.go +++ b/pkg/config/ebpf_tracer_test.go @@ -165,60 +165,60 @@ func TestContextPropagationMode_MarshalText(t *testing.T) { func TestContextPropagationMode_HasMethods(t *testing.T) { tests := []struct { - name string - mode ContextPropagationMode - wantHeaders bool - wantTCP bool - wantIPOptions bool - wantIsEnabled bool + name string + mode ContextPropagationMode + wantHeaders bool + wantTCP bool + wantIPOptions bool + wantIsEnabled bool }{ { - name: "all", - mode: ContextPropagationAll, - wantHeaders: true, - wantTCP: true, - wantIPOptions: true, - wantIsEnabled: true, + name: "all", + mode: ContextPropagationAll, + wantHeaders: true, + wantTCP: true, + wantIPOptions: true, + wantIsEnabled: true, }, { - name: "disabled", - mode: ContextPropagationDisabled, - wantHeaders: false, - wantTCP: false, - wantIPOptions: false, - wantIsEnabled: false, + name: "disabled", + mode: ContextPropagationDisabled, + wantHeaders: false, + wantTCP: false, + wantIPOptions: false, + wantIsEnabled: false, }, { - name: "headers only", - mode: ContextPropagationHeaders, - wantHeaders: true, - wantTCP: false, - wantIPOptions: false, - wantIsEnabled: true, + name: "headers only", + mode: ContextPropagationHeaders, + wantHeaders: true, + wantTCP: false, + wantIPOptions: false, + wantIsEnabled: true, }, { - name: "tcp only", - mode: ContextPropagationTCP, - wantHeaders: false, - wantTCP: true, - wantIPOptions: false, - wantIsEnabled: true, + name: "tcp only", + mode: ContextPropagationTCP, + wantHeaders: false, + wantTCP: true, + wantIPOptions: false, + wantIsEnabled: true, }, { - name: "ip only", - mode: ContextPropagationIPOptions, - wantHeaders: false, - wantTCP: false, - wantIPOptions: true, - wantIsEnabled: true, + name: "ip only", + mode: ContextPropagationIPOptions, + wantHeaders: false, + wantTCP: false, + wantIPOptions: true, + wantIsEnabled: true, }, { - name: "headers and tcp", - mode: ContextPropagationHeaders | ContextPropagationTCP, - wantHeaders: true, - wantTCP: true, - wantIPOptions: false, - wantIsEnabled: true, + name: "headers and tcp", + mode: ContextPropagationHeaders | ContextPropagationTCP, + wantHeaders: true, + wantTCP: true, + wantIPOptions: false, + wantIsEnabled: true, }, } From 5cdd1e10c25f507e462ddeb4ab15b8b3fd5f45ab Mon Sep 17 00:00:00 2001 From: Rafael Roquetto Date: Fri, 21 Nov 2025 18:03:03 -0700 Subject: [PATCH 04/14] Fix test --- pkg/internal/ebpf/tpinjector/tpinjector_test.go | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pkg/internal/ebpf/tpinjector/tpinjector_test.go b/pkg/internal/ebpf/tpinjector/tpinjector_test.go index 6e2a175990..6ef8466fd9 100644 --- a/pkg/internal/ebpf/tpinjector/tpinjector_test.go +++ b/pkg/internal/ebpf/tpinjector/tpinjector_test.go @@ -10,6 +10,7 @@ import ( "time" "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" "go.opentelemetry.io/obi/pkg/appolly/services" "go.opentelemetry.io/obi/pkg/config" @@ -72,7 +73,7 @@ func TestTracer_Constants_InjectFlags(t *testing.T) { }, } err := cfg.EBPF.ContextPropagation.UnmarshalText([]byte(tt.contextPropagation)) - assert.NoError(t, err) + require.NoError(t, err) tracer := New(cfg) constants := tracer.Constants() From 441ccd691dda1b025afb80f1a6b36096b8640188 Mon Sep 17 00:00:00 2001 From: Rafael Roquetto Date: Fri, 21 Nov 2025 19:13:25 -0700 Subject: [PATCH 05/14] Add devdoc --- devdocs/context-propagation.md | 286 +++++++++++++++++++++++++++++++++ 1 file changed, 286 insertions(+) create mode 100644 devdocs/context-propagation.md diff --git a/devdocs/context-propagation.md b/devdocs/context-propagation.md new file mode 100644 index 0000000000..271fbca9a1 --- /dev/null +++ b/devdocs/context-propagation.md @@ -0,0 +1,286 @@ +# Context Propagation Architecture + +This document explains how OpenTelemetry context propagation works in the eBPF instrumentation, including the coordination between different injection layers and the mutual exclusion mechanism. + +## Overview + +Context propagation allows distributed tracing by injecting trace context (trace ID, span ID) into outgoing requests. The eBPF instrumentation supports multiple injection methods organized in a fallback hierarchy: + +1. **HTTP headers** (L7) - `Traceparent:` header in plaintext HTTP requests +2. **TCP options** (L4) - Custom TCP option (kind 25) for any TCP traffic +3. **IP options** (L3) - IPv4 options or IPv6 Destination Options as fallback + +## Configuration + +Context propagation is controlled via `OTEL_EBPF_BPF_CONTEXT_PROPAGATION` which accepts a comma-separated list: +- `headers` - Inject HTTP headers +- `tcp` - Inject TCP options +- `ip` - Inject IP options +- `all` - Enable all methods (default) +- `disabled` - Disable context propagation + +Examples: +- `headers,tcp` - HTTP headers for plaintext HTTP, TCP options otherwise +- `tcp,ip` - TCP options with IP options as fallback +- `tcp` - TCP options only + +## Egress (Sending) Flow + +### Execution Order + +The order in which BPF programs execute varies depending on whether Go uprobes or SSL detection is involved: + +#### Scenario A: Go HTTP or SSL/TLS (uprobes involved) + +1. **uprobes** (Go HTTP client or SSL detection) + - Populate `outgoing_trace_map` with initial trace context + - Set `valid=1` for non-SSL, `valid=0` for SSL + +2. **sk_msg (tpinjector)** + - Runs for packets in sockmap + - Can inject HTTP headers and/or schedule TCP options + - Sets `written=1` when injection succeeds + +3. **kprobe (tcp_sendmsg / protocol_http)** + - Protocol detection and trace setup + - Checks `written` flag to reuse trace info + - Deletes from `outgoing_trace_map` if tpinjector handled it + +4. **TC egress (tctracer)** + - Injects IP options if not handled by upper layers + - Checks `written` flag for mutual exclusion + +#### Scenario B: Plain HTTP (no uprobes, kprobes only) + +1. **sk_msg (tpinjector)** + - Runs first for packets in sockmap + - Protocol detector checks if HTTP + - Can inject HTTP headers and/or schedule TCP options + - Creates new trace info and sets `written=1` + +2. **kprobe (tcp_sendmsg / protocol_http)** + - Protocol detection and trace setup + - Checks `written` flag - if set, reuses trace from tpinjector + - Deletes from `outgoing_trace_map` if tpinjector handled it + +3. **TC egress (tctracer)** + - Injects IP options if not handled by upper layers + - Checks `written` flag for mutual exclusion + +#### Scenario C: Non-HTTP TCP (no uprobes, socket not in sockmap) + +1. **kprobe (tcp_sendmsg)** + - Creates trace info in `outgoing_trace_map` + - Sets `valid=1, written=0` + +2. **TC egress (tctracer)** + - Sees `written=0`, injects IP options as fallback + - Sets `valid=0` after injection + +### Mutual Exclusion Mechanism + +The `written` flag implements mutual exclusion through the natural execution order. The key principle: **only inject via one method per connection**. + +#### Case 1: Traffic in sockmap with Go/SSL uprobes + +**For SSL/TLS:** +``` +1. Uprobe sets valid=0, written=0 in outgoing_trace_map +2. tpinjector (sk_msg) runs: + - Schedules TCP options + - Sees valid=0 (SSL), deletes outgoing_trace_map entry +3. protocol_http runs: + - Lookup fails (entry deleted), skips +4. tctracer runs: + - Lookup fails (entry deleted), no IP injection +Result: TCP options only ✓ +``` + +**For Go HTTP with headers+tcp:** +``` +1. Uprobe sets valid=1, written=0 in outgoing_trace_map +2. tpinjector runs: + - Schedules TCP options + - Injects HTTP headers, sets written=1 +3. protocol_http runs: + - Sees written=1, reuses trace, deletes outgoing_trace_map +4. tctracer runs: + - Lookup fails (entry deleted), no IP injection +Result: HTTP headers + TCP options ✓ +``` + +#### Case 2: Traffic in sockmap without uprobes (plain HTTP via kprobes) + +**For plaintext HTTP with headers+tcp:** +``` +1. tpinjector runs first: + - Protocol detector identifies HTTP + - Schedules TCP options + - Injects HTTP headers + - Creates trace, sets written=1, stores in outgoing_trace_map +2. protocol_http (kprobe) runs: + - Sees written=1, reuses trace from tpinjector + - Deletes outgoing_trace_map +3. tctracer runs: + - Lookup fails (entry deleted), no IP injection +Result: HTTP headers + TCP options ✓ +``` + +**For plaintext HTTP with tcp only:** +``` +1. tpinjector runs first: + - Protocol detector identifies HTTP + - Schedules TCP options, sets written=1 + - Skips HTTP headers (inject_flags check) + - Creates trace, stores in outgoing_trace_map +2. protocol_http (kprobe) runs: + - Sees written=1, reuses trace from tpinjector + - Deletes outgoing_trace_map +3. tctracer runs: + - Lookup fails (entry deleted), no IP injection +Result: TCP options only ✓ +``` + +#### Case 3: Traffic NOT in sockmap (tpinjector doesn't run) + +**For any traffic:** +``` +1. Kprobe sets valid=1, written=0 in outgoing_trace_map +2. tpinjector doesn't run (socket not in sockmap) +3. protocol_http runs: + - Sees written=0, creates new trace + - Does NOT delete outgoing_trace_map +4. tctracer runs: + - Sees written=0, injects IP options + - Sets valid=0 (done) +Result: IP options as fallback ✓ +``` + +#### Case 4: TCP option injection fails + +If `bpf_sk_storage_get()` fails in `schedule_write_tcp_option`, the function returns early **without setting written=1**. This allows IP options to be injected as fallback. + +## Ingress (Receiving) Flow + +### Execution Order + +On ingress, the execution order is different: + +1. **TC ingress (tctracer)** - Parses IP options first +2. **BPF_SOCK_OPS (tpinjector)** - Parses TCP options second +3. **kprobe (tcp_recvmsg / protocol_http)** - Parses HTTP headers last + +### "Last One Wins" Strategy + +Unlike egress (which uses mutual exclusion), ingress uses a **"last one wins"** approach: + +1. **TC ingress** parses IP options (if present) + - Extracts trace_id from IP options + - Generates span_id from TCP seq/ack + - Stores in `incoming_trace_map` + +2. **BPF_SOCK_OPS** parses TCP options (if present) + - Extracts trace_id and span_id from TCP option + - **Overwrites** entry in `incoming_trace_map` + +3. **protocol_http** parses HTTP headers (if present) + - Extracts trace_id, span_id, flags from `Traceparent:` header + - **Overwrites** previous values + +This creates a natural priority hierarchy: +- **IP options**: Lowest priority (most likely to be stripped by middleboxes) +- **TCP options**: Medium priority (better reliability) +- **HTTP headers**: Highest priority (W3C standard, most reliable) + +### Why "Last One Wins" on Ingress? + +1. **Unknown sender behavior**: We don't control what the sender injected +2. **Natural priority**: Execution order matches reliability (most reliable parsed last) +3. **Handles redundancy**: If sender sent multiple methods, we automatically use the best one +4. **Simplicity**: No coordination logic needed between layers + +## The outgoing_trace_map + +`outgoing_trace_map` is a BPF map (type: `BPF_MAP_TYPE_HASH`) that coordinates context propagation between egress layers. It stores `tp_info_pid_t` structs keyed by connection info. + +### tp_info_pid_t::valid (u8) + +State machine tracking the injection lifecycle: +- **0**: Invalid/SSL (don't inject) OR injection complete (set by tctracer after IP injection) +- **1**: First packet seen, needs L4 span ID setup +- **2**: L4 span ID setup done, ready for injection + +**Set to 0:** +- Go uprobes: SSL connections (`go_nethttp.c`) +- Kprobes: SSL connections (`trace_common.h`) +- tctracer: After successful IP option injection (`tctracer.c::encode_data_in_ip_options`) +- trace_common: Conflicting requests or timeouts (`trace_common.h`) + +**Set to 1:** +- tpinjector: Creating new trace (`tpinjector.c::create_trace_info`) +- protocol_http: Creating new trace (`protocol_http.h::protocol_http`) +- protocol_tcp: Creating new trace (`protocol_tcp.h`) + +**Set to 2:** +- tctracer: After populating span ID from TCP seq/ack (`tctracer.c::obi_app_egress`) + +**Checked:** +- tpinjector: Skip protocol detection for SSL (`tpinjector.c::handle_existing_tp_pid`) +- tctracer: First packet handling and injection decision (`tctracer.c::obi_app_egress`) + +### tp_info_pid_t::written (u8) + +Coordination flag for mutual exclusion between egress injection layers: +- **0**: Not yet handled by tpinjector (sk_msg layer) +- **1**: Already handled by tpinjector (TCP options or HTTP headers injected) + +**Purpose**: Implements the fallback hierarchy by preventing lower layers from injecting when higher layers already succeeded. + +**Set to 0:** +- tpinjector: Initializing new trace (`tpinjector.c::create_trace_info`) +- protocol_http: Initializing new trace (`protocol_http.h::protocol_http`) +- Go uprobes: Creating client requests (`go_nethttp.c`) + +**Set to 1:** +- tpinjector: After scheduling TCP options (`tpinjector.c::schedule_write_tcp_option`) +- tpinjector: After injecting HTTP headers (`tpinjector.c::write_http_traceparent`, `tpinjector.c::obi_packet_extender`) + +**Checked:** +- protocol_http: Skip processing if tpinjector handled it (`protocol_http.h::protocol_http`) +- tctracer: Skip IP injection if upper layer handled it (`tctracer.c::obi_app_egress`) + +**Key Behavior**: The `written` flag serves two purposes: +1. **protocol_http optimization**: Reuse existing trace info, avoid regenerating span IDs +2. **tctracer mutual exclusion**: Signal that upper layer already injected context + +## The incoming_trace_map + +`incoming_trace_map` is a BPF map (type: `BPF_MAP_TYPE_HASH`) that stores parsed trace context from incoming packets. It stores `tp_info_pid_t` structs keyed by connection info. + +Unlike `outgoing_trace_map`, there is no coordination between layers - each layer independently parses and overwrites the map entry if context is found, implementing the "last one wins" strategy. + +## Summary + +1. **Egress uses mutual exclusion**: + - Upper layers (tpinjector, protocol_http) delete the `outgoing_trace_map` entry + - Lower layers (tctracer) can't inject if entry is already deleted + - Result: Only one injection method per connection + +2. **Ingress uses "last one wins"**: + - Each layer independently parses if context is present + - Later layers overwrite earlier layers + - Result: Most reliable method takes precedence + +3. **IP options are truly a fallback**: + - On egress: Only injected when TCP options fail or socket isn't in sockmap + - On ingress: Lowest priority, overwritten by TCP options or HTTP headers + +4. **SSL/TLS uses TCP options, not HTTP headers**: + - Can't inject into encrypted payload + - TCP options work before TLS handshake + - tpinjector deletes entry early to skip HTTP detection + +5. **Execution order varies by scenario**: + - Go/SSL: uprobes → tpinjector → kprobe → tctracer + - Plain HTTP (sockmap): tpinjector → kprobe → tctracer + - Non-sockmap: kprobe → tctracer From 2669d93452120adb66ce30fa90a4c87e5951f3ac Mon Sep 17 00:00:00 2001 From: Rafael Roquetto Date: Fri, 21 Nov 2025 19:27:46 -0700 Subject: [PATCH 06/14] Remove deprecated option --- pkg/config/ebpf_tracer.go | 24 ------------------------ pkg/obi/config.go | 18 ++++++------------ pkg/obi/config_test.go | 15 +++++++-------- pkg/obi/os_test.go | 9 ++++++++- 4 files changed, 21 insertions(+), 45 deletions(-) diff --git a/pkg/config/ebpf_tracer.go b/pkg/config/ebpf_tracer.go index 2e90619212..c3ab45f102 100644 --- a/pkg/config/ebpf_tracer.go +++ b/pkg/config/ebpf_tracer.go @@ -4,9 +4,7 @@ package config import ( - "errors" "fmt" - "log/slog" "strings" "time" ) @@ -63,9 +61,6 @@ type EBPFTracer struct { // Must be at least 0 HTTPRequestTimeout time.Duration `yaml:"http_request_timeout" env:"OTEL_EBPF_BPF_HTTP_REQUEST_TIMEOUT" validate:"gte=0"` - // Deprecated: equivalent to ContextPropagationAll - ContextPropagationEnabled bool `yaml:"enable_context_propagation" env:"OTEL_EBPF_BPF_ENABLE_CONTEXT_PROPAGATION" validate:"boolean"` - // Enables distributed context propagation. // Can be a combination of: headers, tcp, ip (e.g., "headers,tcp" or "all") ContextPropagation ContextPropagationMode `yaml:"context_propagation" env:"OTEL_EBPF_BPF_CONTEXT_PROPAGATION"` @@ -135,25 +130,6 @@ type EBPFBufferSizes struct { Postgres uint32 `yaml:"postgres" env:"OTEL_EBPF_BPF_BUFFER_SIZE_POSTGRES" validate:"lte=8192"` } -func (c *EBPFTracer) Validate() error { - // TODO remove after deleting ContextPropagationEnabled - if c.ContextPropagationEnabled && c.ContextPropagation != ContextPropagationDisabled { - return errors.New("ebpf.enable_context_propagation and ebpf.context_propagation in the YAML configuration file or OTEL_EBPF_BPF_ENABLE_CONTEXT_PROPAGATION and OTEL_EBPF_BPF_CONTEXT_PROPAGATION are mutually exclusive") - } - - return nil -} - -func (c *EBPFTracer) IsContextPropagationEnabled() { - // TODO deprecated (REMOVE) - // remove after deleting ContextPropagationEnabled - if c.ContextPropagationEnabled { - slog.Warn("DEPRECATION NOTICE: 'ebpf.enable_context_propagation' configuration option has been " + - "deprecated and will be removed in the future - use 'ebpf.context_propagation' instead") - c.ContextPropagation = ContextPropagationAll - } -} - // HasHeaders returns true if HTTP headers context propagation is enabled func (m ContextPropagationMode) HasHeaders() bool { return m&ContextPropagationHeaders != 0 diff --git a/pkg/obi/config.go b/pkg/obi/config.go index 17cf7f9d0c..faf6adf46d 100644 --- a/pkg/obi/config.go +++ b/pkg/obi/config.go @@ -72,13 +72,12 @@ var DefaultConfig = Config{ ShutdownTimeout: 10 * time.Second, EnforceSysCaps: false, EBPF: config.EBPFTracer{ - BatchLength: 100, - BatchTimeout: time.Second, - HTTPRequestTimeout: 0, - TCBackend: config.TCBackendAuto, - DNSRequestTimeout: 5 * time.Second, - ContextPropagationEnabled: false, - ContextPropagation: config.ContextPropagationDisabled, + BatchLength: 100, + BatchTimeout: time.Second, + HTTPRequestTimeout: 0, + TCBackend: config.TCBackendAuto, + DNSRequestTimeout: 5 * time.Second, + ContextPropagation: config.ContextPropagationDisabled, RedisDBCache: config.RedisDBCacheConfig{ Enabled: false, MaxSize: 1000, @@ -407,9 +406,6 @@ func (c *Config) Validate() error { return ConfigError("you can't enable OTEL internal metrics without enabling OTEL metrics") } - // TODO deprecated (REMOVE) - c.EBPF.IsContextPropagationEnabled() - return nil } @@ -422,9 +418,7 @@ func (c *Config) otelNetO11yEnabled() bool { } func (c *Config) willUseTC() bool { - // remove after deleting ContextPropagationEnabled return c.EBPF.ContextPropagation.HasIPOptions() || - c.EBPF.ContextPropagationEnabled || (c.Enabled(FeatureNetO11y) && c.NetworkFlows.Source == EbpfSourceTC) } diff --git a/pkg/obi/config_test.go b/pkg/obi/config_test.go index 2c28cde982..e4a6d97c22 100644 --- a/pkg/obi/config_test.go +++ b/pkg/obi/config_test.go @@ -122,14 +122,13 @@ discovery: EnforceSysCaps: false, TracePrinter: "json", EBPF: config.EBPFTracer{ - BatchLength: 100, - BatchTimeout: time.Second, - HTTPRequestTimeout: 0, - MaxTransactionTime: 5 * time.Minute, - TCBackend: config.TCBackendAuto, - DNSRequestTimeout: 5 * time.Second, - ContextPropagationEnabled: false, - ContextPropagation: config.ContextPropagationDisabled, + BatchLength: 100, + BatchTimeout: time.Second, + HTTPRequestTimeout: 0, + MaxTransactionTime: 5 * time.Minute, + TCBackend: config.TCBackendAuto, + DNSRequestTimeout: 5 * time.Second, + ContextPropagation: config.ContextPropagationDisabled, RedisDBCache: config.RedisDBCacheConfig{ Enabled: false, MaxSize: 1000, diff --git a/pkg/obi/os_test.go b/pkg/obi/os_test.go index 254c5e62de..f7c7805d5f 100644 --- a/pkg/obi/os_test.go +++ b/pkg/obi/os_test.go @@ -108,6 +108,13 @@ type capTestData struct { useTC bool } +func contextPropagationMode(useTC bool) config.ContextPropagationMode { + if useTC { + return config.ContextPropagationIPOptions + } + return config.ContextPropagationDisabled +} + var capTests = []capTestData{ // core {osCap: unix.CAP_BPF, class: capCore, kernMaj: 6, kernMin: 10, useTC: false}, @@ -150,7 +157,7 @@ func TestCheckOSCapabilities(t *testing.T) { cfg := Config{ NetworkFlows: NetworkConfig{Enable: data.class == capNet, Source: netSource(data.useTC)}, - EBPF: config.EBPFTracer{ContextPropagationEnabled: data.useTC}, + EBPF: config.EBPFTracer{ContextPropagation: contextPropagationMode(data.useTC)}, } if data.class == capApp { // activates app o11y feature From e971b80c6fd6bce41113fc7fa6737b24fcdeb11d Mon Sep 17 00:00:00 2001 From: Rafael Roquetto Date: Fri, 21 Nov 2025 19:32:10 -0700 Subject: [PATCH 07/14] Fix docs linter --- devdocs/context-propagation.md | 18 ++++++++++++++++++ 1 file changed, 18 insertions(+) diff --git a/devdocs/context-propagation.md b/devdocs/context-propagation.md index 271fbca9a1..ff9c85c9c6 100644 --- a/devdocs/context-propagation.md +++ b/devdocs/context-propagation.md @@ -13,6 +13,7 @@ Context propagation allows distributed tracing by injecting trace context (trace ## Configuration Context propagation is controlled via `OTEL_EBPF_BPF_CONTEXT_PROPAGATION` which accepts a comma-separated list: + - `headers` - Inject HTTP headers - `tcp` - Inject TCP options - `ip` - Inject IP options @@ -20,6 +21,7 @@ Context propagation is controlled via `OTEL_EBPF_BPF_CONTEXT_PROPAGATION` which - `disabled` - Disable context propagation Examples: + - `headers,tcp` - HTTP headers for plaintext HTTP, TCP options otherwise - `tcp,ip` - TCP options with IP options as fallback - `tcp` - TCP options only @@ -84,6 +86,7 @@ The `written` flag implements mutual exclusion through the natural execution ord #### Case 1: Traffic in sockmap with Go/SSL uprobes **For SSL/TLS:** + ``` 1. Uprobe sets valid=0, written=0 in outgoing_trace_map 2. tpinjector (sk_msg) runs: @@ -97,6 +100,7 @@ Result: TCP options only ✓ ``` **For Go HTTP with headers+tcp:** + ``` 1. Uprobe sets valid=1, written=0 in outgoing_trace_map 2. tpinjector runs: @@ -112,6 +116,7 @@ Result: HTTP headers + TCP options ✓ #### Case 2: Traffic in sockmap without uprobes (plain HTTP via kprobes) **For plaintext HTTP with headers+tcp:** + ``` 1. tpinjector runs first: - Protocol detector identifies HTTP @@ -127,6 +132,7 @@ Result: HTTP headers + TCP options ✓ ``` **For plaintext HTTP with tcp only:** + ``` 1. tpinjector runs first: - Protocol detector identifies HTTP @@ -144,6 +150,7 @@ Result: TCP options only ✓ #### Case 3: Traffic NOT in sockmap (tpinjector doesn't run) **For any traffic:** + ``` 1. Kprobe sets valid=1, written=0 in outgoing_trace_map 2. tpinjector doesn't run (socket not in sockmap) @@ -188,6 +195,7 @@ Unlike egress (which uses mutual exclusion), ingress uses a **"last one wins"** - **Overwrites** previous values This creates a natural priority hierarchy: + - **IP options**: Lowest priority (most likely to be stripped by middleboxes) - **TCP options**: Medium priority (better reliability) - **HTTP headers**: Highest priority (W3C standard, most reliable) @@ -206,50 +214,60 @@ This creates a natural priority hierarchy: ### tp_info_pid_t::valid (u8) State machine tracking the injection lifecycle: + - **0**: Invalid/SSL (don't inject) OR injection complete (set by tctracer after IP injection) - **1**: First packet seen, needs L4 span ID setup - **2**: L4 span ID setup done, ready for injection **Set to 0:** + - Go uprobes: SSL connections (`go_nethttp.c`) - Kprobes: SSL connections (`trace_common.h`) - tctracer: After successful IP option injection (`tctracer.c::encode_data_in_ip_options`) - trace_common: Conflicting requests or timeouts (`trace_common.h`) **Set to 1:** + - tpinjector: Creating new trace (`tpinjector.c::create_trace_info`) - protocol_http: Creating new trace (`protocol_http.h::protocol_http`) - protocol_tcp: Creating new trace (`protocol_tcp.h`) **Set to 2:** + - tctracer: After populating span ID from TCP seq/ack (`tctracer.c::obi_app_egress`) **Checked:** + - tpinjector: Skip protocol detection for SSL (`tpinjector.c::handle_existing_tp_pid`) - tctracer: First packet handling and injection decision (`tctracer.c::obi_app_egress`) ### tp_info_pid_t::written (u8) Coordination flag for mutual exclusion between egress injection layers: + - **0**: Not yet handled by tpinjector (sk_msg layer) - **1**: Already handled by tpinjector (TCP options or HTTP headers injected) **Purpose**: Implements the fallback hierarchy by preventing lower layers from injecting when higher layers already succeeded. **Set to 0:** + - tpinjector: Initializing new trace (`tpinjector.c::create_trace_info`) - protocol_http: Initializing new trace (`protocol_http.h::protocol_http`) - Go uprobes: Creating client requests (`go_nethttp.c`) **Set to 1:** + - tpinjector: After scheduling TCP options (`tpinjector.c::schedule_write_tcp_option`) - tpinjector: After injecting HTTP headers (`tpinjector.c::write_http_traceparent`, `tpinjector.c::obi_packet_extender`) **Checked:** + - protocol_http: Skip processing if tpinjector handled it (`protocol_http.h::protocol_http`) - tctracer: Skip IP injection if upper layer handled it (`tctracer.c::obi_app_egress`) **Key Behavior**: The `written` flag serves two purposes: + 1. **protocol_http optimization**: Reuse existing trace info, avoid regenerating span IDs 2. **tctracer mutual exclusion**: Signal that upper layer already injected context From 902099b23d36c2b420a08739f295df7965e98986 Mon Sep 17 00:00:00 2001 From: Rafael Roquetto Date: Fri, 21 Nov 2025 19:36:35 -0700 Subject: [PATCH 08/14] Remove magic number --- bpf/tpinjector/tpinjector.c | 9 +++++++-- 1 file changed, 7 insertions(+), 2 deletions(-) diff --git a/bpf/tpinjector/tpinjector.c b/bpf/tpinjector/tpinjector.c index c0a63567e5..a822382580 100644 --- a/bpf/tpinjector/tpinjector.c +++ b/bpf/tpinjector/tpinjector.c @@ -47,6 +47,11 @@ enum { volatile const u32 inject_flags = k_inject_http_headers | k_inject_tcp_options; // default: both enabled +// TCP option kind for OpenTelemetry context propagation +// Kind 25 is unassigned per IANA TCP Parameters registry (released 2000-12-18) +// Better than experimental options (253-254) which must not be shipped as defaults +enum { k_tcp_option_kind_otel = 25 }; + enum { k_tail_write_msg_traceparent = 0 }; SCRATCH_MEM_SIZED(tp_str_buf, 64) @@ -318,7 +323,7 @@ static __always_inline void bpf_sock_ops_write_hdr_cb(struct bpf_sock_ops *skops return; } - struct tp_option opt = {.kind = 25, .len = sizeof(struct tp_option)}; + struct tp_option opt = {.kind = k_tcp_option_kind_otel, .len = sizeof(struct tp_option)}; __builtin_memcpy(opt.trace_id, tp_pid->tp.trace_id, sizeof(opt.trace_id)); __builtin_memcpy(opt.span_id, tp_pid->tp.span_id, sizeof(opt.span_id)); @@ -340,7 +345,7 @@ static __always_inline void bpf_sock_ops_write_hdr_cb(struct bpf_sock_ops *skops static __always_inline void bpf_sock_ops_parse_hdr_cb(struct bpf_sock_ops *skops) { struct tp_option opt = {}; - opt.kind = 25; + opt.kind = k_tcp_option_kind_otel; const long ret = bpf_load_hdr_opt(skops, &opt, sizeof(opt), 0); From f48c39081a0ed6135e7702ac60a93e792d8069bc Mon Sep 17 00:00:00 2001 From: Rafael Roquetto Date: Fri, 21 Nov 2025 19:51:06 -0700 Subject: [PATCH 09/14] Fix tests --- internal/test/integration/docker-compose-nodejs-dist.yml | 2 +- pkg/obi/config_test.go | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/internal/test/integration/docker-compose-nodejs-dist.yml b/internal/test/integration/docker-compose-nodejs-dist.yml index 81bb3bf18e..f42129b1d4 100644 --- a/internal/test/integration/docker-compose-nodejs-dist.yml +++ b/internal/test/integration/docker-compose-nodejs-dist.yml @@ -54,7 +54,7 @@ services: OTEL_EBPF_LOG_LEVEL: "DEBUG" OTEL_EBPF_BPF_DEBUG: "TRUE" OTEL_EBPF_HOSTNAME: "beyla" - OTEL_EBPF_BPF_ENABLE_CONTEXT_PROPAGATION: true + OTEL_EBPF_BPF_CONTEXT_PROPAGATION: "all" OTEL_EBPF_INTERNAL_METRICS_PROMETHEUS_PORT: 8999 OTEL_EBPF_INTERNAL_METRICS_PROMETHEUS_PATH: /metrics OTEL_EBPF_BPF_HTTP_REQUEST_TIMEOUT: "5s" diff --git a/pkg/obi/config_test.go b/pkg/obi/config_test.go index e4a6d97c22..4aac96bffc 100644 --- a/pkg/obi/config_test.go +++ b/pkg/obi/config_test.go @@ -571,11 +571,11 @@ func TestDefaultLegacyExclusionFilter(t *testing.T) { } func TestWillUseTC(t *testing.T) { - env := envMap{"OTEL_EBPF_BPF_ENABLE_CONTEXT_PROPAGATION": "true"} + env := envMap{"OTEL_EBPF_BPF_CONTEXT_PROPAGATION": "ip"} cfg := loadConfig(t, env) assert.True(t, cfg.willUseTC()) - env = envMap{"OTEL_EBPF_BPF_ENABLE_CONTEXT_PROPAGATION": "false"} + env = envMap{"OTEL_EBPF_BPF_CONTEXT_PROPAGATION": "headers"} cfg = loadConfig(t, env) assert.False(t, cfg.willUseTC()) From 5b3f09590cbfc8bbb4c12c621fa117e69af6bf1c Mon Sep 17 00:00:00 2001 From: Rafael Roquetto Date: Fri, 21 Nov 2025 20:11:25 -0700 Subject: [PATCH 10/14] Remove stray comment --- bpf/generictracer/protocol_http.h | 1 - 1 file changed, 1 deletion(-) diff --git a/bpf/generictracer/protocol_http.h b/bpf/generictracer/protocol_http.h index a3fdac6cb9..8bbcbe99a6 100644 --- a/bpf/generictracer/protocol_http.h +++ b/bpf/generictracer/protocol_http.h @@ -80,7 +80,6 @@ http_get_or_create_trace_info(http_connection_metadata_t *meta, set_trace_info_for_connection(conn, TRACE_TYPE_CLIENT, tp_p); // clean up so that TC does not pick it up - // FIXME do we really? bpf_map_delete_elem(&outgoing_trace_map, &e_key); return; } From 714aa78cc5b5f59ceffde67859f6c4b445fced7f Mon Sep 17 00:00:00 2001 From: Rafael Roquetto Date: Mon, 24 Nov 2025 09:56:08 -0700 Subject: [PATCH 11/14] Reuse encode_hex --- bpf/common/trace_util.h | 8 ++++---- bpf/tpinjector/tpinjector.c | 34 ++++++++++------------------------ 2 files changed, 14 insertions(+), 28 deletions(-) diff --git a/bpf/common/trace_util.h b/bpf/common/trace_util.h index af4fc8792d..e2159cc2c8 100644 --- a/bpf/common/trace_util.h +++ b/bpf/common/trace_util.h @@ -34,8 +34,8 @@ static __always_inline void urand_bytes(unsigned char *buf, u32 size) { } } -static __always_inline void decode_hex(unsigned char *dst, const unsigned char *src, int src_len) { - for (int i = 1, j = 0; i < src_len; i += 2) { +static __always_inline void decode_hex(unsigned char *dst, const unsigned char *src, u32 src_len) { + for (u32 i = 1, j = 0; i < src_len; i += 2) { unsigned char p = src[i - 1]; unsigned char q = src[i]; @@ -49,8 +49,8 @@ static __always_inline void decode_hex(unsigned char *dst, const unsigned char * } } -static __always_inline void encode_hex(unsigned char *dst, const unsigned char *src, int src_len) { - for (int i = 0, j = 0; i < src_len; i++) { +static __always_inline void encode_hex(unsigned char *dst, const unsigned char *src, u32 src_len) { + for (u32 i = 0, j = 0; i < src_len; i++) { unsigned char p = src[i]; dst[j++] = hex[(p >> 4) & 0xff]; dst[j++] = hex[p & 0x0f]; diff --git a/bpf/tpinjector/tpinjector.c b/bpf/tpinjector/tpinjector.c index a822382580..43fe3b70ea 100644 --- a/bpf/tpinjector/tpinjector.c +++ b/bpf/tpinjector/tpinjector.c @@ -67,18 +67,6 @@ struct tp_option { unsigned char span_id[SPAN_ID_SIZE_BYTES]; }; -static __always_inline void -encode_hex_skb(unsigned char *dst, const unsigned char *src, u32 src_len) { - -#pragma clang loop unroll(full) - for (u32 i = 0, j = 0; i < src_len; i++) { - unsigned char p = src[i]; - - dst[j++] = hex[(p >> 4) & 0xff]; - dst[j++] = hex[p & 0x0f]; - } -} - static __always_inline const char *tp_string_from_opt(const struct tp_option *opt) { unsigned char *buf = tp_str_buf_mem(); @@ -94,13 +82,13 @@ static __always_inline const char *tp_string_from_opt(const struct tp_option *op *ptr++ = '-'; // Trace ID - encode_hex_skb(ptr, opt->trace_id, TRACE_ID_SIZE_BYTES); + encode_hex(ptr, opt->trace_id, TRACE_ID_SIZE_BYTES); ptr += TRACE_ID_CHAR_LEN; *ptr++ = '-'; // SpanID - encode_hex_skb(ptr, opt->span_id, SPAN_ID_SIZE_BYTES); + encode_hex(ptr, opt->span_id, SPAN_ID_SIZE_BYTES); ptr += SPAN_ID_CHAR_LEN; *ptr++ = '-'; @@ -304,7 +292,7 @@ static __always_inline void bpf_sock_ops_opt_len_cb(struct bpf_sock_ops *skops) const long ret = bpf_reserve_hdr_opt(skops, sizeof(struct tp_option), 0); if (ret != 0) { - bpf_dbg_printk("failed to reserve TCP option: %d", ret); + bpf_dbg_printk("bpf_sock_ops_opt_len_cb: failed to reserve TCP option: %d", ret); return; } } @@ -319,7 +307,7 @@ static __always_inline void bpf_sock_ops_write_hdr_cb(struct bpf_sock_ops *skops const tp_info_pid_t *tp_pid = bpf_sk_storage_get(&sk_tp_info_pid_map, sk, NULL, 0); if (!tp_pid) { - bpf_dbg_printk("tp info not found"); + bpf_dbg_printk("bpf_sock_ops_write_hdr_cb: tp info not found"); return; } @@ -331,14 +319,14 @@ static __always_inline void bpf_sock_ops_write_hdr_cb(struct bpf_sock_ops *skops const long ret = bpf_store_hdr_opt(skops, &opt, sizeof(opt), 0); if (ret != 0) { - bpf_dbg_printk("failed to store option: %d", ret); + bpf_dbg_printk("bpf_sock_ops_write_hdr_cb: failed to store option: %d", ret); } if (k_bpf_debug) { const char *tp_str = tp_string_from_opt(&opt); if (tp_str) { - bpf_dbg_printk("written TP to TCP options: %s", tp_str); + bpf_dbg_printk("bpf_sock_ops_write_hdb_cb: written TP to TCP options: %s", tp_str); } } } @@ -354,7 +342,7 @@ static __always_inline void bpf_sock_ops_parse_hdr_cb(struct bpf_sock_ops *skops } if (ret < 0) { - bpf_dbg_printk("error parsing TCP option = %d", ret); + bpf_dbg_printk("bpf_sock_ops_parse_hdr_cb: error parsing TCP option = %d", ret); return; } @@ -362,7 +350,7 @@ static __always_inline void bpf_sock_ops_parse_hdr_cb(struct bpf_sock_ops *skops const char *tp_str = tp_string_from_opt(&opt); if (tp_str) { - bpf_dbg_printk("found TP in TCP options: %s", tp_str); + bpf_dbg_printk("bpf_sock_ops_parse_hdr_cb: found TP in TCP options: %s", tp_str); } } @@ -542,13 +530,13 @@ make_tp_string_skb(unsigned char *buf, const tp_info_t *tp, const unsigned char *buf++ = '-'; // Trace ID - encode_hex_skb(buf, tp->trace_id, TRACE_ID_SIZE_BYTES); + encode_hex(buf, tp->trace_id, TRACE_ID_SIZE_BYTES); buf += TRACE_ID_CHAR_LEN; *buf++ = '-'; // SpanID - encode_hex_skb(buf, tp->span_id, SPAN_ID_SIZE_BYTES); + encode_hex(buf, tp->span_id, SPAN_ID_SIZE_BYTES); buf += SPAN_ID_CHAR_LEN; *buf++ = '-'; @@ -689,8 +677,6 @@ int obi_packet_extender(struct sk_msg_md *msg) { const connection_info_t conn = get_connection_info(msg); const egress_key_t e_key = make_key(&conn); - bpf_dbg_printk("%s pid = %u", __FUNCTION__, id >> 32); - tp_info_pid_t *tp_pid = get_tp_info_pid(&e_key); // Higher-level uprobes have already set the tp_pid for us (either Go, or SSL) From b7fb3309a806ab3fbe1ed6dae5bc637b57cd7160 Mon Sep 17 00:00:00 2001 From: Rafael Roquetto Date: Mon, 24 Nov 2025 09:56:49 -0700 Subject: [PATCH 12/14] Remove unused aliases --- pkg/config/ebpf_tracer.go | 6 ------ pkg/config/ebpf_tracer_test.go | 27 --------------------------- 2 files changed, 33 deletions(-) diff --git a/pkg/config/ebpf_tracer.go b/pkg/config/ebpf_tracer.go index c3ab45f102..67c3f76e35 100644 --- a/pkg/config/ebpf_tracer.go +++ b/pkg/config/ebpf_tracer.go @@ -26,12 +26,6 @@ const ( ContextPropagationAll = ContextPropagationHeaders | ContextPropagationTCP | ContextPropagationIPOptions ) -// Deprecated aliases for backwards compatibility -const ( - ContextPropagationHeadersOnly = ContextPropagationHeaders - ContextPropagationIPOptionsOnly = ContextPropagationIPOptions -) - // EBPFTracer configuration for eBPF programs type EBPFTracer struct { // Enables logging of eBPF program events diff --git a/pkg/config/ebpf_tracer_test.go b/pkg/config/ebpf_tracer_test.go index 5b9707dfff..93b537e291 100644 --- a/pkg/config/ebpf_tracer_test.go +++ b/pkg/config/ebpf_tracer_test.go @@ -240,33 +240,6 @@ func TestContextPropagationMode_HasMethods(t *testing.T) { } } -func TestContextPropagationMode_BackwardsCompatibility(t *testing.T) { - // Test that old constants still work - tests := []struct { - name string - mode ContextPropagationMode - want ContextPropagationMode - }{ - { - name: "HeadersOnly equals Headers", - mode: ContextPropagationHeadersOnly, - want: ContextPropagationHeaders, - }, - { - name: "IPOptionsOnly equals IPOptions", - mode: ContextPropagationIPOptionsOnly, - want: ContextPropagationIPOptions, - }, - } - - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - if tt.mode != tt.want { - t.Errorf("Backwards compatibility broken: %v != %v", tt.mode, tt.want) - } - }) - } -} func TestContextPropagationMode_TracerLoading(t *testing.T) { // Test which tracers should be loaded for each configuration From cd0221724e0ee9ffa5095a505d89bbc8103c8f42 Mon Sep 17 00:00:00 2001 From: Rafael Roquetto Date: Mon, 24 Nov 2025 10:05:26 -0700 Subject: [PATCH 13/14] Update docs --- devdocs/context-propagation.md | 29 ++++++++++++++++++++--------- 1 file changed, 20 insertions(+), 9 deletions(-) diff --git a/devdocs/context-propagation.md b/devdocs/context-propagation.md index ff9c85c9c6..b15cc9df0d 100644 --- a/devdocs/context-propagation.md +++ b/devdocs/context-propagation.md @@ -99,18 +99,29 @@ The `written` flag implements mutual exclusion through the natural execution ord Result: TCP options only ✓ ``` -**For Go HTTP with headers+tcp:** +**For Go HTTP (plaintext):** + +Go supports two approaches for HTTP header injection: +- **Approach 1 (uprobe)**: Use `bpf_probe_write_user` to inject directly into Go's HTTP buffer +- **Approach 2 (sk_msg)**: Use tpinjector to extend the packet + +The uprobe attempts approach 1 first. If successful, it deletes the `outgoing_trace_map` entry to prevent approach 2 from running: ``` -1. Uprobe sets valid=1, written=0 in outgoing_trace_map -2. tpinjector runs: +1. uprobe_persistConnRoundTrip sets valid=1, written=0 in outgoing_trace_map +2. uprobe_writeSubset attempts bpf_probe_write_user: + - If successful: deletes outgoing_trace_map entry + - If failed: entry remains for tpinjector +3. tpinjector runs (only if entry still exists): - Schedules TCP options - - Injects HTTP headers, sets written=1 -3. protocol_http runs: - - Sees written=1, reuses trace, deletes outgoing_trace_map -4. tctracer runs: - - Lookup fails (entry deleted), no IP injection -Result: HTTP headers + TCP options ✓ + - Injects HTTP headers via sk_msg, sets written=1 +4. protocol_http runs: + - If written=1: reuses trace, deletes outgoing_trace_map + - If written=0: creates new trace +5. tctracer runs: + - If entry deleted: no IP injection + - If entry exists with written=0: injects IP options +Result: HTTP headers (via uprobe OR sk_msg) + TCP options ✓ ``` #### Case 2: Traffic in sockmap without uprobes (plain HTTP via kprobes) From 45e22bc95116366ff1922ab254a52e5f3efb54cf Mon Sep 17 00:00:00 2001 From: Rafael Roquetto Date: Mon, 24 Nov 2025 10:05:46 -0700 Subject: [PATCH 14/14] Fix formatting --- pkg/config/ebpf_tracer_test.go | 1 - 1 file changed, 1 deletion(-) diff --git a/pkg/config/ebpf_tracer_test.go b/pkg/config/ebpf_tracer_test.go index 93b537e291..a88b8e39f8 100644 --- a/pkg/config/ebpf_tracer_test.go +++ b/pkg/config/ebpf_tracer_test.go @@ -240,7 +240,6 @@ func TestContextPropagationMode_HasMethods(t *testing.T) { } } - func TestContextPropagationMode_TracerLoading(t *testing.T) { // Test which tracers should be loaded for each configuration // tpinjector handles: HTTP headers (sk_msg) and TCP options (BPF_SOCK_OPS)