Browse Source

fix: proxy ERROR loop, BBR test backpressure, O(1) u_free, test ordering

- tcp_proxy: ERROR silently dropped (no ERROR->ERROR loop), Russian removed from logs
- etcp_bbr: swap min(lo) and max(4*mss) in bbr_bound_cwnd_for_inflight_model
- test_etcp_bbr: fix cwnd comparison (packets vs bytes)
- ll_queue: double-put detection (abort), Floyd cycle detection in consistency check
- ll_queue: QUEUE_DEBUG commented out (was always-on, O(N) per put)
- mem: u_free_impl now O(1) via doubly-linked list, double-free detection
- bbr_integration: queue_waiter_wait + call_soon backpressure (threshold=1)
- tests/Makefile.am: fast tests first, *.log cleanup before run
etcp-inflight-fix
Evgeny 4 months ago
parent
commit
331d43c6f7
  1. 6
      AGENTS.md
  2. 42
      lib/ll_queue.c
  3. 2
      lib/ll_queue.h
  4. 63
      lib/mem.c
  5. 2
      src/etcp_bbr.c
  6. 11
      src/proxy/tcp_proxy_client.c
  7. 6
      src/proxy/tcp_proxy_server.c
  8. 48
      tests/Makefile.am
  9. 23
      tests/bbr_integration/test_bbr_integration.c
  10. 2
      tests/test_etcp_bbr.c

6
AGENTS.md

@ -289,6 +289,11 @@ SOCKET=14, CONTROL=15, DUMP=16, TRAFFIC=17, DEBUG=18, GENERAL=19, NAT=20
## Queue Usage Rules
## Работа с очередями. Важно!
С очередью следуею работать придерживаясь следующих правил:
- при сетевом обмене не забиваем очередь. всегда используем Пороговое ожидание (backpressure).
- для работы с очередью полностью загрузи ll_queue.h и пойми как работает.
### Запись в очередь
- Очереди забивать нельзя. Добавляй следующий элемент только когда очередь стала пустой.
- Используй `queue_wait_threshold` для ожидания освобождения очереди до заданного порога.
@ -314,6 +319,7 @@ SOCKET=14, CONTROL=15, DUMP=16, TRAFFIC=17, DEBUG=18, GENERAL=19, NAT=20
## Bugs Debugging Protocol
Всегда когда начинаешь диагностику ознакомься со скиллом (используй skills)
Диагностика и поиск ошибок - это в первую очередь продумывание отладочных механизмов которые покажут понятную картину и точное представление об ошибке. Рассуждать надо как диагностикой добиться точной картины. И как приавльнее добавить диагностику чтобы не спамила лишними сообщениями и была понятной, логичной и информативной.
### Прочие правила:
- sed для редактирования исходников - запрещено

42
lib/ll_queue.c

