Browse Source

refactor: remove tcp_conn_push_write, use direct queue put with chunking

External code pushes to write_queue directly (entry_pool + data_pool).
Added queue_set_threshold(write_queue, 32, 0) for write backpressure.
Tests split data into write_chunk_size chunks before queue put.
etcp-inflight-fix
Evgeny 4 months ago
parent
commit
04f6b3e040
  1. 31
      lib/tcp_io.c
  2. 7
      lib/tcp_io.h
  3. 14
      src/proxy/tcp_proxy_server.c
  4. 24
      tests/test_tcp_io.c
  5. 12
      tests/test_tcp_proxy_server.c

31
lib/tcp_io.c

@ -68,6 +68,7 @@ struct tcp_conn* tcp_conn_create(
u_free(tc); return NULL; u_free(tc); return NULL;
} }
queue_set_threshold(tc->read_queue, read_low_water, 0); queue_set_threshold(tc->read_queue, read_low_water, 0);
queue_set_threshold(tc->write_queue, 32, 0);
queue_set_callback(tc->write_queue, write_queue_fetch_cb, tc); queue_set_callback(tc->write_queue, write_queue_fetch_cb, tc);
queue_set_waiter_defer(tc->write_queue, 1); queue_set_waiter_defer(tc->write_queue, 1);
@ -326,33 +327,3 @@ void tcp_conn_set_flushed(struct tcp_conn* tc, void (*on_flushed)(struct tcp_con
if (!tc) return; if (!tc) return;
tc->on_flushed = on_flushed; tc->on_flushed = on_flushed;
} }
// ====================================================================
// Отправка данных в сокет (внешний интерфейс)
// ====================================================================
int tcp_conn_push_write(struct tcp_conn* tc, const uint8_t* data, size_t len) {
if (!tc || !data || len == 0) return -1;
if (tc->error || tc->fin) return -1;
size_t offset = 0;
while (offset < len) {
size_t chunk = len - offset;
if (chunk > tc->write_chunk_size) chunk = tc->write_chunk_size;
struct ll_entry* e = queue_entry_new_from_pool(tc->entry_pool);
uint8_t* buf = memory_pool_alloc(tc->data_pool);
if (!e || !buf) {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_io: push_write alloc failed fd=%d chunk=%zu", (int)tc->sock, chunk);
if (e) queue_entry_free(e);
if (buf) memory_pool_free(tc->data_pool, buf);
return -1;
}
memcpy(buf, data + offset, chunk);
e->dgram = buf;
e->len = (uint16_t)chunk;
queue_data_put(tc->write_queue, e);
offset += chunk;
}
return 0;
}

7
lib/tcp_io.h

@ -85,9 +85,10 @@ struct tcp_conn* tcp_conn_create(
void tcp_conn_destroy(struct tcp_conn* tc); void tcp_conn_destroy(struct tcp_conn* tc);
// Отправка данных в сокет: режет на чанки по write_chunk_size, кладёт в write_queue. // Внешний код пишет данные в tc->write_queue напрямую (queue_data_put).
// Автозабор write_queue сам отправит когда сокет будет готов. // Автозабор write_queue (deferred) сам отправляет когда сокет готов.
int tcp_conn_push_write(struct tcp_conn* tc, const uint8_t* data, size_t len); // Перед push проверять порог: queue_set_threshold в tcp_conn_create (32 entries).
// При заполнении — queue_waiter_wait на освобождение.
// Одноразовый коллбэк: вызывается когда write_queue + write_buf полностью опустели. // Одноразовый коллбэк: вызывается когда write_queue + write_buf полностью опустели.
// После вызова сбрасывается. Установить повторно можно в любой момент. // После вызова сбрасывается. Установить повторно можно в любой момент.

14
src/proxy/tcp_proxy_server.c

@ -238,7 +238,19 @@ int tcp_proxy_server_handle_data(struct UTUN_INSTANCE* inst, struct ETCP_CONN* c
queue_dgram_free(entry); queue_entry_free(entry); return -1; queue_dgram_free(entry); queue_entry_free(entry); return -1;
} }
size_t data_len = entry->len - TCP_PROXY_HDR_SIZE; size_t data_len = entry->len - TCP_PROXY_HDR_SIZE;
if (data_len > 0) tcp_conn_push_write(rc->tc, entry->dgram + TCP_PROXY_HDR_SIZE, data_len); if (data_len > 0) {
struct ll_entry* e = queue_entry_new_from_pool(rc->tc->entry_pool);
uint8_t* buf = memory_pool_alloc(rc->tc->data_pool);
if (e && buf) {
memcpy(buf, entry->dgram + TCP_PROXY_HDR_SIZE, data_len);
e->dgram = buf; e->len = (uint16_t)data_len;
queue_data_put(rc->tc->write_queue, e);
} else {
DEBUG_ERROR(DEBUG_CATEGORY_SOCKET, "tcp_proxy_server: handle_data alloc failed sid=%08x", stream_id);
if (e) queue_entry_free(e);
if (buf) memory_pool_free(rc->tc->data_pool, buf);
}
}
queue_dgram_free(entry); queue_entry_free(entry); queue_dgram_free(entry); queue_entry_free(entry);
return 0; return 0;
} }

