Browse Source

Debug: Add ETCP log_name identifier and improve log format

- Add log_name[16] field to ETCP_CONN structure for connection identification
- Add etcp_update_log_name() function to update identifier when peer_node_id is known
- Update all DEBUG_* calls in etcp.c and etcp_loadbalancer.c to include log_name prefix
- Add DEBUG_CATEGORY_NORMALIZER for packet normalizer debug output
- Change log timestamp format to [hh:mm:ss-mmm.uuu] with microseconds precision
- Reorder debug output: (file:line) function() [log_name] message
- Remove duplicate function names from log messages
- Clean up backup files from pkt_normalizer development
nodeinfo-routing-update
Evgeny 8 months ago
parent
commit
c1890ad7a6
  1. 5
      doc/etcp_arch.md
  2. 56
      lib/debug_config.c
  3. 5
      lib/debug_config.h
  4. 315
      src/etcp.c
  5. 6
      src/etcp.h
  6. 12
      src/etcp_connections.c
  7. 26
      src/etcp_loadbalancer.c
  8. 253
      src/pkt_normalizer.c
  9. 282
      src/pkt_normalizer.c.backup
  10. 313
      src/pkt_normalizer.c1
  11. 288
      src/pkt_normalizer.c2
  12. 295
      src/pkt_normalizer.c3
  13. 2
      src/pkt_normalizer.h

5
doc/etcp_arch.md

@ -15,3 +15,8 @@
- обеспечивает обновление метрик каналов etcp_connections
- вызывается при получении ack
Тесты:
# Правила работы с очередями в тестах:
- для добавления в очередь надо использовать queue_wait_threshold(q,0,0,arg) и добавлять по одному пакету, каждый раз дожидаясь когда очередь станет пустой.
- для получения надо использовать queue_set_callback + queue_resume_callback

56
lib/debug_config.c

