Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
20 changes: 10 additions & 10 deletions ebpftracer/ebpf.go

Large diffs are not rendered by default.

170 changes: 123 additions & 47 deletions ebpftracer/ebpf/l7/clickhouse.c
Original file line number Diff line number Diff line change
Expand Up @@ -9,74 +9,150 @@
#define CLICKHOUSE_SERVER_CODE_EXCEPTION 2
#define CLICKHOUSE_SERVER_CODE_END_OF_STREAM 5

#define CLICKHOUSE_COMPRESSION_NONE 0x02
#define CLICKHOUSE_COMPRESSION_LZ4 0x82
#define CLICKHOUSE_COMPRESSION_ZSTD 0x90

#define CLICKHOUSE_MIN_QUERY_SIZE 40
#define CLICKHOUSE_MAX_USER_SIZE 63
#define CLICKHOUSE_MAX_ADDRESS_SIZE 48

// the max legitimate offset within the query header is 197 (2+36+1+1+63+1+36+1+48+8)
#define CLICKHOUSE_MAX_OFFSET 200 // checking this fixes the program load on 5.4

static __always_inline
int is_clickhouse_uuid(char *buf) {
__u8 u[CLICKHOUSE_QUERY_ID_SIZE];
bpf_read(buf, u);
if (u[8] != '-' || u[13] != '-' || u[18] != '-' || u[23] != '-') {
return 0;
}
return 1;
}

