diff --git a/src/pkt_normalizer.c b/src/pkt_normalizer.c index c3a4b1bd..9effe870 100644 --- a/src/pkt_normalizer.c +++ b/src/pkt_normalizer.c @@ -109,13 +109,15 @@ void pn_unpacker_reset_state(struct PKTNORM* pn) { void pn_packer_send(struct PKTNORM* pn, uint8_t* data, uint16_t len) { if (!pn || !data || len == 0) return; - + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_packer_send: pn=%p, len=%d", pn, len); + struct ll_entry* entry = ll_alloc_lldgram(len); memcpy(entry->dgram, data, len); entry->len = len; entry->dgram_pool = NULL; - queue_data_put(pn->input, entry, 0); + int ret = queue_data_put(pn->input, entry, 0); + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_packer_send: queue_data_put returned %d, input count=%d", ret, queue_entry_count(pn->input)); // Cancel flush timer if active if (pn->flush_timer) { @@ -128,6 +130,7 @@ void pn_packer_send(struct PKTNORM* pn, uint8_t* data, uint16_t len) { static void packer_cb(struct ll_queue* q, void* arg) { struct PKTNORM* pn = (struct PKTNORM*)arg; if (!pn) return; + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_packer: packer_cb"); queue_wait_threshold(pn->etcp->input_queue, 0, 0, etcp_input_ready_cb, pn); } @@ -136,6 +139,7 @@ static void packer_cb(struct ll_queue* q, void* arg) { static void pn_send_to_etcp(struct PKTNORM* pn) { if (!pn || !pn->sndpart.data || !pn->sndpart.len==0) return; + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_packer: pn_send_to_etcp"); // Allocate ETCP_FRAGMENT from rx_pool struct ETCP_FRAGMENT* frag = queue_entry_new_from_pool(pn->etcp->rx_pool); if (!frag) {// drop data @@ -151,6 +155,10 @@ static void pn_send_to_etcp(struct PKTNORM* pn) { queue_data_put(pn->etcp->input_queue, frag, 0); queue_entry_free(&pn->sndpart); + // Сбросить структуру после освобождения + pn->sndpart.dgram = NULL; + pn->sndpart.len = 0; + pn->sndpart.memlen = 0; } // Internal: Renew sndpart buffer @@ -172,6 +180,8 @@ static void etcp_input_ready_cb(struct ll_queue* q, void* arg) { struct PKTNORM* pn = (struct PKTNORM*)arg; if (!pn) return; + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_packer: etcp_input_ready_cb"); + void* data = queue_data_get(pn->input); if (!data) { queue_resume_callback(pn->input); @@ -187,8 +197,9 @@ static void etcp_input_ready_cb(struct ll_queue* q, void* arg) { int remain = pn->frag_size - pn->sndpart.len; if (remain < 3) { - DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_normalizer: part size error, remain=%d", remain); - break; + // Буфер почти полон, отправить его + pn_send_to_etcp(pn); + continue; // Продолжить с новым буфером } if (ptr == 0) { @@ -204,7 +215,10 @@ static void etcp_input_ready_cb(struct ll_queue* q, void* arg) { pn->sndpart.len += n; ptr += n; - pn_buf_renew(pn); + // Проверить, не заполнился ли буфер + if (pn->sndpart.len >= pn->frag_size - 2) { + pn_send_to_etcp(pn); + } } queue_dgram_free(in_dgram); diff --git a/tests/test_etcp_100_packets b/tests/test_etcp_100_packets index 4e0b4a6c..0473c106 100755 Binary files a/tests/test_etcp_100_packets and b/tests/test_etcp_100_packets differ diff --git a/tests/test_etcp_minimal b/tests/test_etcp_minimal index 6e3df713..e739de7b 100755 Binary files a/tests/test_etcp_minimal and b/tests/test_etcp_minimal differ diff --git a/tests/test_etcp_simple_traffic b/tests/test_etcp_simple_traffic index a03b9bf0..86046e58 100755 Binary files a/tests/test_etcp_simple_traffic and b/tests/test_etcp_simple_traffic differ diff --git a/tests/test_etcp_two_instances b/tests/test_etcp_two_instances index fdbd05c5..67208380 100755 Binary files a/tests/test_etcp_two_instances and b/tests/test_etcp_two_instances differ diff --git a/tests/test_pkt_normalizer_etcp b/tests/test_pkt_normalizer_etcp index e1b90923..53c5a1a5 100755 Binary files a/tests/test_pkt_normalizer_etcp and b/tests/test_pkt_normalizer_etcp differ diff --git a/tests/test_pkt_normalizer_etcp.c b/tests/test_pkt_normalizer_etcp.c index 311c6510..66ca38b9 100644 --- a/tests/test_pkt_normalizer_etcp.c +++ b/tests/test_pkt_normalizer_etcp.c @@ -96,13 +96,23 @@ static int is_connection_established(struct UTUN_INSTANCE* inst) { // Send packets from client to server (forward direction) via normalizer static void send_packets_fwd(void) { - if (!client_instance || !client_pn || packets_sent_fwd >= TOTAL_PACKETS) return; + printf("DEBUG send_packets_fwd: client_instance=%p, client_pn=%p, packets_sent_fwd=%d/%d\n", + (void*)client_instance, (void*)client_pn, packets_sent_fwd, TOTAL_PACKETS); + fflush(stdout); + + if (!client_instance || !client_pn || packets_sent_fwd >= TOTAL_PACKETS) { + printf("DEBUG send_packets_fwd: early return (instance=%p, pn=%p, sent=%d)\n", + (void*)client_instance, (void*)client_pn, packets_sent_fwd); + fflush(stdout); + return; + } // Start timing on first packet if (packets_sent_fwd == 0) { clock_gettime(CLOCK_MONOTONIC, &start_time_fwd); phase = 1; printf("Starting forward transfer (client -> server) via normalizer...\n"); + fflush(stdout); } // Send while we have packets @@ -227,29 +237,60 @@ static void monitor_and_send(void* arg) { static int connection_checked = 0; if (!connection_checked) { - if (is_connection_established(client_instance)) { + int client_established = is_connection_established(client_instance); + int server_established = is_connection_established(server_instance); + printf("Connection check: client=%d, server=%d\n", client_established, server_established); + printf("Connections: client=%p, server=%p\n", + (void*)client_instance->connections, (void*)server_instance->connections); + if (client_established) { DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Connection established, starting transmission via normalizer"); connection_checked = 1; // Initialize normalizers after connection is established + printf("DEBUG: client_pn=%p, client_instance->connections=%p\n", + (void*)client_pn, (void*)client_instance->connections); + fflush(stdout); if (!client_pn && client_instance->connections) { + printf("Creating client normalizer...\n"); + fflush(stdout); client_pn = pn_init(client_instance->connections); if (!client_pn) { - printf("Failed to create client normalizer\n"); + printf("Failed to create client normalizer (returned NULL)\n"); + fflush(stdout); test_completed = 2; return; } printf("Client normalizer created (frag_size=%d)\n", client_pn->frag_size); + fflush(stdout); + } else if (!client_pn) { + printf("No client connections available\n"); + fflush(stdout); + } else { + printf("Client normalizer already exists\n"); + fflush(stdout); } + printf("DEBUG: server_pn=%p, server_instance->connections=%p\n", + (void*)server_pn, (void*)server_instance->connections); + fflush(stdout); if (!server_pn && server_instance->connections) { + printf("Creating server normalizer...\n"); + fflush(stdout); server_pn = pn_init(server_instance->connections); if (!server_pn) { - printf("Failed to create server normalizer\n"); + printf("Failed to create server normalizer (returned NULL)\n"); + fflush(stdout); test_completed = 2; return; } printf("Server normalizer created (frag_size=%d)\n\n", server_pn->frag_size); + fflush(stdout); + } else if (!server_pn) { + printf("No server connections available\n"); + fflush(stdout); + } else { + printf("Server normalizer already exists\n"); + fflush(stdout); } } }