@ -2,15 +2,16 @@
* Debug module
*/
#include "debug_config.h"
#include <stdlib.h>
#include <string.h>
#include <errno.h>
#include <ctype.h>
#include <stdarg.h>
#include <stdio.h>
#include <time.h>
#include <strings.h>
#include "debug_config.h"
#include <stdlib.h>
#include <string.h>
#include <errno.h>
#include <ctype.h>
#include <stdarg.h>
#include <stdio.h>
#include <time.h>
#include <strings.h>
#include <sys/time.h>
#include <arpa/inet.h>
#include <netinet/in.h>
@ -219,13 +220,14 @@ void debug_output(debug_level_t level, debug_category_t category,
FILE* output = debug_output_file ? debug_output_file : stdout;
/* Add timestamp if enabled */
time_t now = time(NULL);
struct tm* tm_info = localtime(&now);
char time_str[32];
strftime(time_str, sizeof(time_str), "%Y-%m-%d %H:%M:%S", tm_info);
offset += snprintf(buffer + offset, remaining, "[%s] ", time_str);
remaining = BUFFER_SIZE - offset;
/* Add timestamp with microseconds: hh:mm:ss-xxx.yyy */
struct timeval tv;
gettimeofday(&tv, NULL);
struct tm* tm_info = localtime(&tv.tv_sec);
char time_str[32];
strftime(time_str, sizeof(time_str), "%H:%M:%S", tm_info);
offset += snprintf(buffer + offset, remaining, "[%s-%03ld.%03ld] ", time_str, tv.tv_usec / 1000, tv.tv_usec % 1000);
remaining = BUFFER_SIZE - offset;
/* Add level */
const char* level_name = get_level_name(level);
@ -236,17 +238,17 @@ void debug_output(debug_level_t level, debug_category_t category,
offset += snprintf(buffer + offset, remaining, "[%llu] ", (unsigned long long)category);
remaining = BUFFER_SIZE - offset;
/* Add function name if enabled */
if (g_debug_config.function_name_enabled && function) {
offset += snprintf(buffer + offset, remaining, "%s() ", function);
remaining = BUFFER_SIZE - offset;
}
/* Add file:line if enabled */
if (g_debug_config.file_line_enabled && file) {
offset += snprintf(buffer + offset, remaining, "(%s:%d) ", file, line);
remaining = BUFFER_SIZE - offset;
}
/* Add file:line if enabled */
if (g_debug_config.file_line_enabled && file) {
offset += snprintf(buffer + offset, remaining, "(%s:%d) ", file, line);
remaining = BUFFER_SIZE - offset;
}
/* Add function name if enabled */
if (g_debug_config.function_name_enabled && function) {
offset += snprintf(buffer + offset, remaining, "%s() ", function);
remaining = BUFFER_SIZE - offset;
}
/* Add the actual message */
offset += vsnprintf(buffer + offset, remaining, format, args);

5
lib/debug_config.h

@ -40,8 +40,9 @@ typedef uint64_t debug_category_t;
#define DEBUG_CATEGORY_CONFIG ((debug_category_t)1 << 7) // configuration parsing
#define DEBUG_CATEGORY_TUN ((debug_category_t)1 << 8) // TUN interface
#define DEBUG_CATEGORY_ROUTING ((debug_category_t)1 << 9) // routing table
#define DEBUG_CATEGORY_TIMERS ((debug_category_t)1 << 10) // timer management
#define DEBUG_CATEGORY_ALL ((debug_category_t)0xFFFFFFFFFFFFFFFFULL)
#define DEBUG_CATEGORY_TIMERS ((debug_category_t)1 << 10) // timer management
#define DEBUG_CATEGORY_NORMALIZER ((debug_category_t)1 << 11) // packet normalizer
#define DEBUG_CATEGORY_ALL ((debug_category_t)0xFFFFFFFFFFFFFFFFULL)
/* Debug configuration structure */
typedef struct {

315
src/etcp.c

@ -109,13 +109,17 @@ struct ETCP_CONN* etcp_connection_create(struct UTUN_INSTANCE* instance) {
queue_set_callback(etcp->input_send_q, input_send_q_cb, etcp);
queue_set_callback(etcp->input_wait_ack, wait_ack_cb, etcp);
etcp->link_ready_for_send_fn = etcp_link_ready_callback;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_connection_create: connection initialized. ETCP=%p mtu=%d, next_tx_id=%u",
etcp, etcp->mtu, etcp->next_tx_id);
return etcp;
}
etcp->link_ready_for_send_fn = etcp_link_ready_callback;
// Initialize log_name with local node_id (peer will be updated later when known)
snprintf(etcp->log_name, sizeof(etcp->log_name), "%04llu→????",
(unsigned long long)(instance->node_id % 10000));
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] connection initialized. ETCP=%p mtu=%d, next_tx_id=%u",
etcp->log_name, etcp, etcp->mtu, etcp->next_tx_id);
return etcp;
}
// Close connection with NULL pointer safety (prevents double free)
void etcp_connection_close(struct ETCP_CONN* etcp) {
@ -225,17 +229,26 @@ void etcp_connection_close(struct ETCP_CONN* etcp) {
free(etcp);
}
// Reset connection (stub)
void etcp_conn_reset(struct ETCP_CONN* etcp) {
// Reset IDs, queues, etc. as per protocol.txt
etcp->next_tx_id = 1;
etcp->last_rx_id = 0;
etcp->last_delivered_id = 0;
// Clear inflight, rx_list, etc.
}
// ====================================================================== Отправка данных
// Reset connection (stub)
void etcp_conn_reset(struct ETCP_CONN* etcp) {
// Reset IDs, queues, etc. as per protocol.txt
etcp->next_tx_id = 1;
etcp->last_rx_id = 0;
etcp->last_delivered_id = 0;
// Clear inflight, rx_list, etc.
}
// Update log_name when peer_node_id becomes known
void etcp_update_log_name(struct ETCP_CONN* etcp) {
if (!etcp || !etcp->instance) return;
uint64_t local_id = etcp->instance->node_id % 10000;
uint64_t peer_id = etcp->peer_node_id % 10000;
snprintf(etcp->log_name, sizeof(etcp->log_name), "%04llu→%04llu",
(unsigned long long)local_id, (unsigned long long)peer_id);
}
// ====================================================================== Отправка данных
// Send data through ETCP connection
// Allocates memory from data_pool and places in input queue
@ -243,28 +256,28 @@ void etcp_conn_reset(struct ETCP_CONN* etcp) {
int etcp_send(struct ETCP_CONN* etcp, const void* data, uint16_t len) {
// DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_send: ENTER etcp=%p, data=%p, len=%zu", etcp, data, len);
if (!etcp || !data || len == 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: invalid parameters (etcp=%p, data=%p, len=%zu)", etcp, data, len);
return -1;
}
if (!etcp || !data || len == 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] invalid parameters (etcp=%p, data=%p, len=%zu)", etcp->log_name, etcp, data, len);
return -1;
}
if (!etcp->input_queue) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: input_queue is NULL for etcp=%p", etcp);
return -1;
}
if (!etcp->input_queue) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] input_queue is NULL for etcp=%p", etcp->log_name, etcp);
return -1;
}
// Check length against maximum packet size
if (len > PACKET_DATA_SIZE) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: packet too large (len=%zu, max=%d)", len, PACKET_DATA_SIZE);
return -1;
}
// Check length against maximum packet size
if (len > PACKET_DATA_SIZE) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] packet too large (len=%zu, max=%d)", etcp->log_name, len, PACKET_DATA_SIZE);
return -1;
}
// Allocate packet data from data_pool (following ETCP reception pattern)
uint8_t* packet_data = memory_pool_alloc(etcp->instance->data_pool);
if (!packet_data) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: failed to allocate packet data from data_pool");
return -1;
}
// Allocate packet data from data_pool (following ETCP reception pattern)
uint8_t* packet_data = memory_pool_alloc(etcp->instance->data_pool);
if (!packet_data) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to allocate packet data from data_pool", etcp->log_name);
return -1;
}
// Copy user data to packet buffer
memcpy(packet_data, data, len);
@ -272,87 +285,87 @@ int etcp_send(struct ETCP_CONN* etcp, const void* data, uint16_t len) {
// Create queue entry - this allocates ll_entry + data pointer
struct ETCP_FRAGMENT* pkt = (struct ETCP_FRAGMENT*)queue_entry_new_from_pool(etcp->io_pool);
if (!pkt) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: failed to allocate queue entry");
memory_pool_free(etcp->instance->data_pool, packet_data);
return -1;
}
pkt->seq = 0; // Will be assigned by input_queue_cb
pkt->timestamp = 0; // Will be set by input_queue_cb
pkt->ll.dgram = packet_data; // Point to data_pool allocation
pkt->ll.len = len; // размер packet_data
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_send: created PACKET %p with data %p (len=%zu)", pkt, packet_data, len);
struct ETCP_FRAGMENT* pkt = (struct ETCP_FRAGMENT*)queue_entry_new_from_pool(etcp->io_pool);
if (!pkt) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to allocate queue entry", etcp->log_name);
memory_pool_free(etcp->instance->data_pool, packet_data);
return -1;
}
pkt->seq = 0; // Will be assigned by input_queue_cb
pkt->timestamp = 0; // Will be set by input_queue_cb
pkt->ll.dgram = packet_data; // Point to data_pool allocation
pkt->ll.len = len; // размер packet_data
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] created PACKET %p with data %p (len=%zu)", etcp->log_name, pkt, packet_data, len);
// Add to input queue - input_queue_cb will process it
if (queue_data_put(etcp->input_queue, (struct ll_entry*)pkt, 0) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_send: failed to add to input queue");
memory_pool_free(etcp->instance->data_pool, packet_data);
memory_pool_free(etcp->io_pool, pkt);
return -1;
}
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to add to input queue", etcp->log_name);
memory_pool_free(etcp->instance->data_pool, packet_data);
memory_pool_free(etcp->io_pool, pkt);
return -1;
}
return 0;
}
static void input_queue_try_resume(struct ETCP_CONN* etcp) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "input_queue_try_resume: ENTER etcp=%p", etcp);
static void input_queue_try_resume(struct ETCP_CONN* etcp) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] ENTER etcp=%p", etcp->log_name, etcp);
// если размер input_wait_ack+input_send_q в байтах < optimal_inflight то resume сейчас.
size_t wait_ack_bytes = queue_total_bytes(etcp->input_wait_ack);
size_t send_q_bytes = queue_total_bytes(etcp->input_send_q);
size_t total_bytes = wait_ack_bytes + send_q_bytes;
if (total_bytes < etcp->optimal_inflight) {
queue_resume_callback(etcp->input_send_q);// вызвать лишний раз resume не страшно.
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "input_queue_try_resume: resumed input_send_q callback");
}
}
if (total_bytes < etcp->optimal_inflight) {
queue_resume_callback(etcp->input_send_q);// вызвать лишний раз resume не страшно.
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] resumed input_send_q callback", etcp->log_name);
}
}
void etcp_stats(struct ETCP_CONN* etcp) {
if (!etcp) return;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "ETCP stats for conn=%p:", etcp);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] stats for conn=%p:", etcp->log_name, etcp);
// Queue statistics
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, " Queues:");
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, " input_queue: %zu pkts, %zu bytes",
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] Queues:", etcp->log_name);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] input_queue: %zu pkts, %zu bytes", etcp->log_name,
queue_entry_count(etcp->input_queue), queue_total_bytes(etcp->input_queue));
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, " input_send_q: %zu pkts, %zu bytes",
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] input_send_q: %zu pkts, %zu bytes", etcp->log_name,
queue_entry_count(etcp->input_send_q), queue_total_bytes(etcp->input_send_q));
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, " input_wait_ack: %zu pkts, %zu bytes",
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] input_wait_ack: %zu pkts, %zu bytes", etcp->log_name,
queue_entry_count(etcp->input_wait_ack), queue_total_bytes(etcp->input_wait_ack));
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, " ack_q: %zu pkts",
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] ack_q: %zu pkts", etcp->log_name,
queue_entry_count(etcp->ack_q));
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, " recv_q: %zu pkts",
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] recv_q: %zu pkts", etcp->log_name,
queue_entry_count(etcp->recv_q));
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, " output_queue: %zu pkts, %zu bytes",
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] output_queue: %zu pkts, %zu bytes", etcp->log_name,
queue_entry_count(etcp->output_queue), queue_total_bytes(etcp->output_queue));
// RTT metrics
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, " RTT metrics:");
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, " rtt_last: %u (0.1ms)", etcp->rtt_last);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, " rtt_avg_10: %u (0.1ms)", etcp->rtt_avg_10);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, " rtt_avg_100: %u (0.1ms)", etcp->rtt_avg_100);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, " jitter: %u (0.1ms)", etcp->jitter);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] RTT metrics:", etcp->log_name);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] rtt_last: %u (0.1ms)", etcp->log_name, etcp->rtt_last);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] rtt_avg_10: %u (0.1ms)", etcp->log_name, etcp->rtt_avg_10);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] rtt_avg_100: %u (0.1ms)", etcp->log_name, etcp->rtt_avg_100);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] jitter: %u (0.1ms)", etcp->log_name, etcp->jitter);
// Counters
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, " Counters:");
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, " bytes_sent_total: %u", etcp->bytes_sent_total);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, " retransmissions_count: %u", etcp->retransmissions_count);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, " ack_packets_count: %u", etcp->ack_packets_count);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, " unacked_bytes: %u", etcp->unacked_bytes);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] Counters:", etcp->log_name);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] bytes_sent_total: %u", etcp->log_name, etcp->bytes_sent_total);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] retransmissions_count: %u", etcp->log_name, etcp->retransmissions_count);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] ack_packets_count: %u", etcp->log_name, etcp->ack_packets_count);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] unacked_bytes: %u", etcp->log_name, etcp->unacked_bytes);
// IDs
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, " IDs:");
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, " next_tx_id: %u", etcp->next_tx_id);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, " last_rx_id: %u", etcp->last_rx_id);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, " last_delivered_id:%u", etcp->last_delivered_id);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, " rx_ack_till: %u", etcp->rx_ack_till);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] IDs:", etcp->log_name);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] next_tx_id: %u", etcp->log_name, etcp->next_tx_id);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] last_rx_id: %u", etcp->log_name, etcp->last_rx_id);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] last_delivered_id:%u", etcp->log_name, etcp->last_delivered_id);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] rx_ack_till: %u", etcp->log_name, etcp->rx_ack_till);
}
// Input callback for input_queue (добавление новых кодограмм в стек)
@ -362,23 +375,23 @@ static void input_queue_cb(struct ll_queue* q, void* arg) {
struct ETCP_CONN* etcp = (struct ETCP_CONN*)arg;
struct ETCP_FRAGMENT* in_pkt = (struct ETCP_FRAGMENT*)queue_data_get(q);
if (!in_pkt) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "input_queue_cb: cannot get element (pool=%p etcp=%p)", etcp->inflight_pool, etcp);
queue_resume_callback(q);
return;
}
if (!in_pkt) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] cannot get element (pool=%p etcp=%p)", etcp->log_name, etcp->inflight_pool, etcp);
queue_resume_callback(q);
return;
}
memory_pool_free(etcp->io_pool, in_pkt);// перемещаем из io_pool в inflight_pool
// Create INFLIGHT_PACKET
struct INFLIGHT_PACKET* p = memory_pool_alloc(etcp->inflight_pool);
if (!p) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "input_queue_cb: cannot allocate INFLIGHT_PACKET (pool=%p etcp=%p)", etcp->inflight_pool, etcp);
// Create INFLIGHT_PACKET
struct INFLIGHT_PACKET* p = memory_pool_alloc(etcp->inflight_pool);
if (!p) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] cannot allocate INFLIGHT_PACKET (pool=%p etcp=%p)", etcp->log_name, etcp->inflight_pool, etcp);
queue_entry_free((struct ll_entry*)in_pkt); // Free the ETCP_FRAGMENT
queue_resume_callback(q);
return;
}
}
// Setup inflight packet (based on protocol.txt)
memset(p, 0, sizeof(*p));
@ -389,13 +402,13 @@ static void input_queue_cb(struct ll_queue* q, void* arg) {
p->ll.dgram_pool = in_pkt->ll.dgram_pool;
p->ll.len = in_pkt->ll.len;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "input_queue_cb: input -> inflight (seq=%u, len=%u)", p->seq, p->ll.len);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] input -> inflight (seq=%u, len=%u)", etcp->log_name, p->seq, p->ll.len);
// Add to send queue
if (queue_data_put(etcp->input_send_q, (struct ll_entry*)p, p->seq) != 0) {
memory_pool_free(etcp->inflight_pool, p);
queue_entry_free((struct ll_entry*)in_pkt);
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "input_queue_cb: EXIT (queue put failed)");
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] EXIT (queue put failed)", etcp->log_name);
return;
}
// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "input_queue_cb: successfully moved from input_queue to input_send_q");
@ -429,8 +442,8 @@ static void ack_timeout_check(void* arg) {
struct INFLIGHT_PACKET* pkt = (struct INFLIGHT_PACKET*)current;
uint64_t elapsed = now - pkt->last_timestamp;
if (elapsed > timeout) {
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "ack_timeout_check: timeout for seq=%u, elapsed=%llu, timeout=%llu, send_count=%u. Moving: wait_ack -> send_q",
pkt->seq, (unsigned long long)elapsed, (unsigned long long)timeout, pkt->send_count);
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] ack_timeout_check: timeout for seq=%u, elapsed=%llu, timeout=%llu, send_count=%u. Moving: wait_ack -> send_q",
etcp->log_name, pkt->seq, (unsigned long long)elapsed, (unsigned long long)timeout, pkt->send_count);
// Increment counters
pkt->send_count++;
@ -452,7 +465,7 @@ static void ack_timeout_check(void* arg) {
// shedule timer
int64_t next_timeout=timeout - elapsed;
etcp->retrans_timer = uasync_set_timeout(etcp->instance->ua, next_timeout+10, etcp, ack_timeout_check);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "ack_timeout_check: retransmission timer set for %llu units", next_timeout);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] ack_timeout_check: retransmission timer set for %llu units", etcp->log_name, next_timeout);
return;
}
}
@ -470,11 +483,11 @@ static void wait_ack_cb(struct ll_queue* q, void* arg) {
// вызывается линком когда освобождается или очередью если появляются данные на передачу
struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) {
struct ETCP_LINK* link = etcp_loadbalancer_select_link(etcp);
if (!link) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: no link available");
return NULL;// если линков нет - ждём появления свободного
}
struct ETCP_LINK* link = etcp_loadbalancer_select_link(etcp);
if (!link) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] no link available", etcp->log_name);
return NULL;// если линков нет - ждём появления свободного
}
size_t send_q_size = queue_entry_count(etcp->input_send_q);
@ -488,7 +501,7 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) {
// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: getting packet from input_send_q");
struct INFLIGHT_PACKET* inf_pkt = (struct INFLIGHT_PACKET*)queue_data_get(etcp->input_send_q);
if (inf_pkt) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: prepare udp dgram for send packet %p (seq=%u, len=%u)", inf_pkt, inf_pkt->seq, inf_pkt->ll.len);
// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] prepare udp dgram for send packet %p (seq=%u, len=%u)", etcp->log_name, inf_pkt, inf_pkt->seq, inf_pkt->ll.len);
inf_pkt->last_timestamp=get_current_time_units();
inf_pkt->send_count++;
@ -498,16 +511,16 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) {
size_t ack_q_size = queue_entry_count(etcp->ack_q);
if (!inf_pkt && ack_q_size == 0) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: no data/ack to send");
return NULL;
}
if (!inf_pkt && ack_q_size == 0) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] no data/ack to send", etcp->log_name);
return NULL;
}
struct ETCP_DGRAM* dgram = memory_pool_alloc(etcp->instance->pkt_pool);
if (!dgram) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: failed to allocate ETCP_DGRAM");
return NULL;
}
struct ETCP_DGRAM* dgram = memory_pool_alloc(etcp->instance->pkt_pool);
if (!dgram) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to allocate ETCP_DGRAM", etcp->log_name);
return NULL;
}
dgram->link = link;
dgram->noencrypt_len=0;
@ -539,7 +552,7 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) {
uint16_t dly=get_current_timestamp()-ack_pkt->recv_timestamp;
dgram->data[ptr++]=dly;
dgram->data[ptr++]=dly>>8;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: add ACK N%d dTS=%d", ack_pkt->seq, dly);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] add ACK N%d dTS=%d", etcp->log_name, ack_pkt->seq, dly);
queue_entry_free((struct ll_entry*)ack_pkt);
if (inf_pkt && inf_pkt->ll.len+ptr>=etcp->mtu-10) break;// pkt len (надо просчитать точнее включая все заголовки)
@ -548,9 +561,9 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) {
dgram->data[1]=ptr/8;
if (inf_pkt) {
// фрейм data (0) обязательно в конец
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: ready to send packet with payload (seq=%u, len=%u), ack_size=%d", inf_pkt->seq, inf_pkt->ll.len, dgram->data[1]);
if (inf_pkt) {
// фрейм data (0) обязательно в конец
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] add DATA (seq=%u, len=%u), ack_size=%d", etcp->log_name, inf_pkt->seq, inf_pkt->ll.len, dgram->data[1]);
dgram->data[ptr++]=0;// payload
dgram->data[ptr++]=inf_pkt->seq;
dgram->data[ptr++]=inf_pkt->seq>>8;
@ -559,9 +572,9 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) {
memcpy(&dgram->data[ptr], inf_pkt->ll.dgram, inf_pkt->ll.len); ptr+=inf_pkt->ll.len;
}
else {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_request_pkt: built packet with %d bytes total", dgram->data_len);
}
else {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] only ACK packet with %d bytes total", etcp->log_name, dgram->data_len);
}
dgram->data_len=ptr;
@ -596,10 +609,10 @@ static void ack_response_timer_cb(void* arg) {// проверяем неотпр
// ====================================================================== Прием данных
void etcp_output_try_assembly(struct ETCP_CONN* etcp) {
// пробуем собрать выходную очередь из фрагментов
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_output_try_assembly: etcp=%p, last_delivered_id=%u, recv_q_count=%d",
etcp, etcp->last_delivered_id, queue_entry_count(etcp->recv_q));
void etcp_output_try_assembly(struct ETCP_CONN* etcp) {
// пробуем собрать выходную очередь из фрагментов
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] etcp=%p, last_delivered_id=%u, recv_q_count=%d",
etcp->log_name, etcp, etcp->last_delivered_id, queue_entry_count(etcp->recv_q));
uint32_t next_expected_id = etcp->last_delivered_id + 1;
int delivered_count = 0;
@ -610,11 +623,11 @@ void etcp_output_try_assembly(struct ETCP_CONN* etcp) {
struct ETCP_FRAGMENT* rx_pkt = (struct ETCP_FRAGMENT*)queue_find_data_by_id(etcp->recv_q, next_expected_id);
if (!rx_pkt) {
// No more contiguous packets found
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_output_try_assembly: no packet found for id=%u, stopping", next_expected_id);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] no packet found for id=%u, stopping", etcp->log_name, next_expected_id);
break;
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_output_try_assembly: assembling packet id=%u (len=%u)",
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] assembling packet id=%u (len=%u)", etcp->log_name,
rx_pkt->seq, rx_pkt->ll.len);
// Simply move ETCP_FRAGMENT from recv_q to output_queue - no data copying needed
@ -625,10 +638,10 @@ void etcp_output_try_assembly(struct ETCP_CONN* etcp) {
if (queue_data_put(etcp->output_queue, (struct ll_entry*)rx_pkt, next_expected_id) == 0) {
delivered_bytes += rx_pkt->ll.len;
delivered_count++;
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_output_try_assembly: moved packet id=%u to output_queue",
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] moved packet id=%u to output_queue", etcp->log_name,
next_expected_id);
} else {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_output_try_assembly: failed to add packet id=%u to output_queue",
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to add packet id=%u to output_queue", etcp->log_name,
next_expected_id);
// Put it back in recv_q if we can't add to output_queue
queue_data_put(etcp->recv_q, (struct ll_entry*)rx_pkt, next_expected_id);
@ -640,15 +653,15 @@ void etcp_output_try_assembly(struct ETCP_CONN* etcp) {
next_expected_id++;
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_output_try_assembly: delivered %u contiguous packets (%u bytes), last_delivered_id=%u, output_queue_count=%d",
delivered_count, delivered_bytes, etcp->last_delivered_id, queue_entry_count(etcp->output_queue));
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] delivered %u contiguous packets (%u bytes), last_delivered_id=%u, output_queue_count=%d",
etcp->log_name, delivered_count, delivered_bytes, etcp->last_delivered_id, queue_entry_count(etcp->output_queue));
}
// Process ACK receipt - remove acknowledged packet from inflight queues
void etcp_ack_recv(struct ETCP_CONN* etcp, uint32_t seq, uint16_t ts, uint16_t dts) {
if (!etcp) return;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_ack_recv: processing ACK for seq=%u, ts=%u, dts=%u", seq, ts, dts);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] processing ACK for seq=%u, ts=%u, dts=%u", etcp->log_name, seq, ts, dts);
// Find the acknowledged packet in the wait_ack queue
struct INFLIGHT_PACKET* acked_pkt = (struct INFLIGHT_PACKET*)queue_find_data_by_id(etcp->input_wait_ack, seq);
@ -690,16 +703,16 @@ void etcp_ack_recv(struct ETCP_CONN* etcp, uint32_t seq, uint16_t ts, uint16_t d
}
etcp->jitter = rtt_max - rtt_min;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_ack_recv: RTT updated - last=%u, avg_10=%u, avg_100=%u, jitter=%u",
rtt, etcp->rtt_avg_10, etcp->rtt_avg_100, etcp->jitter);
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] RTT updated - last=%u, avg_10=%u, avg_100=%u, jitter=%u",
etcp->log_name, rtt, etcp->rtt_avg_10, etcp->rtt_avg_100, etcp->jitter);
}
// Update connection statistics
etcp->unacked_bytes -= acked_pkt->ll.len;
etcp->bytes_sent_total += acked_pkt->ll.len;
etcp->ack_packets_count++;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_ack_recv: removed packet seq=%u from wait_ack, unacked_bytes now %u", seq, etcp->unacked_bytes);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] removed packet seq=%u from wait_ack, unacked_bytes now %u", etcp->log_name, seq, etcp->unacked_bytes);
if (acked_pkt->ll.dgram) {
memory_pool_free(etcp->instance->data_pool, acked_pkt->ll.dgram);
@ -709,8 +722,8 @@ void etcp_ack_recv(struct ETCP_CONN* etcp, uint32_t seq, uint16_t ts, uint16_t d
// Try to resume sending more packets if window space opened up
input_queue_try_resume(etcp);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_ack_recv: completed for seq=%u", seq);
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] completed for seq=%u", etcp->log_name, seq);
}
// Process incoming decrypted packet
@ -755,13 +768,13 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) {
p->pkt_timestamp=pkt->timestamp;
p->recv_timestamp=get_current_timestamp();
queue_data_put(etcp->ack_q, (struct ll_entry*)p, p->seq);
if (etcp->ack_resp_timer == NULL) {
etcp->ack_resp_timer = uasync_set_timeout(etcp->instance->ua, ACK_DELAY_TB, etcp, ack_response_timer_cb);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_conn_input: set ack_timer for delayed ACK send");
}
if (etcp->ack_resp_timer == NULL) {
etcp->ack_resp_timer = uasync_set_timeout(etcp->instance->ua, ACK_DELAY_TB, etcp, ack_response_timer_cb);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] set ack_timer for delayed ACK send", etcp->log_name);
}
if ((int32_t)(etcp->last_delivered_id-seq)<0) if (queue_find_data_by_id(etcp->recv_q, seq)==NULL) {// проверяем есть ли пакет с этим seq
uint32_t pkt_len=len-5;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_conn_input: adding packet seq=%u to recv_q (last_delivered_id=%u)", seq, etcp->last_delivered_id);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] adding packet seq=%u to recv_q (last_delivered_id=%u)", etcp->log_name, seq, etcp->last_delivered_id);
// отправляем пакет в очередь на сборку
uint8_t* payload_data = memory_pool_alloc(etcp->instance->data_pool);
struct ETCP_FRAGMENT* rx_pkt = (struct ETCP_FRAGMENT*)queue_entry_new_from_pool(etcp->io_pool);
@ -773,7 +786,7 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) {
// Copy the actual payload data
memcpy(payload_data, data + 5, pkt_len);
queue_data_put(etcp->recv_q, (struct ll_entry*)rx_pkt, seq);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_conn_input: packet seq=%u added to recv_q, calling assembly (last_delivered_id=%u)", seq, etcp->last_delivered_id);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] packet seq=%u added to recv_q, calling assembly (last_delivered_id=%u)", etcp->log_name, seq, etcp->last_delivered_id);
if (etcp->last_delivered_id+1==seq) etcp_output_try_assembly(etcp);// пробуем собрать выходную очередь из фрагментов
}
}