static __always_inline
int is_clickhouse_query(char *buf, __u64 buf_size) {
// Need at least the header bytes to validate ClickHouse protocol
if (buf_size < CLICKHOUSE_QUERY_ID_SIZE+3) {
if (buf_size < CLICKHOUSE_MIN_QUERY_SIZE) {
return 0;
}

__u8 b[CLICKHOUSE_QUERY_ID_SIZE+3];
__u8 b[2];
bpf_read(buf, b);

// First byte must be QUERY command
if (b[0] != CLICKHOUSE_CLIENT_CODE_QUERY) {
return 0;
}

int offset = 0;
if (b[1] == 0) {
offset = 2;
} else if (b[1] == CLICKHOUSE_QUERY_ID_SIZE) {
offset = 2 + CLICKHOUSE_QUERY_ID_SIZE;
} else {
// Invalid query ID length - not ClickHouse
__u64 offset = 2;
if (b[1] == CLICKHOUSE_QUERY_ID_SIZE) {
if (!is_clickhouse_uuid(buf+2)) {
return 0;
}
offset += CLICKHOUSE_QUERY_ID_SIZE;
} else if (b[1] != 0) {
return 0;
}

// Validate query kind
if (b[offset] != CLICKHOUSE_QUERY_KIND_INITIAL && b[offset] != CLICKHOUSE_QUERY_KIND_SECONDARY) {
__u8 kind = 0;
bpf_read(buf+offset, kind);
if (kind != CLICKHOUSE_QUERY_KIND_INITIAL && kind != CLICKHOUSE_QUERY_KIND_SECONDARY) {
return 0;
}

// Additional validation: check for reasonable message structure
// ClickHouse messages should have more structured content after the query kind
if (buf_size > offset + 1) {
// Check if the next bytes look like ClickHouse protocol continuation
// This helps avoid false positives from AMQP frames that happen to match
__u8 next_byte = 0;
bpf_read(buf + offset + 1, next_byte);

// AMQP frames often have specific patterns that differ from ClickHouse
// AMQP Basic.Publish has class=60, method=40 at fixed positions
// If we see AMQP-like patterns, reject this as ClickHouse
if (offset == 2 && buf_size >= 9) {
__u16 potential_class = 0;
__u16 potential_method = 0;
bpf_read(buf + 7, potential_class);
bpf_read(buf + 9, potential_method);

// Check for AMQP Basic class (60) and Publish method (40)
if (bpf_htons(potential_class) == 60 && bpf_htons(potential_method) == 40) {
return 0; // This looks like AMQP, not ClickHouse
}
offset += 1;
__u8 len = 0;
bpf_read(buf+offset, len); // initial_user
if (len > CLICKHOUSE_MAX_USER_SIZE) {
return 0;
}
offset += 1 + len;
if (offset > CLICKHOUSE_MAX_OFFSET || offset >= buf_size) {
return 0;
}
bpf_read(buf+offset, len); // initial_query_id
if (len == CLICKHOUSE_QUERY_ID_SIZE) {
if (offset + 1 + CLICKHOUSE_QUERY_ID_SIZE > buf_size) {
return 0;
}
if (!is_clickhouse_uuid(buf+offset+1)) {
return 0;
}
} else if (len != 0) {
return 0;
}
offset += 1 + len;
if (offset > CLICKHOUSE_MAX_OFFSET || offset >= buf_size) {
return 0;
}
bpf_read(buf+offset, len); // initial_address
if (len > CLICKHOUSE_MAX_ADDRESS_SIZE) {
return 0;
}
if (len > 0) {
if (offset + 1 >= buf_size) {
return 0;
}
__u8 c = 0;
bpf_read(buf+offset+1, c);
if (!((c >= '0' && c <= '9') || c == '[' || c == ':')) {
return 0;
}
}
offset += 1 + len + 8; // interface follows the 8-byte initial_query_start_time_microseconds
if (offset > CLICKHOUSE_MAX_OFFSET || offset >= buf_size) {
return 0;
}
__u8 iface = 0;
bpf_read(buf+offset, iface); // TCP, HTTP, GRPC, MYSQL, POSTGRESQL, LOCAL, TCP_INTERSERVER, PROMETHEUS
if (iface < 1 || iface > 8) {
return 0;
}

return 1;
}
Comment thread
mayankpande88 marked this conversation as resolved.

static __always_inline
int is_clickhouse_response(char *buf, __s32 *status) {
__u8 code = 0;
bpf_read(buf, code);
if (code == CLICKHOUSE_SERVER_CODE_DATA || code == CLICKHOUSE_SERVER_CODE_END_OF_STREAM) {
*status = STATUS_OK;
return 1;
int is_clickhouse_response(char *buf, __u64 buf_size, __s32 *status) {
if (buf_size < 1) {
return 0;
}
__u8 b[3] = {};
if (buf_size < 3) {
bpf_read(buf, b[0]);
} else {
bpf_read(buf, b);
}
if (b[0] == CLICKHOUSE_SERVER_CODE_DATA) {
if (buf_size < 3) {
return 0;
}
if (b[1] != 0) { // temporary table name is always empty
return 0;
}
if (b[2] == 1) { // uncompressed block: BlockInfo field number 1
*status = STATUS_OK;
return 1;
}
if (buf_size < 2+16+1) {
return 0;
}
__u8 method = 0;
bpf_read(buf+2+16, method); // compressed block: compression method follows the 16-byte checksum
if (method == CLICKHOUSE_COMPRESSION_LZ4 || method == CLICKHOUSE_COMPRESSION_ZSTD || method == CLICKHOUSE_COMPRESSION_NONE) {
*status = STATUS_OK;
return 1;
}
return 0;
}
if (code == CLICKHOUSE_SERVER_CODE_EXCEPTION) {
if (b[0] == CLICKHOUSE_SERVER_CODE_EXCEPTION) {
if (buf_size < 1+4) {
return 0;
}
__s32 code = 0;
bpf_read(buf+1, code);
if (code <= 0 || code > 4096) {
return 0;
}
*status = STATUS_FAILED;
return 1;
}
if (b[0] == CLICKHOUSE_SERVER_CODE_END_OF_STREAM && buf_size == 1) {
*status = STATUS_OK;
return 1;
}
return 0;
}
Comment thread
mayankpande88 marked this conversation as resolved.
2 changes: 1 addition & 1 deletion ebpftracer/ebpf/l7/l7.c
Original file line number Diff line number Diff line change
Expand Up @@ -989,7 +989,7 @@ int trace_exit_read(void *ctx, __u64 id, __u32 pid, __u16 is_tls, long int ret)
} else if (e->protocol == PROTOCOL_KAFKA) {
response = is_kafka_response(payload, req->request_id);
} else if (e->protocol == PROTOCOL_CLICKHOUSE) {
response = is_clickhouse_response(payload, &e->status);
response = is_clickhouse_response(payload, ret, &e->status);
if (!response) {
discard_l7_event(e);
return 0; // keeping the query in the map
Expand Down
Loading
Loading