24
tests/test_tcp_io.c

@ -48,6 +48,22 @@ static void reset_counters(void) {
g_last_fin_tc = NULL; g_last_fin_tc = NULL;
} }
static int push_write(struct tcp_conn* tc, const uint8_t* data, size_t len) {
size_t offset = 0;
while (offset < len) {
size_t chunk = len - offset;
if (chunk > tc->write_chunk_size) chunk = tc->write_chunk_size;
struct ll_entry* e = queue_entry_new_from_pool(tc->entry_pool);
uint8_t* buf = memory_pool_alloc(tc->data_pool);
if (!e || !buf) { if (e) queue_entry_free(e); if (buf) memory_pool_free(tc->data_pool, buf); return -1; }
memcpy(buf, data + offset, chunk);
e->dgram = buf; e->len = (uint16_t)chunk;
queue_data_put(tc->write_queue, e);
offset += chunk;
}
return 0;
}
static void test_basic_send_recv(void) { static void test_basic_send_recv(void) {
TEST_START("Basic send and recv via tcp_conn"); TEST_START("Basic send and recv via tcp_conn");
@ -65,8 +81,8 @@ static void test_basic_send_recv(void) {
// Отправляем данные через tcp_conn // Отправляем данные через tcp_conn
uint8_t send_buf[3000]; uint8_t send_buf[3000];
memset(send_buf, 'A', sizeof(send_buf)); memset(send_buf, 'A', sizeof(send_buf));
int ret = tcp_conn_push_write(tc, send_buf, sizeof(send_buf)); int ret = push_write(tc, send_buf, sizeof(send_buf));
ASSERT_EQ(ret, 0, "tcp_conn_push_write first call failed"); ASSERT_EQ(ret, 0, "push_write first call failed");
uasync_poll(ua, 1); // write_cb должен отправить uasync_poll(ua, 1); // write_cb должен отправить
// Читаем с другой стороны // Читаем с другой стороны
@ -126,8 +142,8 @@ static void test_partial_write(void) {
ASSERT_TRUE(big_buf != NULL, "malloc failed"); ASSERT_TRUE(big_buf != NULL, "malloc failed");
memset(big_buf, 'X', 200000); memset(big_buf, 'X', 200000);
int ret = tcp_conn_push_write(tc, big_buf, 200000); int ret = push_write(tc, big_buf, 200000);
ASSERT_EQ(ret, 0, "tcp_conn_push_write large failed"); ASSERT_EQ(ret, 0, "push_write large failed");
uasync_poll(ua, 1); // write_cb отправляет что может uasync_poll(ua, 1); // write_cb отправляет что может
// Дрейним sv[1] и проверяем что все данные приходят // Дрейним sv[1] и проверяем что все данные приходят

12
tests/test_tcp_proxy_server.c

@ -43,6 +43,16 @@ done: close(cli); close(srv);
static void on_fin_cb(struct tcp_conn* tc, void* arg) { (void)tc; (void)arg; } static void on_fin_cb(struct tcp_conn* tc, void* arg) { (void)tc; (void)arg; }
static void on_error_cb(struct tcp_conn* tc, int err, void* arg) { (void)tc; (void)arg; (void)err; } static void on_error_cb(struct tcp_conn* tc, int err, void* arg) { (void)tc; (void)arg; (void)err; }
static int push_write(struct tcp_conn* tc, const uint8_t* data, size_t len) {
struct ll_entry* e = queue_entry_new_from_pool(tc->entry_pool);
uint8_t* buf = memory_pool_alloc(tc->data_pool);
if (!e || !buf) { if (e) queue_entry_free(e); if (buf) memory_pool_free(tc->data_pool, buf); return -1; }
memcpy(buf, data, len);
e->dgram = buf; e->len = (uint16_t)len;
queue_data_put(tc->write_queue, e);
return 0;
}
static void timeout_cb(void* arg) { static void timeout_cb(void* arg) {
(void)arg; (void)arg;
printf("[FAIL] timeout\n"); printf("[FAIL] timeout\n");
@ -84,7 +94,7 @@ int main(void) {
if (phase == 0 && tc->connected) { if (phase == 0 && tc->connected) {
phase = 1; phase = 1;
for (int i = 0; i < PAYLOAD_SIZE; i++) send_buf[i] = (uint8_t)(rand() & 0xFF); for (int i = 0; i < PAYLOAD_SIZE; i++) send_buf[i] = (uint8_t)(rand() & 0xFF);
tcp_conn_push_write(tc, send_buf, PAYLOAD_SIZE); push_write(tc, send_buf, PAYLOAD_SIZE);
} }
if (phase == 1 && !tc->error) { if (phase == 1 && !tc->error) {

Loading…
Cancel
Save