6
src/etcp.h

@ -135,6 +135,9 @@ struct ETCP_CONN {
// Flags
// uint8_t wait_timeout_active; // In wait timeout state - Not used
// Logging identifier (format: "XXXX→XXXX" - last 4 digits of local and peer node_id)
char log_name[16];
};
// Functions
@ -161,6 +164,9 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp);
// Process ACK receipt - remove acknowledged packet from inflight queues
void etcp_ack_recv(struct ETCP_CONN* etcp, uint32_t seq, uint16_t ts, uint16_t dts);
// Update log_name when peer_node_id becomes known
void etcp_update_log_name(struct ETCP_CONN* etcp);
#ifdef __cplusplus
}
#endif

12
src/etcp_connections.c

@ -444,7 +444,7 @@ void etcp_link_close(struct ETCP_LINK* link) {
}
int etcp_encrypt_send(struct ETCP_DGRAM* dgram) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_encrypt_send called, link=%p", dgram ? dgram->link : NULL);
// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_encrypt_send called, link=%p", dgram ? dgram->link : NULL);
// printf("[ETCP DEBUG] etcp_encrypt_send: ENTERING FUNCTION\n");
int errcode=0;
@ -478,7 +478,7 @@ int etcp_encrypt_send(struct ETCP_DGRAM* dgram) {
struct sockaddr_in* sin = (struct sockaddr_in*)addr;
char addr_str[INET_ADDRSTRLEN];
inet_ntop(AF_INET, &sin->sin_addr, addr_str, INET_ADDRSTRLEN);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[ETCP] Sending packet to %s:%d, size=%zd", addr_str, ntohs(sin->sin_port), enc_buf_len + dgram->noencrypt_len);
// DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[ETCP] Sending packet to %s:%d, size=%zd", addr_str, ntohs(sin->sin_port), enc_buf_len + dgram->noencrypt_len);
}
ssize_t sent = sendto(dgram->link->conn->fd, enc_buf, enc_buf_len + dgram->noencrypt_len, 0, (struct sockaddr*)addr, addr_len);
@ -486,7 +486,7 @@ int etcp_encrypt_send(struct ETCP_DGRAM* dgram) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "sendto failed, errno=%d", errno);
dgram->link->send_errors++; errcode=3; goto es_err;
} else {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "sendto succeeded, sent=%zd bytes to port %d", sent, ntohs(((struct sockaddr_in*)addr)->sin_port));
// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "sendto succeeded, sent=%zd bytes to port %d", sent, ntohs(((struct sockaddr_in*)addr)->sin_port));
dgram->link->total_encrypted += sent;
}
return (int)sent;
@ -496,7 +496,7 @@ es_err:
}
static void etcp_connections_read_callback(int fd, void* arg) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_connections_read_callback fd=%d, socket=%p", fd, arg);
// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_connections_read_callback fd=%d, socket=%p", fd, arg);
// !!!!!! DANGER: в этой функции ПРЕДЕЛЬНАЯ АККУРАТНОСТЬ. Если кажется что не туда указатель то невнимательно аланизировал !!!!!
// НЕ РУИНИТЬ (uint8_t*)&pkt->timestamp - это правильно !!!!
//
@ -586,6 +586,7 @@ static void etcp_connections_read_callback(int fd, void* arg) {
if (!conn) { errorcode=55; DEBUG_ERROR(DEBUG_CATEGORY_CONNECTION, "etcp_connections_read_callback: failed to create connection"); goto ec_fr; }// облом
memcpy(&conn->crypto_ctx, &sc, sizeof(sc));// добавляем ключ
conn->peer_node_id=peer_id;
etcp_update_log_name(conn); // Update log_name with peer_node_id
char buf[128];
addr_to_string(&addr, buf, sizeof(buf));
DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "New connection from %s peer_id=%ld etcp=%p", buf, peer_id, conn);
@ -663,6 +664,7 @@ static void etcp_connections_read_callback(int fd, void* arg) {
// DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Received INIT_RESPONSE from server_node_id=%llu, mtu=%d", (unsigned long long)server_node_id, link->mtu);
link->etcp->peer_node_id = server_node_id; // If not set
etcp_update_log_name(link->etcp); // Update log_name with peer_node_id
// Mark link as initialized
// DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "Setting link->initialized=1, link=%p, is_server=%d", link, link->is_server);
@ -757,6 +759,7 @@ int init_connections(struct UTUN_INSTANCE* instance) {
// For now, set peer node ID to indicate we have peer key
// The actual peer key will be exchanged during connection establishment
etcp_conn->peer_node_id = 1; // Simple indicator
etcp_update_log_name(etcp_conn); // Update log_name with peer_node_id
DEBUG_INFO(DEBUG_CATEGORY_CRYPTO, "init_connections: setting peer public key for client %s", client->name);
// Set peer public key (assuming hex format)
@ -826,6 +829,7 @@ int init_connections(struct UTUN_INSTANCE* instance) {
if (strlen(client->peer_public_key_hex) > 0) {
// For now, set peer node ID to indicate we have peer key
// The actual peer key will be exchanged during connection establishment
etcp_update_log_name(etcp_conn); // Update log_name with peer_node_id
etcp_conn->peer_node_id = 1; // Simple indicator
DEBUG_INFO(DEBUG_CATEGORY_CRYPTO, "init_connections: setting peer public key for client %s", client->name);

26
src/etcp_loadbalancer.c

@ -26,8 +26,8 @@ static void shaper_timer_cb(void* arg); // Shaper wait callback
// Select link for transmission
struct ETCP_LINK* etcp_loadbalancer_select_link(struct ETCP_CONN* etcp) {
if (!etcp || !etcp->links) {
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_loadbalancer_select_link: invalid parameters (etcp=%p, links=%p)",
etcp, etcp ? etcp->links : NULL);
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] invalid parameters (etcp=%p, links=%p)",
etcp ? etcp->log_name : "????→????", etcp, etcp ? etcp->links : NULL);
return NULL;
}
@ -63,10 +63,10 @@ struct ETCP_LINK* etcp_loadbalancer_select_link(struct ETCP_CONN* etcp) {
}
if (best) {
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "etcp_loadbalancer_select_link: selected link %p (load_tb=%llu)",
best, (unsigned long long)min_load_tb);
// DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] selected link %p (load_tb=%llu)",
// etcp->log_name, best, (unsigned long long)min_load_tb);
} else {
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_loadbalancer_select_link: no suitable link found");
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] no suitable link found", etcp->log_name);
}
return best;
@ -81,7 +81,7 @@ void etcp_loadbalancer_send(struct ETCP_DGRAM* dgram) {
struct ETCP_CONN* etcp = dgram->link ? dgram->link->etcp : NULL;
if (!etcp) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_loadbalancer_send: no ETCP_CONN associated with dgram");
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] no ETCP_CONN associated with dgram", etcp->log_name);
memory_pool_free(etcp->instance->pkt_pool, dgram);
// free(dgram);
return;
@ -91,7 +91,7 @@ void etcp_loadbalancer_send(struct ETCP_DGRAM* dgram) {
if (!dgram->link) {
dgram->link = etcp_loadbalancer_select_link(etcp);
if (!dgram->link) {
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "etcp_loadbalancer_send: no link available, dropping dgram");
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] no link available, dropping dgram", etcp->log_name);
memory_pool_free(etcp->instance->pkt_pool, dgram);
// free(dgram); // Assume free; adjust if pooled
return;
@ -101,18 +101,18 @@ void etcp_loadbalancer_send(struct ETCP_DGRAM* dgram) {
struct ETCP_LINK* link = dgram->link;
size_t pkt_size = dgram->data_len + sizeof(uint16_t); // Include timestamp
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "etcp_loadbalancer_send: sending dgram on link=%p, size=%zu", link, pkt_size);
// DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] sending dgram on link=%p, size=%zu", etcp->log_name, link, pkt_size);
// Encrypt and send (from etcp_connections.c)
int send_result = etcp_encrypt_send(dgram);
if (send_result < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "etcp_loadbalancer_send: encrypt/send failed (%d)", send_result);
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] encrypt/send failed (%d)", etcp->log_name, send_result);
// Don't return here, let the function continue to free dgram at the end
}
// Update shaper after successful send
if (link->bandwidth == 0) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "etcp_loadbalancer_send: unlimited bandwidth, skipping shaper update");
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] unlimited bandwidth, skipping shaper update", etcp->log_name);
// Don't return here, let the function continue to free dgram at the end
} else {
// Time to transmit (ns per byte = 8e9 / bw for bits/sec)
@ -133,14 +133,14 @@ void etcp_loadbalancer_send(struct ETCP_DGRAM* dgram) {
if (link->shaper_load_time_tb >= now_tb + SHAPER_BURST_DELAY_TB) {
uint64_t wait_tb = link->shaper_load_time_tb - now_tb;
link->shaper_timer = uasync_set_timeout(link->etcp->instance->ua, wait_tb, link, shaper_timer_cb);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_loadbalancer_send: scheduled shaper timer (wait_tb=%llu)", (unsigned long long)wait_tb);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] scheduled shaper timer (wait_tb=%llu)", etcp->log_name, (unsigned long long)wait_tb);
}
// Inactivity correction
if (link->shaper_load_time_tb < now_tb-ALLOWED_DELTA) link->shaper_load_time_tb = now_tb-ALLOWED_DELTA;
}
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "etcp_loadbalancer_send: updated load_tb=%llu, sub=%llu, state=%u",
(unsigned long long)link->shaper_load_time_tb, (unsigned long long)link->shaper_sub_nanotime, link->shaper_timer==NULL?1:0);
// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] updated load_tb=%llu, sub=%llu, state=%u", etcp->log_name,
// (unsigned long long)link->shaper_load_time_tb, (unsigned long long)link->shaper_sub_nanotime, link->shaper_timer==NULL?1:0);
memory_pool_free(etcp->instance->pkt_pool, dgram);
// free(dgram); // Free the dgram in all cases - we own it

