You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
412 lines
14 KiB
412 lines
14 KiB
// pkt_normalizer.c - Implementation of packet normalizer for ETCP |
|
#include "pkt_normalizer.h" |
|
#include "etcp.h" // For ETCP_CONN and related structures |
|
#include "etcp_api.h" // For etcp_recv callback |
|
#include "routing.h" // For routing_add_conn/routing_del_conn |
|
#include "utun_instance.h" // For UTUN_INSTANCE |
|
#include "ll_queue.h" // For queue operations |
|
#include "u_async.h" // For UASYNC |
|
#include <stdlib.h> |
|
#include <string.h> |
|
#include <stdio.h> // For debugging (can be removed if not needed) |
|
#include "../lib/debug_config.h" // For DEBUG macros |
|
#include "../lib/mem.h" |
|
|
|
// Forward declarations |
|
static void packer_cb(struct ll_queue* q, void* arg); |
|
static void pn_flush_cb(void* arg); |
|
static void etcp_input_ready_cb(struct ll_queue* q, void* arg); |
|
static void pn_unpacker_cb(struct ll_queue* q, void* arg); |
|
static void pn_send_to_etcp(struct PKTNORM* pn); |
|
|
|
// Initialization |
|
struct PKTNORM* pn_init(struct ETCP_CONN* etcp) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, ""); |
|
if (!etcp) return NULL; |
|
|
|
if (etcp->mtu<200) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_NORMALIZER, "pn_init: ETCP MTU error=%d", etcp->mtu); |
|
} |
|
|
|
struct PKTNORM* pn = u_calloc(1, sizeof(struct PKTNORM)); |
|
if (!pn) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_NORMALIZER, "pn_init: calloc failed"); |
|
return NULL; |
|
} |
|
|
|
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "pn_init:[%s] init %p, mtu=%d", etcp->log_name, etcp, etcp->mtu); |
|
|
|
|
|
pn->etcp = etcp; |
|
pn->ua = etcp->instance->ua; |
|
pn->frag_size = etcp->mtu - ACK_REZERV - UDP_HDR_SIZE - UDP_SC_HDR_SIZE; |
|
int min_frag = UDP_HDR_SIZE + UDP_SC_HDR_SIZE + ACK_REZERV; |
|
if (etcp->mtu < min_frag || pn->frag_size > etcp->mtu) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_NORMALIZER, "pn_init: MTU %d too small (min %d)", etcp->mtu, min_frag); |
|
pn->frag_size = etcp->mtu; |
|
} |
|
pn->tx_wait_time = 10; |
|
|
|
pn->input = queue_new(pn->ua, 0, 0, 0, "pn_input"); // No hash needed |
|
if (!pn->input) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_NORMALIZER, "pn_init: queue_new(input) failed"); |
|
pn_deinit(pn); |
|
return NULL; |
|
} |
|
|
|
pn->output = queue_new(pn->ua, 0, 0, 0, "pn_output"); // No hash needed |
|
if (!pn->output) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_NORMALIZER, "pn_init: queue_new(output) failed"); |
|
pn_deinit(pn); |
|
return NULL; |
|
} |
|
|
|
queue_set_callback(pn->input, packer_cb, pn); |
|
queue_set_callback(etcp->output_queue, pn_unpacker_cb, pn); |
|
queue_set_callback(pn->output, etcp_int_recv, etcp); |
|
|
|
pn->data = NULL; |
|
pn->recvpart = NULL; |
|
pn->flush_timer = NULL; |
|
|
|
pn->input_waiter_handle.internal = NULL; |
|
|
|
return pn; |
|
} |
|
|
|
// Deinitialization |
|
void pn_deinit(struct PKTNORM* pn) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, ""); |
|
if (!pn) return; |
|
|
|
// Unregister from routing module |
|
if (pn->etcp) { |
|
routing_del_conn(pn->etcp); |
|
} |
|
|
|
// Drain and free queues |
|
if (pn->input) { |
|
struct ll_entry* entry; |
|
while ((entry = queue_data_get(pn->input)) != NULL) { |
|
if (entry->dgram) { |
|
u_free(entry->dgram); |
|
} |
|
queue_entry_free(entry); |
|
} |
|
queue_free(pn->input); |
|
} |
|
if (pn->output) { |
|
struct ll_entry* entry; |
|
while ((entry = queue_data_get(pn->output)) != NULL) { |
|
if (entry->dgram) { |
|
u_free(entry->dgram); |
|
} |
|
queue_entry_free(entry); |
|
} |
|
queue_free(pn->output); |
|
} |
|
|
|
if (pn->flush_timer) { |
|
uasync_call_soon_cancel(pn->ua, pn->flush_timer); |
|
} |
|
|
|
if (pn->data) { |
|
memory_pool_free(pn->etcp->instance->data_pool, pn->data); |
|
} |
|
if (pn->recvpart) { |
|
queue_dgram_free(pn->recvpart); |
|
queue_entry_free(pn->recvpart); |
|
} |
|
|
|
// Cancel waiter and clear callbacks on ETCP queues to prevent use-after-free |
|
// during etcp_connection_close() when drain_and_free_fragment_queue() is called |
|
if (pn->etcp) { |
|
if (pn->etcp->input_queue) { |
|
queue_waiter_cancel(pn->etcp->input_queue, &pn->input_waiter_handle); |
|
} |
|
if (pn->etcp->output_queue) { |
|
queue_set_callback(pn->etcp->output_queue, NULL, NULL); |
|
} |
|
} |
|
|
|
u_free(pn); |
|
} |
|
|
|
// Reset unpacker state |
|
void pn_unpacker_reset_state(struct PKTNORM* pn) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, ""); |
|
if (!pn) return; |
|
|
|
// Reset recvpart |
|
if (pn->recvpart) { |
|
queue_dgram_free(pn->recvpart); |
|
queue_entry_free(pn->recvpart); |
|
pn->recvpart = NULL; |
|
} |
|
} |
|
|
|
// Reset packer and unpacker state (for reconnection) |
|
void pn_reset(struct PKTNORM* pn) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, ""); |
|
if (!pn) return; |
|
|
|
// Cancel flush timer |
|
if (pn->flush_timer) { |
|
uasync_call_soon_cancel(pn->ua, pn->flush_timer); |
|
pn->flush_timer = NULL; |
|
} |
|
|
|
// Reset packer state |
|
if (pn->data) { |
|
memory_pool_free(pn->etcp->instance->data_pool, pn->data); |
|
pn->data = NULL; |
|
} |
|
pn->data_ptr = 0; |
|
pn->data_size = 0; |
|
|
|
// Free pending if any |
|
if (pn->pending) { |
|
queue_dgram_free(pn->pending); |
|
queue_entry_free(pn->pending); |
|
pn->pending = NULL; |
|
} |
|
pn->pending_in_ptr = 0; |
|
|
|
// Reset unpacker state |
|
pn_unpacker_reset_state(pn); |
|
|
|
queue_resume_callback(pn->input); |
|
queue_resume_callback(pn->output); |
|
} |
|
|
|
// Send data to packer (copies and adds to input queue or pending, triggering callback) используется только в юниттесте |
|
void pn_packer_send(struct PKTNORM* pn, uint8_t* data, uint16_t len) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, ""); |
|
if (!pn || !data || len == 0) return; |
|
|
|
struct ll_entry* entry = ll_alloc_lldgram(len); |
|
if (!entry) { |
|
pn->alloc_errors++; |
|
DEBUG_ERROR(DEBUG_CATEGORY_NORMALIZER, "pn_packer_send: ll_alloc_lldgram failed"); |
|
return; |
|
} |
|
memcpy(entry->dgram, data, len); |
|
entry->len = len; |
|
entry->dgram_pool = NULL; |
|
|
|
// Cancel flush timer if active |
|
if (pn->flush_timer) { |
|
uasync_call_soon_cancel(pn->ua, pn->flush_timer); |
|
pn->flush_timer = NULL; |
|
} |
|
|
|
DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, "PUT to input"); |
|
int ret = queue_data_put(pn->input, entry); |
|
pn->in_total_pkts++; |
|
pn->in_total_bytes += len; |
|
DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, "PUT to input end"); |
|
// DEBUG_INFO(DEBUG_CATEGORY_NORMALIZER, "pn_packer_send: queue_data_put returned %d, input count=%d", ret, queue_entry_count(pn->input)); |
|
} |
|
|
|
|
|
// Internal: Packer callback |
|
static void packer_cb(struct ll_queue* q, void* arg) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, ""); |
|
struct PKTNORM* pn = (struct PKTNORM*)arg; |
|
if (!pn) return; |
|
DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, "input_q->pn: waiting etcp input threshold"); |
|
queue_waiter_wait(pn->etcp->input_queue, &pn->input_waiter_handle, etcp_input_ready_cb, pn); |
|
} |
|
|
|
// Helper to send block to ETCP as ETCP_FRAGMENT |
|
static void pn_send_to_etcp(struct PKTNORM* pn) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, ""); |
|
if (!pn || !pn->data || pn->data_ptr == 0) return; |
|
|
|
// DEBUG_INFO(DEBUG_CATEGORY_NORMALIZER, "pn_packer: pn_send_to_etcp"); |
|
// Allocate ETCP_FRAGMENT from io_pool |
|
struct ETCP_FRAGMENT* frag = (struct ETCP_FRAGMENT*)queue_entry_new_from_pool(pn->etcp->io_pool); |
|
if (!frag) {// drop data |
|
DEBUG_ERROR(DEBUG_CATEGORY_NORMALIZER, "pn_packer: send to etcp alloc error"); |
|
pn->alloc_errors++; |
|
pn->data_ptr = 0; |
|
return; |
|
} |
|
|
|
frag->seq = 0; |
|
frag->timestamp = 0; |
|
frag->ll.dgram = pn->data; |
|
frag->ll.len = pn->data_ptr; |
|
frag->ll.dgram_pool = pn->etcp->instance->data_pool; |
|
frag->ll.memlen = pn->etcp->instance->data_pool->object_size; |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_NORMALIZER, "pn->etcp: size=%d memlen=%d frag_size=%d", frag->ll.len, frag->ll.memlen, pn->frag_size); |
|
if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP)) log_dump(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP, "NORM->ETCP", pn->data, frag->ll.len); |
|
|
|
queue_data_put(pn->etcp->input_queue, (struct ll_entry*)frag); |
|
// Сбросить структуру (dgram передан во фрагмент, не освобождаем) |
|
pn->data = NULL; |
|
pn->data_ptr = 0; |
|
} |
|
|
|
// Internal: Renew sndpart buffer |
|
static void pn_buf_renew(struct PKTNORM* pn) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, ""); |
|
if (pn->data) { |
|
int remain = pn->data_size - pn->data_ptr; |
|
if (remain < 3) pn_send_to_etcp(pn); |
|
} |
|
if (!pn->data) { |
|
pn->data = memory_pool_alloc(pn->etcp->instance->data_pool); |
|
if (!pn->data) { |
|
pn->alloc_errors++; |
|
DEBUG_ERROR(DEBUG_CATEGORY_NORMALIZER, "pn_buf_renew: memory_pool_alloc failed"); |
|
return; |
|
} |
|
int size=pn->etcp->instance->data_pool->object_size; |
|
if (size>pn->frag_size) size=pn->frag_size; |
|
pn->data_size = size; |
|
pn->data_ptr=0; |
|
// DEBUG_INFO(DEBUG_CATEGORY_NORMALIZER, "pn_packer: new bufer size=%d bytes",size); |
|
} |
|
} |
|
|
|
// Internal: Process input when etcp->input_queue is ready (empty) |
|
static void etcp_input_ready_cb(struct ll_queue* q, void* arg) { |
|
struct PKTNORM* pn = (struct PKTNORM*)arg; |
|
if (!pn || !pn->etcp) return; |
|
// if (pn->etcp->initialized==0) return;// начинаем обработку очередей только когда готов etcp |
|
|
|
DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, "before get"); |
|
struct ll_entry* in_dgram = queue_data_get(pn->input); |
|
DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, "after get"); |
|
if (!in_dgram) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_NORMALIZER, "Q EMPTY"); |
|
queue_resume_callback(pn->input); |
|
return; |
|
} |
|
|
|
if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP)) log_dump(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP, "->NORM", in_dgram->dgram, in_dgram->len); |
|
|
|
pn_buf_renew(pn); |
|
// DEBUG_INFO(DEBUG_CATEGORY_NORMALIZER, "pn_packer: new pkt hdrpos=%d",pn->data_ptr); |
|
if (!pn->data) goto exit; // Allocation failed |
|
pn->data[pn->data_ptr++] = in_dgram->len & 0xFF; |
|
pn->data[pn->data_ptr++] = (in_dgram->len >> 8) & 0xFF; |
|
|
|
uint16_t in_ptr = 0; |
|
while (in_ptr < in_dgram->len) { |
|
int remain = pn->data_size - pn->data_ptr; |
|
int avail = in_dgram->len - in_ptr; |
|
if (avail < remain) remain = avail; |
|
// DEBUG_INFO(DEBUG_CATEGORY_NORMALIZER, "pn_packer: copy %d bytes (in_ptr=%d, out_ptr=%d)",remain, in_ptr, pn->data_ptr); |
|
memcpy(pn->data + pn->data_ptr, in_dgram->dgram + in_ptr, remain); |
|
pn->data_ptr += remain; |
|
in_ptr += remain; |
|
pn_buf_renew(pn); |
|
if (!pn->data) goto exit; |
|
} |
|
|
|
exit: |
|
queue_dgram_free(in_dgram); |
|
queue_entry_free(in_dgram); |
|
|
|
// Cancel flush timer if active |
|
if (pn->flush_timer) { |
|
uasync_call_soon_cancel(pn->ua, pn->flush_timer); |
|
pn->flush_timer = NULL; |
|
} |
|
|
|
// Set flush timer if no more input |
|
if (queue_entry_count(pn->input) == 0) { |
|
pn->flush_timer = uasync_call_soon(pn->ua/*, pn->tx_wait_time*/, pn, pn_flush_cb); |
|
} |
|
|
|
queue_resume_callback(pn->input); |
|
} |
|
|
|
// Internal: Flush callback on timeout |
|
static void pn_flush_cb(void* arg) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, ""); |
|
struct PKTNORM* pn = (struct PKTNORM*)arg; |
|
if (!pn) return; |
|
|
|
pn->flush_timer = NULL; |
|
pn_send_to_etcp(pn); |
|
} |
|
|
|
|
|
// Internal: Unpacker callback (assembles fragments into original packets) |
|
static void pn_unpacker_cb(struct ll_queue* q, void* arg) { |
|
DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, ""); |
|
struct PKTNORM* pn = (struct PKTNORM*)arg; |
|
if (!pn) return; |
|
|
|
while (1) { |
|
void* data = queue_data_get(pn->etcp->output_queue); |
|
if (!data) break; |
|
|
|
struct ETCP_FRAGMENT* frag = (struct ETCP_FRAGMENT*)data; // Since data is ll_entry* |
|
uint8_t* payload = frag->ll.dgram; |
|
uint16_t len = frag->ll.len; |
|
uint16_t ptr = 0; |
|
if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP)) log_dump(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP, "ETCP->NORM", payload, len); |
|
|
|
DEBUG_DEBUG(DEBUG_CATEGORY_NORMALIZER, "unpacking fragment len=%d", len); |
|
|
|
while (ptr < len) { |
|
if (!pn->recvpart) { |
|
// Need length header for new packet |
|
if (len - ptr < 2) { |
|
pn->logic_errors++; |
|
// Incomplete header, reset |
|
DEBUG_ERROR(DEBUG_CATEGORY_NORMALIZER, "reset state"); |
|
// pn_unpacker_reset_state(pn); |
|
etcp_conn_reinit(pn->etcp); |
|
break; |
|
} |
|
uint16_t part_size = payload[ptr] | (payload[ptr + 1] << 8); |
|
// DEBUG_INFO(DEBUG_CATEGORY_NORMALIZER, "pn_unpacker: new fragment pkt_len=%d (at %d)", part_size, ptr); |
|
ptr += 2; |
|
|
|
if (part_size<1 || part_size>16384) { |
|
pn->logic_errors++; |
|
DEBUG_ERROR(DEBUG_CATEGORY_NORMALIZER, "PART_SIZE ERROR!!! %d", part_size); |
|
etcp_conn_reinit(pn->etcp); |
|
break; |
|
} |
|
|
|
pn->recvpart = ll_alloc_lldgram(part_size); |
|
if (!pn->recvpart) { |
|
pn->alloc_errors++; |
|
DEBUG_ERROR(DEBUG_CATEGORY_NORMALIZER, "ll_alloc_lldgram failed"); |
|
break; |
|
} |
|
pn->recvpart->len = 0; |
|
} |
|
|
|
uint16_t rem = pn->recvpart->memlen - pn->recvpart->len;// осталось собрать байт |
|
uint16_t avail = len - ptr;// доступно байт сейчас |
|
uint16_t cp = (rem < avail) ? rem : avail; |
|
DEBUG_DEBUG(DEBUG_CATEGORY_NORMALIZER, "copy: remain=%d avail=%d in_ptr=%d out_ptr=%d", rem, avail, ptr, pn->recvpart->len); |
|
memcpy(pn->recvpart->dgram + pn->recvpart->len, payload + ptr, cp); |
|
pn->recvpart->len += cp; |
|
ptr += cp; |
|
|
|
if (pn->recvpart->len == pn->recvpart->memlen) { |
|
DEBUG_DEBUG(DEBUG_CATEGORY_NORMALIZER, "unpacked dgram (size=%d)", pn->recvpart->len); |
|
if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP)) log_dump(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP, "NORM->", pn->recvpart->dgram, pn->recvpart->len); |
|
queue_data_put(pn->output, pn->recvpart); |
|
pn->out_total_pkts++; |
|
pn->out_total_bytes += pn->recvpart->len; |
|
pn->recvpart = NULL; |
|
} |
|
} |
|
|
|
// Free the fragment - dgram was malloc'd in pn_send_to_etcp |
|
memory_pool_free(pn->etcp->instance->data_pool, frag->ll.dgram); |
|
memory_pool_free(pn->etcp->io_pool, frag); |
|
} |
|
|
|
queue_resume_callback(q); |
|
}
|
|
|