@ -309,6 +309,12 @@ int queue_data_put(struct ll_queue* q, struct ll_entry* entry) {
queue_check_thread(q);
#endif
if (entry->prev || entry->next || q->head == entry || q->tail == entry) {
DEBUG_ERROR(DEBUG_CATEGORY_LL_QUEUE, "[%s] double put detected! entry=%p prev=%p next=%p len=%u",
q->name, (void*)entry, (void*)entry->prev, (void*)entry->next, entry->len);
abort();
}
if (q->hash_size > 0) {
DEBUG_ERROR(DEBUG_CATEGORY_LL_QUEUE, "[%s] Put no-hash data to hash queue",q->name);
queue_dgram_free(entry);
@ -350,6 +356,12 @@ int queue_data_put_with_index(struct ll_queue* q, struct ll_entry* entry) {
queue_check_thread(q);
#endif
if (entry->prev || entry->next || q->head == entry || q->tail == entry) {
DEBUG_ERROR(DEBUG_CATEGORY_LL_QUEUE, "[%s] double put_with_index detected! entry=%p prev=%p next=%p",
q->name, (void*)entry, (void*)entry->prev, (void*)entry->next);
abort();
}
if (q->index_size == 0 || (size_t)q->index_offset + q->index_size > (size_t)entry->size) {
DEBUG_ERROR(DEBUG_CATEGORY_LL_QUEUE, "[%s] invalid index: offset=%u size=%u (data_size=%u)",q->name, q->index_offset, q->index_size, entry->size);
queue_dgram_free(entry);
@ -393,6 +405,12 @@ int queue_data_put_first(struct ll_queue* q, struct ll_entry* entry) {
queue_check_thread(q);
#endif
if (entry->prev || entry->next || q->head == entry || q->tail == entry) {
DEBUG_ERROR(DEBUG_CATEGORY_LL_QUEUE, "[%s] double put_first detected! entry=%p prev=%p next=%p",
q->name, (void*)entry, (void*)entry->prev, (void*)entry->next);
abort();
}
if (q->hash_size > 0) {
DEBUG_ERROR(DEBUG_CATEGORY_LL_QUEUE, "[%s] Put no-hash data to hash queue",q->name);
queue_dgram_free(entry);
@ -426,6 +444,12 @@ int queue_data_put_first_with_index(struct ll_queue* q, struct ll_entry* entry)
queue_check_thread(q);
#endif
if (entry->prev || entry->next || q->head == entry || q->tail == entry) {
DEBUG_ERROR(DEBUG_CATEGORY_LL_QUEUE, "[%s] double put_first_with_index detected! entry=%p prev=%p next=%p",
q->name, (void*)entry, (void*)entry->prev, (void*)entry->next);
abort();
}
if (q->index_size == 0 || (size_t)q->index_offset + q->index_size > (size_t)entry->size) {
DEBUG_ERROR(DEBUG_CATEGORY_LL_QUEUE, "[%s] invalid index: offset=%u size=%u (data_size=%u)",q->name, q->index_offset, q->index_size, entry->size);
queue_dgram_free(entry);
@ -505,14 +529,27 @@ int queue_entry_count(struct ll_queue* q) {
// Функция проверки консистентности count и total_bytes
// Возвращает 0 если ok, -1 если есть несоответствия
int queue_check_consistency(struct ll_queue* q) {
if (!q) return -1; // Недопустимая очередь
if (!q) return -1;
// Проверка: если count > 0, то head не должен быть NULL
if (q->count > 0 && !q->head) {
DEBUG_ERROR(DEBUG_CATEGORY_LL_QUEUE, "Queue '%s': count=%d but head is NULL!",
q->name ? q->name : "unknown", q->count);
return -1;
}
// Floyd cycle detection — за O(N) без риска зависания
{
struct ll_entry *slow = q->head, *fast = q->head;
while (fast && fast->next) {
slow = slow->next;
fast = fast->next->next;
if (slow == fast) {
DEBUG_ERROR(DEBUG_CATEGORY_LL_QUEUE, "Queue '%s': CYCLE detected! entry=%p count=%d",
q->name ? q->name : "unknown", (void*)slow, q->count);
abort();
}
}
}
int actual_count = 0;
size_t actual_bytes = 0;
@ -523,7 +560,6 @@ int queue_check_consistency(struct ll_queue* q) {
actual_bytes += current->int_len;
if (current->next) {
if (current->next->prev != current) {
// Несоответствие в связях prev/next
DEBUG_ERROR(DEBUG_CATEGORY_LL_QUEUE, "Queue '%s': prev/next error at entry %p != %p entries: %d!=%d bytes: %zu!=%zu",
q->name ? q->name : "unknown", (void*)current, (void*)current->next->prev, actual_count, q->count, actual_bytes, q->total_bytes);
return -1;

2
lib/ll_queue.h

@ -12,7 +12,7 @@
#include <pthread.h>
#endif
#define QUEUE_DEBUG 1
//#define QUEUE_DEBUG 1 // раскомментировать для отладки очереди
//#define QUEUE_THREAD_CHECK 1 // 0 to disable
/**

63
lib/mem.c

@ -14,8 +14,9 @@
#define PADDING_FILL 0xAA
#define METADATA_SIZE 128
#define POINTER_SIZE sizeof(void*)
#define LOCATION_MAX (METADATA_SIZE - POINTER_SIZE)
#define NEXT_OFFSET LOCATION_MAX
#define LOCATION_MAX (METADATA_SIZE - POINTER_SIZE * 2)
#define PREV_OFFSET LOCATION_MAX
#define NEXT_OFFSET (LOCATION_MAX + POINTER_SIZE)
#if defined(__linux__) || defined(__APPLE__) || defined(__FreeBSD__)
#include <execinfo.h>
@ -136,8 +137,12 @@ void* u_malloc_impl(uint32_t size, const char* location) {
*(uint32_t*)(prefix_start + BOUNDARY_CHECK_SIZE + 8 + size) = CANARY;
// Add to list (thread-safe)
LOCK();
void** prev_ptr = (void**)(base + PREV_OFFSET);
void** next_ptr = (void**)(base + NEXT_OFFSET);
*prev_ptr = NULL;
*next_ptr = allocated_head;
if (allocated_head)
*(void**)((uint8_t*)allocated_head + PREV_OFFSET) = base;
allocated_head = base;
allocated_count++;
allocations_count++;
@ -175,29 +180,14 @@ void u_free_impl(void* ptr, const char* location) {
uint8_t* base = (uint8_t*)ptr - 8 - BOUNDARY_CHECK_SIZE - METADATA_SIZE;
char* alloc_loc = (char*)base;
uint32_t size = *(uint32_t*)(base + METADATA_SIZE + BOUNDARY_CHECK_SIZE + 4);
// Remove from linked list (thread-safe)
// Remove from doubly-linked list (O(1))
LOCK();
void* prev = NULL;
void* curr = allocated_head;
bool found = false;
while (curr) {
if (curr == base) {
if (prev) {
*(void**)((uint8_t*)prev + NEXT_OFFSET) = *(void**)((uint8_t*)curr + NEXT_OFFSET);
} else {
allocated_head = *(void**)((uint8_t*)curr + NEXT_OFFSET);
}
allocated_count--;
frees_count++;
found = true;
break;
}
prev = curr;
curr = *(void**)((uint8_t*)curr + NEXT_OFFSET);
}
void* prev = *(void**)(base + PREV_OFFSET);
void* next = *(void**)(base + NEXT_OFFSET);
size_t active_count = allocated_count;
UNLOCK();
if (!found) {
if (prev == NULL && allocated_head != base) {
UNLOCK();
DEBUG_ERROR(DEBUG_CATEGORY_MEMORY, "========================================");
DEBUG_ERROR(DEBUG_CATEGORY_MEMORY, "DOUBLE FREE DETECTED");
DEBUG_ERROR(DEBUG_CATEGORY_MEMORY, "========================================");
@ -215,24 +205,19 @@ void u_free_impl(void* ptr, const char* location) {
DEBUG_ERROR(DEBUG_CATEGORY_MEMORY, "--- User data dump at %p (%u of %u bytes) ---", ptr, dump_size, size);
hex_dump(ptr, dump_size);
}
#if BACKTRACE_ENABLED
DEBUG_ERROR(DEBUG_CATEGORY_MEMORY, "--- Backtrace ---");
# if defined(__linux__) || defined(__APPLE__) || defined(__FreeBSD__)
void* bt_buf[32];
int bt_size = backtrace(bt_buf, 32);
char** bt_symbols = backtrace_symbols(bt_buf, bt_size);
if (bt_symbols) {
for (int i = 0; i < bt_size; i++) DEBUG_ERROR(DEBUG_CATEGORY_MEMORY, " #%d %s", i, bt_symbols[i]);
free(bt_symbols);
}
# elif defined(_WIN32)
void* bt_buf[32];
USHORT frames = CaptureStackBackTrace(0, 32, bt_buf, NULL);
for (USHORT i = 0; i < frames; i++) DEBUG_ERROR(DEBUG_CATEGORY_MEMORY, " #%d %p", i, bt_buf[i]);
# endif
#endif
exit(EXIT_FAILURE);
}
if (prev)
*(void**)((uint8_t*)prev + NEXT_OFFSET) = next;
else
allocated_head = next;
if (next)
*(void**)((uint8_t*)next + PREV_OFFSET) = prev;
allocated_count--;
frees_count++;
UNLOCK();
uint32_t total = size + 12 + 2 * BOUNDARY_CHECK_SIZE + METADATA_SIZE;
memset(base, 0, total);
free(base);

2
src/etcp_bbr.c

@ -779,8 +779,8 @@ static void bbr_bound_cwnd_for_inflight_model(struct bbr* bbr, uint32_t* cwnd, u
(bbr->mode == BBR_PROBE_BW && bbr->cycle_idx == BBR_BW_PROBE_CRUISE))
cap = bbr_inflight_with_headroom(bbr);
cap = (uint32_t)(cap < bbr->inflight_lo ? cap : bbr->inflight_lo);
cap = (uint32_t)(cap > BBR_CWND_MIN_TARGET * mss ? cap : BBR_CWND_MIN_TARGET * mss);
cap = (uint32_t)(cap < bbr->inflight_lo ? cap : bbr->inflight_lo);
*cwnd = (uint32_t)(cap < *cwnd ? cap : *cwnd);
}

11
src/proxy/tcp_proxy_client.c

@ -427,7 +427,7 @@ static struct tcp_proxy_client_conn* tcp_proxy_client_find_conn(struct tcp_proxy
static void tcp_proxy_client_handle_data(struct tcp_proxy_client* p, struct ETCP_CONN* conn, uint32_t stream_id, struct ll_entry* entry) {
struct tcp_proxy_client_conn* pc = tcp_proxy_client_find_conn(p, stream_id);
if (!pc || pc->rem_closed || pc->error || pc->tun_closed) {
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "PROXY DATA sid=%08x — no conn/closed, шлём ERROR", stream_id);
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "PROXY DATA sid=%08x — no conn/closed, sending ERROR", stream_id);
tcp_proxy_client_send_msg(p->inst, p->via_node_id, TCP_PROXY_SUBCMD_ERROR, stream_id, NULL, 0);
queue_dgram_free(entry); queue_entry_free(entry); return;
}
@ -448,7 +448,7 @@ static void tcp_proxy_client_handle_close(struct tcp_proxy_client* p, uint32_t s
stream_id, pc->tun_closed);
pc->rem_closed = 1;
} else {
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "PROXY CLOSE sid=%08x — no conn, шлём ERROR", stream_id);
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "PROXY CLOSE sid=%08x — no conn, sending ERROR", stream_id);
tcp_proxy_client_send_msg(p->inst, p->via_node_id, TCP_PROXY_SUBCMD_ERROR, stream_id, NULL, 0);
}
}
@ -460,8 +460,7 @@ static void tcp_proxy_client_handle_error(struct tcp_proxy_client* p, uint32_t s
stream_id, pc->tun_closed, pc->rem_closed);
pc->error = 1;
} else {
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "PROXY ERROR sid=%08x — no conn, шлём ERROR", stream_id);
tcp_proxy_client_send_msg(p->inst, p->via_node_id, TCP_PROXY_SUBCMD_ERROR, stream_id, NULL, 0);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY ERROR sid=%08x — no conn, silent drop", stream_id);
}
}
@ -496,6 +495,10 @@ void tcp_proxy_client_etcp_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entr
if (subcmd == TCP_PROXY_SUBCMD_CLOSE) { tcp_proxy_client_handle_close(proxy, stream_id); queue_dgram_free(entry); queue_entry_free(entry); return; }
if (subcmd == TCP_PROXY_SUBCMD_ERROR) { tcp_proxy_client_handle_error(proxy, stream_id); queue_dgram_free(entry); queue_entry_free(entry); return; }
}
if (subcmd == TCP_PROXY_SUBCMD_ERROR) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "PROXY ERROR sid=%08x — conn not found, silent drop", stream_id);
queue_dgram_free(entry); queue_entry_free(entry); return;
}
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "TCP proxy client: unhandled subcmd=%02x sid=%08x", subcmd, stream_id);
queue_dgram_free(entry); queue_entry_free(entry);
}

6
src/proxy/tcp_proxy_server.c

@ -194,7 +194,7 @@ static void tcp_proxy_server_sock_read_cb(socket_t sock, void* arg) {
if (ret == 0) {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:RECV fd=%d sid=%08x len=%zd total=%d", (int)rc->sock, rc->stream_id, n, tcp_proxy_server_conn_total(rc));
} else {
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:PAUSE fd=%d sid=%08x len=%zd — backpressure, пауза", (int)rc->sock, rc->stream_id, n);
DEBUG_INFO(DEBUG_CATEGORY_SOCKET, "SOCK:PAUSE fd=%d sid=%08x len=%zd — backpressure, paused", (int)rc->sock, rc->stream_id, n);
rc->pause_buf = u_malloc(n);
if (rc->pause_buf) { memcpy(rc->pause_buf, buf, n); rc->pause_len = n; }
else DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "SOCK:PAUSE malloc(%zd) failed sid=%08x", n, rc->stream_id);
@ -330,7 +330,7 @@ int tcp_proxy_server_handle_data(struct UTUN_INSTANCE* inst, struct ETCP_CONN* c
struct tcp_proxy_server* ctx = &inst->tcp_proxy_server;
struct tcp_proxy_server_conn* rc = tcp_proxy_server_find_conn(ctx, stream_id);
if (!rc || rc->sock == SOCKET_INVALID || rc->error) {
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "TPS handle_data: нет/error conn для sid=%08x, шлём ERROR", stream_id);
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "TPS handle_data: no/error conn for sid=%08x, sending ERROR", stream_id);
if (conn) tcp_proxy_server_send_msg(inst, conn->peer_node_id, TCP_PROXY_SUBCMD_ERROR, stream_id, NULL, 0);
queue_dgram_free(entry); queue_entry_free(entry); return -1;
}
@ -370,7 +370,7 @@ void tcp_proxy_server_handle_close(struct UTUN_INSTANCE* inst, uint32_t stream_i
}
prev = &rc->next;
}
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "TCP proxy server: CLOSE sid=%08x — no conn, шлём ERROR", stream_id);
DEBUG_WARN(DEBUG_CATEGORY_SOCKET, "TCP proxy server: CLOSE sid=%08x — no conn, sending ERROR", stream_id);
tcp_proxy_server_send_msg(inst, inst->node_id, TCP_PROXY_SUBCMD_ERROR, stream_id, NULL, 0);
}

48
tests/Makefile.am

@ -2,45 +2,45 @@
# All available tests (check_PROGRAMS runs via automake check-TESTS)
check_PROGRAMS = \
test_bbr_integration \
test_etcp_bbr \
test_ll_queue \
test_serialize \
test_packet_dump \
test_debug_categories \
test_config_debug \
test_memory_pool_and_config \
test_etcp_api \
test_etcp_crypto \
test_pkt_normalizer_standalone \
test_route_lib \
test_radix \
test_route6_lib \
test_etcp_router_unit \
test_etcp_minimal \
test_ipv6_sockets \
test_u_async_timeouts \
test_u_async_comprehensive \
test_etcp_two_instances \
test_etcp_simple_traffic \
test_ipv6_sockets \
test_etcp_minimal \
test_etcp_100_packets \
test_etcp_reconnect \
test_pkt_normalizer_etcp \
test_pkt_normalizer_standalone \
test_etcp_api \
test_ll_queue \
test_serialize \
test_intensive_memory_pool \
test_memory_pool_and_config \
test_packet_dump \
test_u_async_comprehensive \
test_lwip_tcp \
test_u_async_performance \
test_u_async_timeouts \
test_debug_categories \
test_config_debug \
test_route_lib \
test_bgp_route_exchange \
test_etcp_router \
test_etcp_bbr \
test_etcp_ping \
test_route_ping \
test_nat_detection \
test_nat_engine \
test_nat_transport \
test_nat_stress \
test_tcp_proxy_client \
test_lwip_tcp \
test_etcp_router \
test_etcp_router_unit \
test_tcp_proxy_server \
test_udp_proxy \
test_icmp_proxy \
test_radix \
test_route6_lib \
test_tcp_proxy_client \
test_bgp_route_exchange \
test_bbr_integration \
test_intensive_memory_pool \
bench_timeout_heap \
bench_uasync_timeouts
@ -385,7 +385,7 @@ TEST_LOG_DIR = $(top_builddir)/tests/logs
# Custom check target that runs tests and saves logs to files
check-local: $(check_PROGRAMS)
@$(MKDIR_P) $(TEST_LOG_DIR); \
@rm -f $(TEST_LOG_DIR)/*.log 2>/dev/null || true; $(MKDIR_P) $(TEST_LOG_DIR); \
G='\033[0;32m'; R='\033[0;31m'; Y='\033[0;33m'; N='\033[0m'; \
passed=0; \
failed=0; \

23
tests/bbr_integration/test_bbr_integration.c

@ -41,10 +41,11 @@
#define CLI_PORT 21000
#define PAYLOAD_SIZE 1200
#define TEST_DURATION_MS 15000
#define TEST_DURATION_MS 5000
#define METRICS_TB 1000 /* 100ms in 0.1ms units */
#define SEND_TIMER_TB 1 /* 0.1ms re-schedule */
#define BURST_MAX 64
#define SEND_QUEUE_THRESHOLD 1 /* pause when normalizer->input >= 1 */
#define EMU_DELAY_MS 50
#define EMU_JITTER_MS 12
@ -92,6 +93,7 @@ struct test_ctx {
FILE* log_file;
void* metrics_timer;
struct queue_waiter_handle waiter;
};
/* ===== Time helper ===== */
@ -175,11 +177,25 @@ static void send_timer_cb(void* arg) {
send_burst((struct test_ctx*)arg);
}
static void send_waiter_cb(struct ll_queue* q, void* arg) {
(void)q;
struct test_ctx* ctx = (struct test_ctx*)arg;
uasync_call_soon(ctx->ua, ctx, send_timer_cb);
}
static void send_burst(struct test_ctx* ctx) {
if (ctx->test_done) return;
struct ETCP_CONN* conn = ctx->sender->connections;
if (!conn || !conn->initialized) return;
if (conn->normalizer && conn->normalizer->input) {
if (queue_entry_count(conn->normalizer->input) >= SEND_QUEUE_THRESHOLD) {
queue_waiter_wait(conn->normalizer->input, &ctx->waiter,
send_waiter_cb, ctx);
return;
}
}
int sent = 0;
while (sent < BURST_MAX) {
struct ll_entry* e = ll_alloc_lldgram(1 + PAYLOAD_SIZE);
@ -375,6 +391,9 @@ int main(void) {
t0 = now_us();
while ((now_us() - t0) < 500000ULL) uasync_poll(ctx.ua, 1);
if (ctx.sender->connections && ctx.sender->connections->normalizer && ctx.sender->connections->normalizer->input)
queue_set_threshold(ctx.sender->connections->normalizer->input, 0, 0);
/* Start test */
printf("\n=== Starting traffic (%d seconds) ===\n\n", TEST_DURATION_MS / 1000);
ctx.log_file = fopen("test_bbr_integration.log", "w");
@ -487,6 +506,8 @@ int main(void) {
}
/* Cleanup */
if (ctx.sender->connections && ctx.sender->connections->normalizer && ctx.sender->connections->normalizer->input)
queue_waiter_cancel(ctx.sender->connections->normalizer->input, &ctx.waiter);
dummynet_destroy(ctx.dn);
utun_instance_destroy(ctx.sender);
utun_instance_destroy(ctx.receiver);

2
tests/test_etcp_bbr.c

@ -142,7 +142,7 @@ static int test_probe_rtt(void)
s.now_tb = 10000 + 5000U * 10; // +5 seconds in 0.1ms
run_ack(&s, 1400, 1000, 0, 0, 5600);
if (s.mode != BBR_PROBE_RTT) FAIL("mode not PROBE_RTT");
if (g_cwnd > BBR_CWND_MIN_TARGET + 1400) FAIL("cwnd not capped");
if (g_cwnd > BBR_CWND_MIN_TARGET * 1400 + 1400) FAIL("cwnd not capped");
PASS(); return 0;
}

Loading…
Cancel
Save