253
src/pkt_normalizer.c

@ -6,7 +6,7 @@
#include <stdlib.h>
#include <string.h>
#include <stdio.h> // For debugging (can be removed if not needed)
#include "debug_config.h" // Assuming this for DEBUG_ERROR
#include "../lib/debug_config.h" // For DEBUG macros
// Forward declarations
static void packer_cb(struct ll_queue* q, void* arg);
@ -24,9 +24,8 @@ struct PKTNORM* pn_init(struct ETCP_CONN* etcp) {
pn->etcp = etcp;
pn->ua = etcp->instance->ua;
// frag_size = MTU минус резерв для дополнительных опций (ACK и прочие заголовки)
pn->frag_size = etcp->mtu - 100;
pn->tx_wait_time = 1; // 1ms timer for packet coalescing
pn->frag_size = etcp->mtu - 100; // Use MTU as fixed packet size (adjust if headers need subtraction)
pn->tx_wait_time = 10;
pn->input = queue_new(pn->ua, 0); // No hash needed
pn->output = queue_new(pn->ua, 0); // No hash needed
@ -41,7 +40,6 @@ struct PKTNORM* pn_init(struct ETCP_CONN* etcp) {
pn->data = NULL;
pn->recvpart = NULL;
pn->recvpart_rem = 0;
pn->flush_timer = NULL;
return pn;
@ -96,14 +94,13 @@ void pn_unpacker_reset_state(struct PKTNORM* pn) {
queue_entry_free(pn->recvpart);
pn->recvpart = NULL;
}
pn->recvpart_rem = 0;
}
// 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) {
if (!pn || !data || len == 0) return;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_packer_send: pn=%p, len=%d", pn, len);
DEBUG_INFO(DEBUG_CATEGORY_NORMALIZER, "pn_packer_send: pn=%p, len=%d", pn, len);
struct ll_entry* entry = ll_alloc_lldgram(len);
if (!entry) return;
@ -112,7 +109,7 @@ void pn_packer_send(struct PKTNORM* pn, uint8_t* data, uint16_t len) {
entry->dgram_pool = NULL;
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));
DEBUG_INFO(DEBUG_CATEGORY_NORMALIZER, "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) {
@ -125,25 +122,20 @@ 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;
int input_count = queue_entry_count(pn->input);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_packer: packer_cb called, input_count=%d", input_count);
DEBUG_INFO(DEBUG_CATEGORY_NORMALIZER, "pn_packer: packer_cb");
// Process all available packets immediately
etcp_input_ready_cb(q, pn);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_packer: packer_cb finished");
queue_wait_threshold(pn->etcp->input_queue, 0, 0, etcp_input_ready_cb, pn);
}
// Helper to send block to ETCP as ETCP_FRAGMENT
static void pn_send_to_etcp(struct PKTNORM* pn) {
if (!pn || !pn->data || pn->data_ptr == 0) return;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_packer: pn_send_to_etcp");
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_ETCP, "pn_packer: send to etcp alloc error");
DEBUG_ERROR(DEBUG_CATEGORY_NORMALIZER, "pn_packer: send to etcp alloc error");
pn->alloc_errors++;
pn->data_ptr = 0;
return;
@ -152,26 +144,27 @@ static void pn_send_to_etcp(struct PKTNORM* pn) {
frag->seq = 0;
frag->timestamp = 0;
frag->ll.dgram = pn->data;
frag->ll.dgram_pool = pn->etcp->instance->data_pool;
frag->ll.len = pn->data_ptr;
frag->ll.memlen = pn->etcp->instance->data_pool->object_size;
queue_data_put(pn->etcp->input_queue, (struct ll_entry*)frag, 0);
// Сбросить структуру (dgram передан во фрагмент, не освобождаем)
pn->data = NULL;
}
// Internal: Renew buffer for packer
// Internal: Renew sndpart buffer
static void pn_buf_renew(struct PKTNORM* pn) {
if (pn->data && pn->data_ptr > 0) {
pn_send_to_etcp(pn);
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);
size_t size = pn->etcp->instance->data_pool->object_size;
if (size > pn->frag_size) size = pn->frag_size;
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;
pn->data_ptr=0;
DEBUG_INFO(DEBUG_CATEGORY_NORMALIZER, "pn_packer: new bufer size=%d bytes",size);
}
}
@ -180,118 +173,44 @@ static void etcp_input_ready_cb(struct ll_queue* q, void* arg) {
struct PKTNORM* pn = (struct PKTNORM*)arg;
if (!pn) return;
int processed = 0;
int input_count = queue_entry_count(pn->input);
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_packer: etcp_input_ready_cb started, input_count=%d", input_count);
// Process all available packets from pn->input
struct ll_entry* in_dgram;
int loop_count = 0;
while ((in_dgram = queue_data_get(pn->input)) != NULL) {
loop_count++;
processed++;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_packer: processing packet %d, len=%d", loop_count, in_dgram->len);
// Ensure buffer is allocated
if (!pn->data) {
pn_buf_renew(pn);
if (!pn->data) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "pn_packer: failed to allocate buffer");
queue_dgram_free(in_dgram);
queue_entry_free(in_dgram);
continue;
}
}
// Check if packet fits in current buffer (need 2 bytes for size header + packet data)
int space_needed = in_dgram->len + 2;
int space_available = pn->data_size - pn->data_ptr;
// DEBUG_INFO(DEBUG_CATEGORY_NORMALIZER, "pn_packer: etcp_input_ready_cb");
struct ll_entry* in_dgram = queue_data_get(pn->input);
if (!in_dgram) { queue_resume_callback(pn->input); return; }
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 (space_available < space_needed && pn->data_ptr > 0) {
// Not enough space - flush current buffer first
pn_send_to_etcp(pn);
pn_buf_renew(pn);
space_available = pn->data_size - pn->data_ptr;
}
exit:
queue_dgram_free(in_dgram);
queue_entry_free(in_dgram);
if (space_available < space_needed) {
// Packet is too big for single fragment - need to split across multiple fragments
// First fragment: [2 bytes total_size][data_part1]
// Next fragments: [data_part2][data_part3]...
int total_size = in_dgram->len;
int bytes_copied = 0;
int is_first_fragment = 1;
while (bytes_copied < total_size) {
if (!pn->data) {
pn_buf_renew(pn);
if (!pn->data) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "pn_packer: failed to allocate buffer for split");
break;
}
}
int remaining = total_size - bytes_copied;
int chunk;
if (is_first_fragment) {
// First fragment: include 2-byte size header
int space_for_data = pn->data_size - pn->data_ptr - 2;
chunk = (remaining < space_for_data) ? remaining : space_for_data;
// Write size header
pn->data[pn->data_ptr++] = total_size & 0xFF;
pn->data[pn->data_ptr++] = (total_size >> 8) & 0xFF;
// Copy first chunk of data
memcpy(pn->data + pn->data_ptr, in_dgram->dgram, chunk);
pn->data_ptr += chunk;
bytes_copied += chunk;
is_first_fragment = 0;
} else {
// Subsequent fragments: only data, no header
chunk = (remaining < pn->data_size - pn->data_ptr) ? remaining : pn->data_size - pn->data_ptr;
memcpy(pn->data + pn->data_ptr, in_dgram->dgram + bytes_copied, chunk);
pn->data_ptr += chunk;
bytes_copied += chunk;
}
// Send this fragment
pn_send_to_etcp(pn);
pn_buf_renew(pn);
}
queue_dgram_free(in_dgram);
queue_entry_free(in_dgram);
continue;
}
// Cancel flush timer if active
if (pn->flush_timer) {
uasync_cancel_timeout(pn->ua, pn->flush_timer);
pn->flush_timer = NULL;
}
// Write size header (2 bytes, little-endian)
pn->data[pn->data_ptr++] = in_dgram->len & 0xFF;
pn->data[pn->data_ptr++] = (in_dgram->len >> 8) & 0xFF;
// Write packet data
memcpy(pn->data + pn->data_ptr, in_dgram->dgram, in_dgram->len);
pn->data_ptr += in_dgram->len;
queue_dgram_free(in_dgram);
queue_entry_free(in_dgram);
// Check if buffer is almost full (less than 100 bytes remaining)
if (pn->data_ptr + 100 >= pn->data_size) {
// Buffer is almost full - send immediately
pn_send_to_etcp(pn);
pn_buf_renew(pn);
} else {
// Set flush timer to send after tx_wait_time ms
if (pn->flush_timer) {
uasync_cancel_timeout(pn->ua, pn->flush_timer);
}
pn->flush_timer = uasync_set_timeout(pn->ua, pn->tx_wait_time, pn, pn_flush_cb);
}
// Set flush timer if no more input
if (queue_entry_count(pn->input) == 0) {
pn->flush_timer = uasync_set_timeout(pn->ua, pn->tx_wait_time, pn, pn_flush_cb);
}
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_packer: etcp_input_ready_cb finished, processed=%d", processed);
queue_resume_callback(pn->input);
}
@ -300,7 +219,6 @@ static void pn_flush_cb(void* arg) {
struct PKTNORM* pn = (struct PKTNORM*)arg;
if (!pn) return;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_flush_cb: called, pn=%p", pn);
pn->flush_timer = NULL;
pn_send_to_etcp(pn);
}
@ -322,75 +240,54 @@ void pn_flush(struct PKTNORM* pn) {
// Internal: Unpacker callback (assembles fragments into original packets)
static void pn_unpacker_cb(struct ll_queue* q, void* arg) {
struct PKTNORM* pn = (struct PKTNORM*)arg;
if (!pn) {
queue_resume_callback(q);
return;
}
// Process all available fragments from output_queue
struct ETCP_FRAGMENT* frag;
while ((frag = (struct ETCP_FRAGMENT*)queue_data_get(pn->etcp->output_queue)) != NULL) {
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;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_unpacker: frag len=%d, recvpart=%p, recvpart_rem=%d",
len, (void*)pn->recvpart, pn->recvpart_rem);
DEBUG_INFO(DEBUG_CATEGORY_NORMALIZER, "pn_unpacker: unpacking fragment len=%d", len);
while (ptr < len) {
if (pn->recvpart_rem==0) {
// New packet - read size header (2 bytes)
if (!pn->recvpart) {
// Need length header for new packet
if (len - ptr < 2) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "pn_unpacker: incomplete header, len=%d, ptr=%d", len, ptr);
// Incomplete header, reset
pn_unpacker_reset_state(pn);
break;
}
uint16_t pkt_len = payload[ptr] | (payload[ptr + 1] << 8);
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;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_unpacker: new packet, pkt_len=%d", pkt_len);
if (pkt_len == 0) {
// Пустой пакет - пропускаем
continue;
}
pn->recvpart = ll_alloc_lldgram(pkt_len);
pn->recvpart = ll_alloc_lldgram(part_size);
if (!pn->recvpart) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "pn_unpacker: failed to alloc recvpart");
pn_unpacker_reset_state(pn);
break;
}
pn->recvpart->len = 0;
pn->recvpart_rem = pkt_len; // Сколько байт осталось собрать
}
// Копируем данные в recvpart
uint16_t rem = pn->recvpart_rem;
uint16_t avail = len - ptr;
uint16_t rem = pn->recvpart->memlen - pn->recvpart->len;// осталось собрать байт
uint16_t avail = len - ptr;// доступно байт сейчас
uint16_t cp = (rem < avail) ? rem : avail;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_unpacker: copying cp=%d bytes (rem=%d, avail=%d)", cp, rem, avail);
if (pn->recvpart) memcpy(pn->recvpart->dgram + pn->recvpart->len, payload + ptr, cp);
DEBUG_INFO(DEBUG_CATEGORY_NORMALIZER, "pn_unpacker: 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;
pn->recvpart_rem -= cp;
ptr += cp;
// Если пакет полностью собран
if (pn->recvpart_rem == 0) {
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_unpacker: packet complete, len=%d", pn->recvpart->len);
if (pn->recvpart) {
queue_data_put(pn->output, pn->recvpart, 0);
pn->recvpart = NULL; // Сбросить указатель после передачи
}
if (pn->recvpart->len == pn->recvpart->memlen) {
queue_data_put(pn->output, pn->recvpart, 0);
pn->recvpart = NULL;
}
}
// Free the fragment using ll_queue API
queue_dgram_free(&frag->ll);
queue_entry_free(&frag->ll);
// 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);

282
src/pkt_normalizer.c.backup

@ -1,282 +0,0 @@
// pkt_normalizer.c - Implementation of packet normalizer for ETCP
#include "pkt_normalizer.h"
#include "etcp.h" // For ETCP_CONN and related structures
#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 "debug_config.h" // Assuming this for DEBUG_ERROR
// 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) {
if (!etcp) return NULL;
struct PKTNORM* pn = calloc(1, sizeof(struct PKTNORM));
if (!pn) return NULL;
pn->etcp = etcp;
pn->ua = etcp->instance->ua;
pn->frag_size = etcp->mtu - 100; // Use MTU as fixed packet size (adjust if headers need subtraction)
pn->tx_wait_time = 10;
pn->input = queue_new(pn->ua, 0); // No hash needed
pn->output = queue_new(pn->ua, 0); // No hash needed
if (!pn->input || !pn->output) {
pn_pair_deinit(pn);
return NULL;
}
queue_set_callback(pn->input, packer_cb, pn);
queue_set_callback(etcp->output_queue, pn_unpacker_cb, pn);
pn->data = NULL;
pn->recvpart = NULL;
pn->flush_timer = NULL;
return pn;
}
// Deinitialization
void pn_pair_deinit(struct PKTNORM* pn) {
if (!pn) return;
// Drain and free queues
if (pn->input) {
struct ll_entry* entry;
while ((entry = queue_data_get(pn->input)) != NULL) {
if (entry->dgram) {
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) {
free(entry->dgram);
}
queue_entry_free(entry);
}
queue_free(pn->output);
}
if (pn->flush_timer) {
uasync_cancel_timeout(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);
}
free(pn);
}
// Reset unpacker state
void pn_unpacker_reset_state(struct PKTNORM* pn) {
if (!pn) return;
if (pn->recvpart) {
queue_dgram_free(pn->recvpart);
queue_entry_free(pn->recvpart);
pn->recvpart = NULL;
}
}
// 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) {
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);
if (!entry) return;
memcpy(entry->dgram, data, len);
entry->len = len;
entry->dgram_pool = NULL;
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) {
uasync_cancel_timeout(pn->ua, pn->flush_timer);
pn->flush_timer = NULL;
}
}
// Internal: Packer callback
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);
}
// Helper to send block to ETCP as ETCP_FRAGMENT
static void pn_send_to_etcp(struct PKTNORM* pn) {
if (!pn || !pn->data || pn->data_ptr == 0) return;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "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_ETCP, "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.dgram_pool = pn->etcp->instance->data_pool;
frag->ll.len = pn->data_ptr;
frag->ll.memlen = pn->etcp->instance->data_pool->object_size;
queue_data_put(pn->etcp->input_queue, (struct ll_entry*)frag, 0);
// Сбросить структуру (dgram передан во фрагмент, не освобождаем)
pn->data = NULL;
}
// Internal: Renew sndpart buffer
static void pn_buf_renew(struct PKTNORM* pn) {
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);
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;
}
}
// 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) return;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_packer: etcp_input_ready_cb");
struct ll_entry* in_dgram = queue_data_get(pn->input);
if (!in_dgram) { queue_resume_callback(pn->input); return; }
uint16_t in_ptr = 0;//
while (in_ptr < in_dgram->len) {
pn_buf_renew(pn);
if (!pn->data) break; // Allocation failed
int remain = pn->data_size - pn->data_ptr;
if (remain<3) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "pn_packer: fatal logic error");
pn->logic_errors++;
break;
}
if (in_ptr == 0) {
pn->data[pn->data_ptr++] = in_dgram->len & 0xFF;
pn->data[pn->data_ptr++] = (in_dgram->len >> 8) & 0xFF;
remain -= 2;
}
memcpy(pn->data + pn->data_ptr, in_dgram->dgram + in_ptr, remain);
pn->data_ptr += remain;
in_ptr += remain;
}
queue_dgram_free(in_dgram);
queue_entry_free(in_dgram);
// Cancel flush timer if active
if (pn->flush_timer) {
uasync_cancel_timeout(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_set_timeout(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) {
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) {
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;
while (ptr < len) {
if (!pn->recvpart) {
// Need length header for new packet
if (len - ptr < 2) {
// Incomplete header, reset
pn_unpacker_reset_state(pn);
break;
}
uint16_t pkt_len = payload[ptr] | (payload[ptr + 1] << 8);
ptr += 2;
pn->recvpart = ll_alloc_lldgram(pkt_len);
if (!pn->recvpart) {
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;
memcpy(pn->recvpart->dgram + pn->recvpart->len, payload + ptr, cp);
pn->recvpart->len += cp;
ptr += cp;
if (pn->recvpart->len == pn->recvpart->memlen) {
queue_data_put(pn->output, pn->recvpart, 0);
pn->recvpart = NULL;
}
}
// Free the fragment using ll_queue API
queue_dgram_free((struct ll_entry*)frag);
queue_entry_free((struct ll_entry*)frag);
}
queue_resume_callback(q);
}

313
src/pkt_normalizer.c1

@ -1,313 +0,0 @@
// pkt_normalizer.c - Implementation of packet normalizer for ETCP
#include "pkt_normalizer.h"
#include "etcp.h" // For ETCP_CONN and related structures
#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 "debug_config.h" // Assuming this for DEBUG_ERROR
// 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) {
if (!etcp) return NULL;
struct PKTNORM* pn = calloc(1, sizeof(struct PKTNORM));
if (!pn) return NULL;
pn->etcp = etcp;
pn->ua = etcp->instance->ua;
pn->frag_size = etcp->mtu - 100; // Use MTU as fixed packet size (adjust if headers need subtraction)
pn->tx_wait_time = 10;
pn->input = queue_new(pn->ua, 0); // No hash needed
pn->output = queue_new(pn->ua, 0); // No hash needed
if (!pn->input || !pn->output) {
pn_pair_deinit(pn);
return NULL;
}
queue_set_callback(pn->input, packer_cb, pn);
queue_set_callback(etcp->output_queue, pn_unpacker_cb, pn);
pn->data = NULL;
pn->recvpart = NULL;
pn->recvpart_rem = 0;
pn->flush_timer = NULL;
return pn;
}
// Deinitialization
void pn_pair_deinit(struct PKTNORM* pn) {
if (!pn) return;
// Drain and free queues
if (pn->input) {
struct ll_entry* entry;
while ((entry = queue_data_get(pn->input)) != NULL) {
if (entry->dgram) {
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) {
free(entry->dgram);
}
queue_entry_free(entry);
}
queue_free(pn->output);
}
if (pn->flush_timer) {
uasync_cancel_timeout(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);
}
free(pn);
}
// Reset unpacker state
void pn_unpacker_reset_state(struct PKTNORM* pn) {
if (!pn) return;
if (pn->recvpart) {
queue_dgram_free(pn->recvpart);
queue_entry_free(pn->recvpart);
pn->recvpart = NULL;
}
pn->recvpart_rem = 0;
}
// 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) {
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);
if (!entry) return;
memcpy(entry->dgram, data, len);
entry->len = len;
entry->dgram_pool = NULL;
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) {
uasync_cancel_timeout(pn->ua, pn->flush_timer);
pn->flush_timer = NULL;
}
}
// Internal: Packer callback
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);
}
// Helper to send block to ETCP as ETCP_FRAGMENT
static void pn_send_to_etcp(struct PKTNORM* pn) {
if (!pn || !pn->data || pn->data_ptr == 0) return;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_packer: pn_send_to_etcp, len=%d", pn->data_ptr);
// 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_ETCP, "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.dgram_pool = pn->etcp->instance->data_pool;
frag->ll.len = pn->data_ptr;
frag->ll.memlen = pn->etcp->instance->data_pool->object_size;
queue_data_put(pn->etcp->input_queue, (struct ll_entry*)frag, 0);
// Сбросить структуру (dgram передан во фрагмент, не освобождаем)
pn->data = NULL;
}
// Internal: Renew sndpart buffer
static void pn_buf_renew(struct PKTNORM* pn) {
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);
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;
}
}
// 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) return;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_packer: etcp_input_ready_cb");
struct ll_entry* in_dgram = queue_data_get(pn->input);
if (!in_dgram) { queue_resume_callback(pn->input); return; }
uint16_t in_ptr = 0;//
while (in_ptr < in_dgram->len) {
pn_buf_renew(pn);
if (!pn->data) break; // Allocation failed
int remain = pn->data_size - pn->data_ptr;
if (remain<3) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "pn_packer: fatal logic error");
pn->logic_errors++;
break;
}
if (in_ptr == 0) {
pn->data[pn->data_ptr++] = in_dgram->len & 0xFF;
pn->data[pn->data_ptr++] = (in_dgram->len >> 8) & 0xFF;
remain -= 2;
}
memcpy(pn->data + pn->data_ptr, in_dgram->dgram + in_ptr, remain);
pn->data_ptr += remain;
in_ptr += remain;
}
queue_dgram_free(in_dgram);
queue_entry_free(in_dgram);
// Cancel flush timer if active
if (pn->flush_timer) {
uasync_cancel_timeout(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_set_timeout(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) {
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) {
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;
uint8_t* payload = frag->ll.dgram;
uint16_t len = frag->ll.len;
uint16_t ptr = 0;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_unpacker: frag len=%d, recvpart=%p, recvpart_rem=%d",
len, (void*)pn->recvpart, pn->recvpart_rem);
while (ptr < len) {
// Если ждем заголовок нового пакета (recvpart == NULL и recvpart_rem == 0)
if (!pn->recvpart && pn->recvpart_rem == 0) {
// Читаем длину пакета из первых 2 байт
if (len - ptr < 2) {
// Неполный заголовок - сбрасываем
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "pn_unpacker: incomplete header, len=%d, ptr=%d", len, ptr);
pn_unpacker_reset_state(pn);
break;
}
uint16_t pkt_len = payload[ptr] | (payload[ptr + 1] << 8);
ptr += 2;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_unpacker: new packet, pkt_len=%d", pkt_len);
if (pkt_len == 0) {
// Пустой пакет - пропускаем
continue;
}
// Проверка на максимальный размер пакета (10KB + 2 байта заголовка)
if (pkt_len > 10002) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "pn_unpacker: invalid pkt_len=%d, resetting", pkt_len);
pn_unpacker_reset_state(pn);
break;
}
pn->recvpart = ll_alloc_lldgram(pkt_len);
if (!pn->recvpart) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "pn_unpacker: failed to alloc recvpart");
break;
}
pn->recvpart->len = 0;
pn->recvpart_rem = pkt_len; // Сколько байт осталось собрать
}
// Копируем данные в recvpart
uint16_t rem = pn->recvpart_rem;
uint16_t avail = len - ptr;
uint16_t cp = (rem < avail) ? rem : avail;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_unpacker: copying cp=%d bytes (rem=%d, avail=%d)", cp, rem, avail);
memcpy(pn->recvpart->dgram + pn->recvpart->len, payload + ptr, cp);
pn->recvpart->len += cp;
pn->recvpart_rem -= cp;
ptr += cp;
// Если пакет полностью собран
if (pn->recvpart_rem == 0) {
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_unpacker: packet complete, len=%d", pn->recvpart->len);
queue_data_put(pn->output, pn->recvpart, 0);
pn->recvpart = NULL;
// Следующий фрагмент будет начинаться с заголовка нового пакета
}
}
// Free the fragment using ll_queue API
queue_dgram_free((struct ll_entry*)frag);
queue_entry_free((struct ll_entry*)frag);
}
queue_resume_callback(q);
}

288
src/pkt_normalizer.c2

@ -1,288 +0,0 @@
// pkt_normalizer.c - Implementation of packet normalizer for ETCP
#include "pkt_normalizer.h"
#include "etcp.h" // For ETCP_CONN and related structures
#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 "debug_config.h" // Assuming this for DEBUG_ERROR
// 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) {
if (!etcp) return NULL;
struct PKTNORM* pn = calloc(1, sizeof(struct PKTNORM));
if (!pn) return NULL;
pn->etcp = etcp;
pn->ua = etcp->instance->ua;
pn->frag_size = etcp->mtu - 100; // Use MTU as fixed packet size (adjust if headers need subtraction)
pn->tx_wait_time = 10;
pn->input = queue_new(pn->ua, 0); // No hash needed
pn->output = queue_new(pn->ua, 0); // No hash needed
if (!pn->input || !pn->output) {
pn_pair_deinit(pn);
return NULL;
}
queue_set_callback(pn->input, packer_cb, pn);
queue_set_callback(etcp->output_queue, pn_unpacker_cb, pn);
pn->data = NULL;
pn->recvpart = NULL;
pn->flush_timer = NULL;
return pn;
}
// Deinitialization
void pn_pair_deinit(struct PKTNORM* pn) {
if (!pn) return;
// Drain and free queues
if (pn->input) {
struct ll_entry* entry;
while ((entry = queue_data_get(pn->input)) != NULL) {
if (entry->dgram) {
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) {
free(entry->dgram);
}
queue_entry_free(entry);
}
queue_free(pn->output);
}
if (pn->flush_timer) {
uasync_cancel_timeout(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);
}
free(pn);
}
// Reset unpacker state
void pn_unpacker_reset_state(struct PKTNORM* pn) {
if (!pn) return;
if (pn->recvpart) {
queue_dgram_free(pn->recvpart);
queue_entry_free(pn->recvpart);
pn->recvpart = NULL;
}
}
// 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) {
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);
if (!entry) return;
memcpy(entry->dgram, data, len);
entry->len = len;
entry->dgram_pool = NULL;
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) {
uasync_cancel_timeout(pn->ua, pn->flush_timer);
pn->flush_timer = NULL;
}
}
// Internal: Packer callback
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);
}
// Helper to send block to ETCP as ETCP_FRAGMENT
static void pn_send_to_etcp(struct PKTNORM* pn) {
if (!pn || !pn->data || pn->data_ptr == 0) return;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "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_ETCP, "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.memlen = pn->etcp->instance->data_pool->object_size;
queue_data_put(pn->etcp->input_queue, (struct ll_entry*)frag, 0);
// Сбросить структуру (dgram передан во фрагмент, не освобождаем)
pn->data = NULL;
}
// Internal: Renew sndpart buffer
static void pn_buf_renew(struct PKTNORM* pn) {
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);
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;
}
}
// 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) return;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_packer: etcp_input_ready_cb");
struct ll_entry* in_dgram = queue_data_get(pn->input);
if (!in_dgram) { queue_resume_callback(pn->input); return; }
uint16_t in_ptr = 0;//
while (in_ptr < in_dgram->len) {
pn_buf_renew(pn);
if (!pn->data) break; // Allocation failed
int remain = pn->data_size - pn->data_ptr;
if (remain<3) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "pn_packer: fatal logic error");
pn->logic_errors++;
break;
}
if (in_ptr == 0) {
pn->data[pn->data_ptr++] = in_dgram->len & 0xFF;
pn->data[pn->data_ptr++] = (in_dgram->len >> 8) & 0xFF;
remain -= 2;
}
memcpy(pn->data + pn->data_ptr, in_dgram->dgram + in_ptr, remain);
pn->data_ptr += remain;
in_ptr += remain;
}
queue_dgram_free(in_dgram);
queue_entry_free(in_dgram);
// Cancel flush timer if active
if (pn->flush_timer) {
uasync_cancel_timeout(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_set_timeout(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) {
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) {
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;
while (ptr < len) {
if (!pn->recvpart) {
// Need length header for new packet
if (len - ptr < 2) {
// Incomplete header, reset
pn_unpacker_reset_state(pn);
break;
}
uint16_t pkt_len = payload[ptr] | (payload[ptr + 1] << 8);
ptr += 2;
pn->recvpart = ll_alloc_lldgram(pkt_len);
if (!pn->recvpart) {
break;
}
pn->recvpart->len = 0;
} else {
// We are in the middle of assembling a packet
// Skip the 2-byte length header at the start of subsequent fragments
if (ptr == 0) {
ptr += 2;
if (ptr >= len) break; // Fragment contains only the length header
}
}
uint16_t rem = pn->recvpart->memlen - pn->recvpart->len;
uint16_t avail = len - ptr;
uint16_t cp = (rem < avail) ? rem : avail;
memcpy(pn->recvpart->dgram + pn->recvpart->len, payload + ptr, cp);
pn->recvpart->len += cp;
ptr += cp;
if (pn->recvpart->len == pn->recvpart->memlen) {
queue_data_put(pn->output, pn->recvpart, 0);
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);
}

295
src/pkt_normalizer.c3

@ -1,295 +0,0 @@
// pkt_normalizer.c - Implementation of packet normalizer for ETCP
#include "pkt_normalizer.h"
#include "etcp.h" // For ETCP_CONN and related structures
#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 "debug_config.h" // Assuming this for DEBUG_ERROR
// 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) {
if (!etcp) return NULL;
struct PKTNORM* pn = calloc(1, sizeof(struct PKTNORM));
if (!pn) return NULL;
pn->etcp = etcp;
pn->ua = etcp->instance->ua;
pn->frag_size = etcp->mtu - 100; // Use MTU as fixed packet size (adjust if needed)
pn->tx_wait_time = 10;
pn->input = queue_new(pn->ua, 0); // No hash needed
pn->output = queue_new(pn->ua, 0); // No hash needed
if (!pn->input || !pn->output) {
pn_pair_deinit(pn);
return NULL;
}
queue_set_callback(pn->input, packer_cb, pn);
queue_set_callback(etcp->output_queue, pn_unpacker_cb, pn);
pn->data = NULL;
pn->recvpart = NULL;
pn->recvpart_rem = 0;
pn->flush_timer = NULL;
return pn;
}
// Deinitialization
void pn_pair_deinit(struct PKTNORM* pn) {
if (!pn) return;
// Drain and free queues
if (pn->input) {
struct ll_entry* entry;
while ((entry = queue_data_get(pn->input)) != NULL) {
if (entry->dgram) {
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) {
free(entry->dgram);
}
queue_entry_free(entry);
}
queue_free(pn->output);
}
if (pn->flush_timer) {
uasync_cancel_timeout(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);
}
free(pn);
}
// Reset unpacker state
void pn_unpacker_reset_state(struct PKTNORM* pn) {
if (!pn) return;
if (pn->recvpart) {
queue_dgram_free(pn->recvpart);
queue_entry_free(pn->recvpart);
pn->recvpart = NULL;
}
pn->recvpart_rem = 0;
}
// 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) {
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);
if (!entry) return;
memcpy(entry->dgram, data, len);
entry->len = len;
entry->dgram_pool = NULL;
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) {
uasync_cancel_timeout(pn->ua, pn->flush_timer);
pn->flush_timer = NULL;
}
}
// Internal: Packer callback
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);
}
// Helper to send block to ETCP as ETCP_FRAGMENT
static void pn_send_to_etcp(struct PKTNORM* pn) {
if (!pn || !pn->data || pn->data_ptr == 0) return;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "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_ETCP, "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.dgram_pool = pn->etcp->instance->data_pool;
frag->ll.len = pn->data_ptr;
frag->ll.memlen = pn->etcp->instance->data_pool->object_size;
queue_data_put(pn->etcp->input_queue, (struct ll_entry*)frag, 0);
// Сбросить структуру (dgram передан во фрагмент, не освобождаем)
pn->data = NULL;
}
// Internal: Renew buffer for packer
static void pn_buf_renew(struct PKTNORM* pn) {
if (pn->data && pn->data_ptr > 0) {
pn_send_to_etcp(pn);
}
if (!pn->data) {
pn->data = memory_pool_alloc(pn->etcp->instance->data_pool);
size_t 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;
}
}
// 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) return;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_packer: etcp_input_ready_cb");
struct ll_entry* in_dgram = queue_data_get(pn->input);
if (!in_dgram) { queue_resume_callback(pn->input); return; }
uint16_t in_ptr = 0;
while (in_ptr < in_dgram->len) {
pn_buf_renew(pn);
if (!pn->data) break; // Allocation failed
int remain = pn->data_size - pn->data_ptr;
if (remain < 3) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "pn_packer: fatal logic error");
pn->logic_errors++;
break;
}
if (in_ptr == 0) {
pn->data[pn->data_ptr++] = in_dgram->len & 0xFF;
pn->data[pn->data_ptr++] = (in_dgram->len >> 8) & 0xFF;
remain -= 2;
}
memcpy(pn->data + pn->data_ptr, in_dgram->dgram + in_ptr, remain);
pn->data_ptr += remain;
in_ptr += remain;
}
queue_dgram_free(in_dgram);
queue_entry_free(in_dgram);
// Cancel flush timer if active
if (pn->flush_timer) {
uasync_cancel_timeout(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_set_timeout(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) {
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) {
struct PKTNORM* pn = (struct PKTNORM*)arg;
if (!pn) return;
struct ETCP_FRAGMENT* frag = (struct ETCP_FRAGMENT*)queue_data_get(pn->etcp->output_queue);
if (!frag) break;
uint8_t* payload = frag->ll.dgram;
uint16_t len = frag->ll.len;
uint16_t ptr = 0;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_unpacker: frag len=%d, recvpart=%p, recvpart_rem=%d",
len, (void*)pn->recvpart, pn->recvpart_rem);
while (ptr < len) {
if (pn->recvpart_rem==0) {
if (len - ptr < 2) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "pn_unpacker: incomplete header, len=%d, ptr=%d", len, ptr);
pn_unpacker_reset_state(pn);
break;
}
uint16_t pkt_len = payload[ptr] | (payload[ptr + 1] << 8);
ptr += 2;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_unpacker: new packet, pkt_len=%d", pkt_len);
if (pkt_len == 0) {
// Пустой пакет - пропускаем
continue;
}
pn->recvpart = ll_alloc_lldgram(pkt_len);
if (!pn->recvpart) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "pn_unpacker: failed to alloc recvpart");
}
pn->recvpart->len = 0;
pn->recvpart_rem = pkt_len; // Сколько байт осталось собрать
}
// Копируем данные в recvpart
uint16_t rem = pn->recvpart_rem;
uint16_t avail = len - ptr;
uint16_t cp = (rem < avail) ? rem : avail;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_unpacker: copying cp=%d bytes (rem=%d, avail=%d)", cp, rem, avail);
if (pn->recvpart) memcpy(pn->recvpart->dgram + pn->recvpart->len, payload + ptr, cp);
pn->recvpart->len += cp;
pn->recvpart_rem -= cp;
ptr += cp;
// Если пакет полностью собран
if (pn->recvpart_rem == 0) {
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "pn_unpacker: packet complete, len=%d", pn->recvpart->len);
if (pn->recvpart) queue_data_put(pn->output, pn->recvpart, 0);
}
}
// Free the fragment using ll_queue API
queue_dgram_free(&frag->ll);
queue_entry_free(&frag->ll);
queue_resume_callback(q);
}

2
src/pkt_normalizer.h

@ -29,7 +29,7 @@ struct PKTNORM {
// unpacker:
struct ll_entry* recvpart; // блок ожидающий заполнение
uint16_t recvpart_rem; // сколько байт осталось собрать (0 = ждём заголовок нового пакета)
// uint16_t recvpart_rem; // сколько байт осталось собрать (0 = ждём заголовок нового пакета)
// stats:
uint32_t alloc_errors;

Loading…
Cancel
Save