Browse Source

Refactor: Add dgram fields to ll_queue, remove ref_count, integrate pkt_normalizer

- ll_queue: Add dgram, dgram_free_fn, dgram_pool, len fields to ll_entry
- ll_queue: Remove ref_count, simplify memory management
- etcp: Add normalizer pointer to ETCP_CONN struct
- pkt_normalizer: Major refactoring for queue integration
- tests: Update for new ll_entry structure
- cleanup: Remove backup files
nodeinfo-routing-update
Evgeny 8 months ago
parent
commit
dd4644776d
  1. 22
      lib/ll_queue.c
  2. 9
      lib/ll_queue.h
  3. 1496
      src/etcp.c
  4. 12
      src/etcp.h
  5. 873
      src/pkt_normalizer.c
  6. 104
      src/pkt_normalizer.h
  7. 369
      src/secure_channel.c.bak
  8. 388
      src/secure_channel.c1
  9. 68
      src/secure_channel.h1
  10. 202
      tests/test_etcp_100_packets.c
  11. 252
      tests/test_etcp_simple_traffic.c

22
lib/ll_queue.c

@ -114,7 +114,7 @@ void* queue_data_new(size_t data_size) {
memset(entry, 0, sizeof(struct ll_entry) + data_size); memset(entry, 0, sizeof(struct ll_entry) + data_size);
entry->size = data_size; entry->size = data_size;
entry->ref_count = 1; // Единственная ссылка - на сам элемент entry->len = 0;
entry->pool = NULL; // Выделено через malloc entry->pool = NULL; // Выделено через malloc
// DEBUG_DEBUG(DEBUG_CATEGORY_LL_QUEUE, "queue_data_new: created entry %p, size=%zu", entry, data_size); // DEBUG_DEBUG(DEBUG_CATEGORY_LL_QUEUE, "queue_data_new: created entry %p, size=%zu", entry, data_size);
@ -130,7 +130,7 @@ void* queue_data_new_from_pool(struct memory_pool* pool) {
memset(entry, 0, pool->object_size); memset(entry, 0, pool->object_size);
entry->size = pool->object_size - sizeof(struct ll_entry); entry->size = pool->object_size - sizeof(struct ll_entry);
entry->ref_count = 1; // Единственная ссылка - на сам элемент entry->len = 0;
entry->pool = pool; // Выделено из пула entry->pool = pool; // Выделено из пула
// DEBUG_DEBUG(DEBUG_CATEGORY_LL_QUEUE, "queue_data_new_from_pool: created entry %p from pool %p", entry, pool); // DEBUG_DEBUG(DEBUG_CATEGORY_LL_QUEUE, "queue_data_new_from_pool: created entry %p from pool %p", entry, pool);
@ -143,14 +143,10 @@ void queue_data_free(void* data) {
struct ll_entry* entry = data_to_entry(data); struct ll_entry* entry = data_to_entry(data);
// DEBUG_DEBUG(DEBUG_CATEGORY_LL_QUEUE, "queue_data_free: freeing entry %p, ref_count=%d", entry, entry->ref_count); if (entry->pool) {
memory_pool_free(entry->pool, entry);
if (--entry->ref_count <= 0) { } else {
if (entry->pool) { free(entry);
memory_pool_free(entry->pool, entry);
} else {
free(entry);
}
} }
} }
@ -226,7 +222,6 @@ int queue_data_put(struct ll_queue* q, void* data, uint32_t id) {
q->count++; q->count++;
q->total_bytes += entry->size; q->total_bytes += entry->size;
// ВАЖНО: НЕ увеличиваем ref_count - элемент просто находится в очереди
size_t send_q_bytes = queue_total_bytes(q); size_t send_q_bytes = queue_total_bytes(q);
// DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "check total bytes: new_q_len=%d element_size:%d", send_q_bytes, entry->size); // DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "check total bytes: new_q_len=%d element_size:%d", send_q_bytes, entry->size);
@ -270,7 +265,6 @@ int queue_data_put_first(struct ll_queue* q, void* data, uint32_t id) {
q->count++; q->count++;
q->total_bytes += entry->size; q->total_bytes += entry->size;
// ВАЖНО: НЕ увеличиваем ref_count - элемент просто находится в очереди
add_to_hash(q, entry); add_to_hash(q, entry);
@ -304,8 +298,6 @@ void* queue_data_get(struct ll_queue* q) {
remove_from_hash(q, entry); remove_from_hash(q, entry);
// ВАЖНО: НЕ уменьшаем ref_count - просто удаляем из очереди
// DEBUG_DEBUG(DEBUG_CATEGORY_LL_QUEUE, "queue_data_get: got entry %p (id=%u), count=%d", entry, entry->id, q->count); // DEBUG_DEBUG(DEBUG_CATEGORY_LL_QUEUE, "queue_data_get: got entry %p (id=%u), count=%d", entry, entry->id, q->count);
// Приостановить коллбэки для предотвращения рекурсии // Приостановить коллбэки для предотвращения рекурсии
@ -390,8 +382,6 @@ int queue_remove_data(struct ll_queue* q, void* data) {
q->count--; q->count--;
q->total_bytes -= entry->size; q->total_bytes -= entry->size;
// ВАЖНО: НЕ уменьшаем ref_count - просто удаляем из очереди
entry->next = NULL; entry->next = NULL;
entry->prev = NULL; entry->prev = NULL;

9
lib/ll_queue.h

@ -22,8 +22,13 @@ typedef void (*queue_callback_fn)(struct ll_queue* q, void* arg);
struct ll_entry { struct ll_entry {
struct ll_entry* next; // Указатель на следующий элемент в очереди struct ll_entry* next; // Указатель на следующий элемент в очереди
struct ll_entry* prev; // Указатель на предыдущий элемент в очереди struct ll_entry* prev; // Указатель на предыдущий элемент в очереди
size_t size; // Размер данных элемента (байт) uint16_t size; // Размер доступной памяти после блока ll_entry - т.е. data[size]. используется для добавления доп. параметров
int ref_count; // Счетчик ссылок (1 при создании, не изменяется при помещении/изъятии из очереди)
uint16_t len; // размер пакета
uint8_t* dgram; // данные пакета
void (*dgram_free_fn)(uint8_t* data, void* arg); // функция освобождения блока
struct memory_pool* dgram_pool; // Пул, из которого выделен этот элемент (NULL, если выделен через malloc)
uint32_t id; // Идентификатор для хеш-поиска uint32_t id; // Идентификатор для хеш-поиска
struct ll_entry* hash_next; // Следующий в хеш-цепочке struct ll_entry* hash_next; // Следующий в хеш-цепочке
struct memory_pool* pool; // Пул, из которого выделен этот элемент (NULL, если выделен через malloc) struct memory_pool* pool; // Пул, из которого выделен этот элемент (NULL, если выделен через malloc)

1496
src/etcp.c

File diff suppressed because it is too large Load Diff

12
src/etcp.h

@ -11,6 +11,11 @@
extern "C" { extern "C" {
#endif #endif
#include "pkt_normalizer.h"
// In struct ETCP_CONN, add:
//struct pn_pair* normalizer;
//!!!!!!!!!!! надо переделать ll_queue чтобы возвращала не дата а свою структуру //!!!!!!!!!!! надо переделать ll_queue чтобы возвращала не дата а свою структуру
// Forward declarations // Forward declarations
@ -34,7 +39,7 @@ uint64_t get_current_time_units(void);
// в этот список пакет добавляется когда перемещается из input_queue в input_send_q, при этом к пакету добавляется struct INFLIGHT_PACKET из inflight_pool. // в этот список пакет добавляется когда перемещается из input_queue в input_send_q, при этом к пакету добавляется struct INFLIGHT_PACKET из inflight_pool.
// пакет полностью удаляется когда приходит ACK (либо conn_reset/close) // пакет полностью удаляется когда приходит ACK (либо conn_reset/close)
struct INFLIGHT_PACKET { struct INFLIGHT_PACKET {// выделяется из etcp->inflight_pool
struct ll_entry ll; struct ll_entry ll;
struct ETCP_LINK* last_link; // Last sent link struct ETCP_LINK* last_link; // Last sent link
uint64_t last_timestamp; // Last send timestamp uint64_t last_timestamp; // Last send timestamp
@ -42,15 +47,13 @@ struct INFLIGHT_PACKET {
uint8_t send_count; // Number of sends uint8_t send_count; // Number of sends
uint8_t retrans_req_count; // Number of retrans requests uint8_t retrans_req_count; // Number of retrans requests
uint8_t state; // WAIT_ACK or WAIT_SEND uint8_t state; // WAIT_ACK or WAIT_SEND
uint8_t* pkt_data;
}; };
// Список пакетов для сборки. собирается в ll_queue (используем быстрый поиск с хешем) // Список пакетов для сборки. собирается в ll_queue (используем быстрый поиск с хешем)
struct ETCP_FRAGMENT { struct ETCP_FRAGMENT {// выделяется из пула etcp->rx_pool
struct ll_entry ll; struct ll_entry ll;
uint32_t seq; uint32_t seq;
uint16_t timestamp; uint16_t timestamp;
uint8_t* pkt_data;
}; };
struct ACK_PACKET { struct ACK_PACKET {
@ -72,6 +75,7 @@ struct ETCP_CONN {
// Crypto and state // Crypto and state
struct secure_channel crypto_ctx; struct secure_channel crypto_ctx;
struct pn_pair* normalizer;
// Peer info // Peer info
uint64_t peer_node_id; // Peer node ID uint64_t peer_node_id; // Peer node ID

873
src/pkt_normalizer.c

@ -1,606 +1,355 @@
// pkt_normalizer.c - Implementation of packet normalizer for ETCP
#include "pkt_normalizer.h" #include "pkt_normalizer.h"
#include "../lib/u_async.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 <stdlib.h>
#include <string.h> #include <string.h>
#include <stdint.h> #include <stdio.h> // For debugging (can be removed if not needed)
#include <stdio.h>
static void packer_handler(struct ll_queue* q, void* arg); // Internal helper to convert void* data to struct ll_entry*
static void unpacker_handler(struct ll_queue* q, void* arg); static inline struct ll_entry* data_to_entry(void* data) {
static void send_buf(struct pn_struct* pn); if (!data) return NULL;
static int get_header(uint8_t* header, size_t L); return (struct ll_entry*)data;
/* Calculate fragment size from mtu */ }
struct pn_struct* pkt_normalizer_init(uasync_t* ua, int is_packer, int mtu) {
struct pn_struct* pn = malloc(sizeof(struct pn_struct)); // Forward declarations
static void packer_cb(struct ll_queue* q, void* arg);
static void flush_cb(void* arg);
static void input_ready_cb(struct ll_queue* q, void* arg);
// 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; if (!pn) return NULL;
pn->ua = ua;
pn->input = queue_new(ua, 0); // No memory pool for now, no hash table pn->etcp = etcp;
if (!pn->input) { pn->ua = etcp->instance->ua;
free(pn); pn->pkt_size = etcp->mtu; // Use MTU as fixed packet size (adjust if headers need subtraction)
return NULL;
} pn->input = queue_new(pn->ua, 0); // No hash needed
pn->output = queue_new(ua, 0); // No memory pool for now, no hash table pn->output = queue_new(pn->ua, 0); // No hash needed
if (!pn->output) { pn->pending = queue_new(pn->ua, 0);
queue_free(pn->input); pn->send_pending = queue_new(pn->ua, 0);
free(pn);
if (!pn->input || !pn->output || !pn->pending || !pn->send_pending) {
pn_pair_deinit(pn);
return NULL; return NULL;
} }
pn->is_packer = is_packer;
if (is_packer) { pn->in_buf = calloc(1, pn->pkt_size);
// Calculate fragment size from mtu: fragment = mtu - 100 pn->cap = pn->pkt_size; // For packer buffer (fixed size)
int fragment_size = mtu - ETCP_OVERHEAD; pn->len = 0; // Current length in packer buffer
if (fragment_size < 256) fragment_size = 256; // Minimum sane value
pn->u.packer.cap = fragment_size; pn->out_buf = calloc(1, pn->pkt_size * 10); // Arbitrary large buffer for assembly
pn->u.packer.buf = malloc(pn->u.packer.cap); pn->cap = pn->pkt_size * 10; // Initial capacity for unpacker (can reallocate if needed)
if (!pn->u.packer.buf) { pn->len = 0;
queue_free(pn->input); pn->total_len = 0;
queue_free(pn->output); pn->in_fragment = 0;
free(pn);
return NULL; // Set callback for automatic processing when items are added to input
} queue_set_callback(pn->input, packer_cb, pn);
pn->u.packer.len = 0;
pn->u.packer.error_count = 0; pn->flush_timer = NULL;
queue_set_callback(pn->input, packer_handler, pn);
} else {
pn->u.unpacker.buf = NULL;
pn->u.unpacker.len = 0;
pn->u.unpacker.total_len = 0;
pn->u.unpacker.cap = 0;
pn->u.unpacker.error_count = 0;
queue_set_callback(pn->input, unpacker_handler, pn);
}
return pn; return pn;
} }
void pkt_normalizer_deinit(struct pn_struct* pn) { // Deinitialization
void pn_pair_deinit(struct PKTNORM* pn) {
if (!pn) return; if (!pn) return;
queue_free(pn->input);
queue_free(pn->output); // Drain and free queues
if (pn->is_packer) { if (pn->input) {
free(pn->u.packer.buf); void* data;
} else { while ((data = queue_data_get(pn->input)) != NULL) {
free(pn->u.unpacker.buf); queue_data_free(data);
free(pn->u.unpacker.service_buf); }
queue_free(pn->input);
} }
free(pn); if (pn->output) {
} void* data;
struct pkt_normalizer_pair* pkt_normalizer_pair_init(uasync_t* ua, int mtu) { while ((data = queue_data_get(pn->output)) != NULL) {
struct pkt_normalizer_pair* pair = malloc(sizeof(struct pkt_normalizer_pair)); queue_data_free(data);
if (!pair) return NULL; }
pair->packer = pkt_normalizer_init(ua, 1, mtu); queue_free(pn->output);
if (!pair->packer) {
free(pair);
return NULL;
} }
pair->unpacker = pkt_normalizer_init(ua, 0, mtu); if (pn->pending) {
if (!pair->unpacker) { void* data;
pkt_normalizer_deinit(pair->packer); while ((data = queue_data_get(pn->pending)) != NULL) {
free(pair); queue_data_free(data);
return NULL; }
queue_free(pn->pending);
} }
return pair; if (pn->send_pending) {
} void* data;
void pkt_normalizer_pair_deinit(struct pkt_normalizer_pair* pair) { while ((data = queue_data_get(pn->send_pending)) != NULL) {
if (!pair) return; struct ETCP_FRAGMENT* frag = data_to_entry(data);
pkt_normalizer_deinit(pair->packer); if (frag->ll.dgram) memory_pool_free(pn->etcp->instance->data_pool, frag->ll.dgram);
pkt_normalizer_deinit(pair->unpacker); memory_pool_free(pn->etcp->rx_pool, frag);
free(pair); }
} queue_free(pn->send_pending);
static int get_header(uint8_t* header, size_t L) {
if (L > 1535) return -1;
if (L <= 239) {
header[0] = (uint8_t)L;
return 1;
} else {
uint8_t high = (uint8_t)(L >> 8);
if (high > 5) return -1;
header[0] = 0xF0 + high;
header[1] = (uint8_t)(L & 0xFF);
return 2;
} }
}
/* Сбросить состояние сборки фрагментов */ if (pn->flush_timer) {
static void reset_fragment_state(struct pn_struct* pn) { uasync_cancel_timeout(pn->ua, pn->flush_timer);
if (!pn->is_packer) {
pn->u.unpacker.len = 0;
pn->u.unpacker.total_len = 0;
pn->u.unpacker.in_fragment = 0;
} }
free(pn->in_buf);
free(pn->out_buf);
free(pn);
} }
/* Таймаут для сборки фрагментов */
static void send_buf(struct pn_struct* pn) { // Reset unpacker state
if (pn->u.packer.len == 0) return; void pn_unpacker_reset_state(struct PKTNORM* pn) {
size_t payload_len = pn->u.packer.len; if (!pn) return;
uint8_t* out = queue_data_new(2 + payload_len); // Новый API pn->len = 0;
if (!out) return; pn->total_len = 0;
*(uint16_t*)out = (uint16_t)payload_len; pn->in_fragment = 0;
memcpy(out + 2, pn->u.packer.buf, payload_len);
queue_data_put(pn->output, out, 0); // Новый API, ID=0 для внутренних пакетов
pn->u.packer.len = 0;
} }
static void packer_handler(struct ll_queue* q, void* arg) {
struct pn_struct* pn = arg; // Send data to packer (copies and adds to input queue or pending, triggering callback)
size_t max = (size_t)1400; void pn_packer_send(struct PKTNORM* pn, uint8_t* data, uint16_t len) {
if (!pn || !data || len == 0) return;
// Получить данные из очереди
void* pkt_data = queue_data_get(q); void* entry_data = queue_data_new(len);
if (!pkt_data) { if (!entry_data) return;
queue_resume_callback(q);
return; struct ll_entry* entry = data_to_entry(entry_data);
} memcpy(entry->data, data, len);
entry->len = len;
// Для определения размера нужно знать структуру, но в новом API размер в самих данных entry->size = len;
// В pkt_normalizer все данные имеют 2-байтовый заголовок с размером
uint16_t* size_ptr = (uint16_t*)pkt_data; // To minimize delay, add to input only if empty, else to pending and set waiter if not set
size_t L = *size_ptr; if (queue_entry_count(pn->input) == 0 && pn->input->waiter.callback == NULL) {
uint8_t* pkt_payload = (uint8_t*)pkt_data + 2; queue_data_put(pn->input, entry_data, 0);
uint8_t header[2];
int hsize = get_header(header, L);
size_t needed = (size_t)hsize + L;
if (hsize < 0 || needed > max) {
// Fragment
if (pn->u.packer.len > 0) {
send_buf(pn);
}
size_t remaining = L;
size_t pos = 0;
int fragment_count = 0;
while (remaining > 0) {
size_t chunk;
size_t payload_len;
uint8_t* fout;
uint8_t* fd;
uint8_t frag_header[2];
int frag_hsize;
if (fragment_count == 0) {
// Первый фрагмент: FF + общая длина (2 байта)
chunk = remaining > (max - 5) ? (max - 5) : remaining; // 2+1+2+chunk <= max
payload_len = 1 + 2 + chunk; // FF + total_len + data
fout = queue_data_new(2 + payload_len);
if (!fout) {
break;
}
fd = fout;
*(uint16_t*)fd = (uint16_t)payload_len;
fd += 2;
*fd++ = 0xFF;
*fd++ = (uint8_t)(L >> 8); // старший байт общей длины
*fd++ = (uint8_t)(L & 0xFF); // младший байт общей длины
} else {
// Не первый фрагмент
if (remaining <= max - 3) {
// Это последний возможный фрагмент (помещается в один пакет с префиксом FE)
// Пытаемся отправить как обычный блок
frag_hsize = get_header(frag_header, remaining);
if (frag_hsize > 0 && (size_t)frag_hsize + remaining + 2 <= max) {
// Успешно: обычный блок
payload_len = frag_hsize + remaining;
chunk = remaining;
fout = queue_data_new(2 + payload_len);
if (!fout) {
break;
}
fd = fout;
*(uint16_t*)fd = (uint16_t)payload_len;
fd += 2;
memcpy(fd, frag_header, frag_hsize);
fd += frag_hsize;
} else {
// Не удалось отправить как обычный блок - разбиваем на 2 фрагмента
// 1. FE фрагмент с частью данных
// 2. Обычный блок с оставшимися данными
// Находим максимальный размер для FE фрагмента
size_t max_fe_data = max - 3; // 2 байта длины + 0xFE
if (max_fe_data > remaining) {
max_fe_data = remaining;
}
// Пробуем различные размеры, начиная с максимального
size_t fe_data_size = 0;
for (size_t try_fe = max_fe_data; try_fe > 0; try_fe--) {
size_t try_regular = remaining - try_fe;
if (try_regular == 0) continue; // Нужно отправить что-то как обычный блок
uint8_t test_header[2];
int hsize = get_header(test_header, try_regular);
if (hsize <= 0) continue;
if ((size_t)hsize + try_regular + 2 <= max) {
fe_data_size = try_fe;
break;
}
}
if (fe_data_size == 0) {
// Не удалось найти разбиение - ошибка
pn->u.packer.error_count++;
// Отправляем как FE (нарушение спецификации, но это крайний случай)
chunk = remaining > (max - 3) ? (max - 3) : remaining;
payload_len = 1 + chunk; // FE + data
fout = queue_data_new(2 + payload_len);
if (!fout) break;
fd = (fout);
*(uint16_t*)fd = (uint16_t)payload_len;
fd += 2;
*fd++ = 0xFE;
} else {
// Отправляем FE фрагмент
chunk = fe_data_size;
payload_len = 1 + chunk; // FE + data
fout = queue_data_new(2 + payload_len);
if (!fout) break;
fd = (fout);
*(uint16_t*)fd = (uint16_t)payload_len;
fd += 2;
*fd++ = 0xFE;
memcpy(fd, pkt_data + pos, chunk);
queue_data_put(pn->output, fout, 0);
pos += chunk;
remaining -= chunk;
fragment_count++;
// Обновляем оставшиеся данные для обычного блока
// (цикл продолжит обработку на следующей итерации)
continue;
}
}
} else {
// Промежуточный фрагмент, отправляем как FE
chunk = remaining > (max - 3) ? (max - 3) : remaining;
payload_len = 1 + chunk; // FE + data
fout = queue_data_new(2 + payload_len);
if (!fout) {
break;
}
fd = fout;
*(uint16_t*)fd = (uint16_t)payload_len;
fd += 2;
*fd++ = 0xFE;
}
}
memcpy(fd, pkt_data + pos, chunk);
queue_data_put(pn->output, fout, 0);
pos += chunk;
remaining -= chunk;
fragment_count++;
}
} else { } else {
if (pn->u.packer.len + needed > max) { queue_data_put(pn->pending, entry_data, 0);
send_buf(pn); if (pn->input->waiter.callback == NULL) {
queue_wait_threshold(pn->input, 0, 0, input_ready_cb, pn);
} }
// Add to buffer
uint8_t* p = pn->u.packer.buf + pn->u.packer.len;
memcpy(p, header, (size_t)hsize);
memcpy(p + hsize, pkt_data, L);
pn->u.packer.len += needed;
} }
queue_data_free(pkt_data);
if (pn->u.packer.len > 0) {
send_buf(pn);
}
queue_resume_callback(q);
}
void pkt_normalizer_set_service_callback(struct pn_struct* pn, pkt_normalizer_service_callback_t callback, void* user_data) {
if (!pn) return;
pn->service_callback = callback;
pn->service_callback_user_data = user_data;
} }
void pkt_normalizer_reset_service_state(struct pn_struct* pn) {
if (!pn || pn->is_packer) return; // Internal: Callback when input queue becomes empty, move next from pending
if (pn->u.unpacker.in_service) { static void input_ready_cb(struct ll_queue* q, void* arg) {
// Deliver pending service packet struct PKTNORM* pn = (struct PKTNORM*)arg;
if (pn->service_callback) {
pn->service_callback(pn->service_callback_user_data, void* data = queue_data_get(pn->pending);
pn->u.unpacker.service_type, if (data) {
pn->u.unpacker.service_buf, queue_data_put(q, data, 0); // This will trigger packer_cb async if needed
pn->u.unpacker.service_len); if (queue_entry_count(pn->pending) > 0) {
queue_wait_threshold(q, 0, 0, input_ready_cb, pn);
} }
free(pn->u.unpacker.service_buf);
pn->u.unpacker.service_buf = NULL;
pn->u.unpacker.service_len = 0;
pn->u.unpacker.service_cap = 0;
pn->u.unpacker.in_service = 0;
} }
} }
void pkt_normalizer_reset_state(struct pn_struct* pn) {
if (!pn) return; static void etcp_input_ready_cb(struct ll_queue* q, void* arg) {
if (pn->is_packer) { struct PKTNORM* pn = (struct PKTNORM*)arg;
// Flush packer buffer
if (pn->u.packer.len > 0) { void* data = queue_data_get(pn->send_pending);
send_buf(pn); if (data) {
queue_data_put(q, data, 0);
if (queue_entry_count(pn->send_pending) > 0) {
queue_wait_threshold(q, 0, 0, etcp_input_ready_cb, pn);
} }
}
}
// Internal: Send a fixed-size chunk to etcp->input_queue (or send_pending)
static void send_chunk(struct PKTNORM* pn, uint8_t* chunk_data, uint16_t chunk_len) {
// Alloc ETCP_FRAGMENT for the chunk (as payload for etcp)
struct ETCP_FRAGMENT* frag = memory_pool_alloc(pn->etcp->rx_pool);
if (!frag) return;
frag->ll.dgram = memory_pool_alloc(pn->etcp->instance->data_pool);
if (!frag->ll.dgram) {
memory_pool_free(pn->etcp->rx_pool, frag);
return;
}
memcpy(frag->ll.dgram, chunk_data, chunk_len);
frag->ll.size = chunk_len;
frag->ll.len = chunk_len;
frag->seq = 0; // Not used for input
frag->timestamp = 0;
struct ll_queue* eq = pn->etcp->input_queue;
// Add to etcp input_queue only if empty, else to send_pending and set waiter
if (queue_entry_count(eq) == 0 && eq->waiter.callback == NULL) {
queue_data_put(eq, frag, 0);
} else { } else {
// Reset unpacker fragment state queue_data_put(pn->send_pending, frag, 0);
reset_fragment_state(pn); if (eq->waiter.callback == NULL) {
// Reset service state queue_wait_threshold(eq, 0, 0, etcp_input_ready_cb, pn);
pkt_normalizer_reset_service_state(pn); }
} }
} }
int pkt_normalizer_send_service(struct pn_struct* pn, uint8_t type, const void* data, size_t len) {
if (!pn || !pn->is_packer) return -1; // Internal: Packer callback (aggregates small, fragments large, sends chunks)
// Service packet ограничен 256 байтами всего static void packer_cb(struct ll_queue* q, void* arg) {
if (len > 256 - 2) return -1; // 2 байта на заголовок (0xFC + тип) struct PKTNORM* pn = (struct PKTNORM*)arg;
size_t max = (size_t)1400; if (!pn) return;
if (max < 3) return -1;
// Размер сервисного пакета: 2 байта длины + 1 байт 0xFC + 1 байт тип + данные pn->len = 0; // Start with empty buffer for this batch
// Если не помещается в один фрагмент - используем продолжение 0xFD
size_t total_service_len = 1 + 1 + len; // 0xFC + type + data while (q->head) {
size_t pos = 0; void* data = queue_data_get(q);
while (total_service_len > 0) { struct ll_entry* entry = data_to_entry(data);
// Определяем размер куска для этого пакета uint16_t item_len = entry->len;
size_t chunk;
uint8_t service_header; if (item_len + 2 > pn->pkt_size) {
if (pos == 0) { // Large item: fragment (must be alone in buffer)
// Первый пакет: 0xFC + тип + часть данных if (pn->len > 0) {
// Максимум данных в первом пакете: max - 2 (длина) - 2 (0xFC+тип) // Flush current aggregated small items first
size_t max_first_data = max - 4; memset(pn->in_buf + pn->len, 0, pn->pkt_size - pn->len); // Pad
if (max_first_data > len) max_first_data = len; send_chunk(pn, pn->in_buf, pn->pkt_size);
chunk = max_first_data; pn->len = 0;
service_header = 0xFC; }
// Fragment the large item
uint8_t* ptr = entry->data;
uint32_t remaining = item_len;
// First fragment
*(uint16_t*)(pn->in_buf) = 0xffff;
*(uint32_t*)(pn->in_buf + 2) = remaining;
uint16_t chunk = pn->pkt_size - 6;
memcpy(pn->in_buf + 6, ptr, chunk);
send_chunk(pn, pn->in_buf, pn->pkt_size);
ptr += chunk;
remaining -= chunk;
// Continuation fragments
while (remaining > 0) {
chunk = (remaining < pn->pkt_size - 2) ? remaining : pn->pkt_size - 2;
*(uint16_t*)(pn->in_buf) = 0x0000; // Continuation marker
memcpy(pn->in_buf + 2, ptr, chunk);
memset(pn->in_buf + 2 + chunk, 0, pn->pkt_size - 2 - chunk); // Pad
send_chunk(pn, pn->in_buf, pn->pkt_size);
ptr += chunk;
remaining -= chunk;
}
} else { } else {
// Продолжение: 0xFD + данные // Small item: try to aggregate
// Максимум данных: max - 2 (длина) - 1 (0xFD) if (pn->len + item_len + 2 > pn->pkt_size) {
size_t max_cont_data = max - 3; // Flush current buffer
size_t remaining = len - pos; memset(pn->in_buf + pn->len, 0, pn->pkt_size - pn->len); // Pad
if (max_cont_data > remaining) max_cont_data = remaining; send_chunk(pn, pn->in_buf, pn->pkt_size);
chunk = max_cont_data; pn->len = 0;
service_header = 0xFD; }
}
if (chunk == 0) break; // Add to buffer
size_t payload_len = 1 + chunk + (pos == 0 ? 1 : 0); // +1 байт типа для первого пакета *(uint16_t*)(pn->in_buf + pn->len) = item_len;
struct ll_entry* entry = queue_data_new(2 + payload_len); memcpy(pn->in_buf + pn->len + 2, entry->data, item_len);
if (!entry) return -1; pn->len += item_len + 2;
uint8_t* d = (entry);
*(uint16_t*)d = (uint16_t)payload_len;
d += 2;
*d++ = service_header;
if (pos == 0) {
*d++ = type;
} }
memcpy(d, (const uint8_t*)data + pos, chunk);
queue_data_put(pn->output, entry, 0); queue_data_free(data); // Free the entry
pos += chunk;
total_service_len -= chunk + (pos == chunk ? 2 : 1); // корректно вычитаем заголовки
} }
return 0;
} // If remaining in buffer, set timeout to flush if queues empty
int pkt_normalizer_get_error_count(const struct pn_struct* pn) { if (pn->len > 0) {
if (!pn) return 0; if (pn->flush_timer) {
if (pn->is_packer) { uasync_cancel_timeout(pn->ua, pn->flush_timer);
return pn->u.packer.error_count; }
pn->flush_timer = uasync_set_timeout(pn->ua, 10, pn, flush_cb); // 1ms = 10 * 0.1ms
} }
return pn->u.unpacker.error_count;
queue_resume_callback(q);
} }
void pkt_normalizer_reset_error_count(struct pn_struct* pn) {
// Internal: Flush callback on timeout
static void flush_cb(void* arg) {
struct PKTNORM* pn = (struct PKTNORM*)arg;
if (!pn) return; if (!pn) return;
if (pn->is_packer) {
pn->u.packer.error_count = 0; pn->flush_timer = NULL;
} else {
pn->u.unpacker.error_count = 0; if (pn->len > 0 && queue_entry_count(pn->input) == 0 && queue_entry_count(pn->pending) == 0) {
memset(pn->in_buf + pn->len, 0, pn->pkt_size - pn->len); // Pad
send_chunk(pn, pn->in_buf, pn->pkt_size);
pn->len = 0;
} }
} }
void pkt_normalizer_flush(struct pn_struct* pn) {
if (!pn || !pn->is_packer) return; // Internal: Add assembled data to etcp->output_queue as ETCP_FRAGMENT
if (pn->u.packer.len > 0) { static void add_to_output(struct PKTNORM* pn, uint8_t* app_data, uint16_t app_len) {
send_buf(pn); struct ETCP_FRAGMENT* frag = memory_pool_alloc(pn->etcp->rx_pool);
if (!frag) return;
frag->ll.dgram = memory_pool_alloc(pn->etcp->instance->data_pool);
if (!frag->ll.dgram) {
memory_pool_free(pn->etcp->rx_pool, frag);
return;
} }
memcpy(frag->ll.dgram, app_data, app_len);
frag->ll.size = app_len;
frag->ll.len = app_len;
frag->seq = 0;
frag->timestamp = 0;
queue_data_put(pn->etcp->output_queue, frag, 0);
} }
static void unpacker_handler(struct ll_queue* q, void* arg) {
struct pn_struct* pn = arg; // Internal: Process incoming fixed-size payload chunk from etcp
while (queue_entry_count(q) > 0) { void pn_unpacker_input(struct PKTNORM* pn, uint8_t* data, uint16_t len) {
uint8_t* entry = queue_data_get(q); if (!pn || len != pn->pkt_size) return;
uint8_t* data = entry;
uint16_t payload_len = *(uint16_t*)data; uint16_t header = *(uint16_t*)data;
if (entry != NULL && *(uint16_t*)entry != 2 + (size_t)payload_len) {
queue_data_free(entry); if (pn->in_fragment) {
continue; // Expect continuation
if (header != 0x0000) {
// Error: invalid continuation
pn_unpacker_reset_state(pn);
return;
} }
uint8_t* cg = data + 2;
size_t cg_pos = 0; uint16_t chunk = pn->pkt_size - 2;
size_t cg_len = (size_t)payload_len; if (pn->len + chunk > pn->cap) {
while (cg_pos < cg_len) { pn->cap *= 2;
uint8_t byte = cg[cg_pos++]; pn->out_buf = realloc(pn->out_buf, pn->cap);
if (byte == 0xFF) { }
/* Начало нового фрагментированного пакета */ memcpy(pn->out_buf + pn->len, data + 2, chunk);
if (pn->u.unpacker.in_fragment) { pn->len += chunk;
/* Не завершен предыдущий фрагмент - ошибка */
pn->u.unpacker.error_count++; if (pn->len >= pn->total_len) {
reset_fragment_state(pn); // Assembly complete: add to output_queue
} add_to_output(pn, pn->out_buf, pn->total_len);
/* Проверить, что есть 2 байта для общей длины */ pn_unpacker_reset_state(pn);
if (cg_pos + 2 > cg_len) { }
pn->u.unpacker.error_count++; } else {
goto err; if (header == 0xffff) {
} // Start of fragment
/* Прочитать общую длину */ pn->total_len = *(uint32_t*)(data + 2);
uint16_t total_len = ((uint16_t)cg[cg_pos] << 8) | cg[cg_pos + 1]; uint16_t chunk = pn->pkt_size - 6;
cg_pos += 2; if (pn->total_len <= chunk) {
size_t chunk_len = cg_len - cg_pos; // Invalid or degenerate case
/* Выделить буфер при необходимости */ pn_unpacker_reset_state(pn);
size_t new_len = chunk_len; return;
if (new_len > pn->u.unpacker.cap) {
size_t new_cap = pn->u.unpacker.cap ? pn->u.unpacker.cap * 2 : 4096;
if (new_cap < new_len) new_cap = new_len;
pn->u.unpacker.buf = realloc(pn->u.unpacker.buf, new_cap);
pn->u.unpacker.cap = new_cap;
}
memcpy(pn->u.unpacker.buf, cg + cg_pos, chunk_len);
pn->u.unpacker.len = chunk_len;
pn->u.unpacker.total_len = total_len;
pn->u.unpacker.in_fragment = 1;
cg_pos += chunk_len;
/* Проверить, не собрали ли уже весь пакет */
if (pn->u.unpacker.len >= pn->u.unpacker.total_len) {
if (pn->u.unpacker.len == pn->u.unpacker.total_len) {
uint8_t* out = queue_data_new(pn->u.unpacker.total_len);
if (out) {
memcpy(out, pn->u.unpacker.buf, pn->u.unpacker.total_len);
queue_data_put(pn->output, out, 0);
}
} else {
/* Слишком много данных - ошибка */
pn->u.unpacker.error_count++;
}
reset_fragment_state(pn);
}
continue;
}
if (byte == 0xFE) {
/* Продолжение фрагментированного пакета */
if (!pn->u.unpacker.in_fragment) {
/* Не было начала фрагмента - ошибка */
pn->u.unpacker.error_count++;
goto err;
}
size_t chunk_len = cg_len - cg_pos;
size_t new_len = pn->u.unpacker.len + chunk_len;
if (new_len > pn->u.unpacker.cap) {
size_t new_cap = pn->u.unpacker.cap ? pn->u.unpacker.cap * 2 : 4096;
if (new_cap < new_len) new_cap = new_len;
pn->u.unpacker.buf = realloc(pn->u.unpacker.buf, new_cap);
pn->u.unpacker.cap = new_cap;
}
memcpy(pn->u.unpacker.buf + pn->u.unpacker.len, cg + cg_pos, chunk_len);
pn->u.unpacker.len = new_len;
cg_pos += chunk_len;
/* Проверить, не собрали ли уже весь пакет */
if (pn->u.unpacker.len >= pn->u.unpacker.total_len) {
if (pn->u.unpacker.len == pn->u.unpacker.total_len) {
uint8_t* out = queue_data_new(pn->u.unpacker.total_len);
if (out) {
memcpy(out, pn->u.unpacker.buf, pn->u.unpacker.total_len);
queue_data_put(pn->output, out, 0);
}
} else {
/* Слишком много данных - ошибка */
pn->u.unpacker.error_count++;
}
reset_fragment_state(pn);
}
continue;
}
if (byte == 0xFC || byte == 0xFD) {
/* Service packet */
if (byte == 0xFC) {
/* Start of service packet */
if (pn->u.unpacker.in_service) {
/* Previous service packet finished - deliver it */
if (pn->service_callback) {
pn->service_callback(pn->service_callback_user_data,
pn->u.unpacker.service_type,
pn->u.unpacker.service_buf,
pn->u.unpacker.service_len);
}
free(pn->u.unpacker.service_buf);
pn->u.unpacker.service_buf = NULL;
pn->u.unpacker.service_len = 0;
pn->u.unpacker.service_cap = 0;
pn->u.unpacker.in_service = 0;
}
/* Read service type */
if (cg_pos >= cg_len) goto err;
uint8_t service_type = cg[cg_pos++];
pn->u.unpacker.service_type = service_type;
pn->u.unpacker.in_service = 1;
pn->u.unpacker.service_len = 0;
} else {
/* 0xFD - continuation */
if (!pn->u.unpacker.in_service) {
/* No service packet started - error */
pn->u.unpacker.error_count++;
goto err;
}
}
/* Read data */
size_t data_len = cg_len - cg_pos;
if (data_len > 0) {
size_t new_len = pn->u.unpacker.service_len + data_len;
if (new_len > 256) {
/* Service packet too long - error */
pn->u.unpacker.error_count++;
free(pn->u.unpacker.service_buf);
pn->u.unpacker.service_buf = NULL;
pn->u.unpacker.service_len = 0;
pn->u.unpacker.service_cap = 0;
pn->u.unpacker.in_service = 0;
goto err;
}
if (new_len > pn->u.unpacker.service_cap) {
size_t new_cap = pn->u.unpacker.service_cap ? pn->u.unpacker.service_cap * 2 : 256;
if (new_cap < new_len) new_cap = new_len;
if (new_cap > 256) new_cap = 256;
uint8_t* new_buf = realloc(pn->u.unpacker.service_buf, new_cap);
if (!new_buf) {
pn->u.unpacker.error_count++;
goto err;
}
pn->u.unpacker.service_buf = new_buf;
pn->u.unpacker.service_cap = new_cap;
}
memcpy(pn->u.unpacker.service_buf + pn->u.unpacker.service_len, cg + cg_pos, data_len);
pn->u.unpacker.service_len = new_len;
cg_pos += data_len;
}
/* Check if this is the end of service packet (end of payload) */
if (cg_pos >= cg_len) {
/* End of current payload, but service packet may continue in next transport packet */
continue;
} else {
/* There is more data in this payload after service packet - error */
pn->u.unpacker.error_count++;
free(pn->u.unpacker.service_buf);
pn->u.unpacker.service_buf = NULL;
pn->u.unpacker.service_len = 0;
pn->u.unpacker.service_cap = 0;
pn->u.unpacker.in_service = 0;
goto err;
}
}
/* Обычная запись (не фрагмент) */
size_t L;
if (byte <= 0xEF) {
L = byte;
} else if (byte >= 0xF0 && byte <= 0xF5) {
if (cg_pos >= cg_len) goto err;
uint8_t ext = cg[cg_pos++];
L = ((size_t)(byte - 0xF0) << 8) | ext;
} else {
/* Недопустимый байт */
goto err;
} }
if (cg_pos + L > cg_len) goto err; memcpy(pn->out_buf, data + 6, chunk);
if (pn->u.unpacker.in_fragment) { pn->len = chunk;
/* Это последний фрагмент в виде обычной записи */ pn->in_fragment = 1;
size_t new_len = pn->u.unpacker.len + L; } else {
if (new_len > pn->u.unpacker.cap) { // Normal aggregated packet: process multiple items (ignore pad)
size_t new_cap = pn->u.unpacker.cap ? pn->u.unpacker.cap * 2 : 4096; uint8_t* ptr = data;
if (new_cap < new_len) new_cap = new_len; size_t processed = 0;
pn->u.unpacker.buf = realloc(pn->u.unpacker.buf, new_cap); while (processed < pn->pkt_size) {
pn->u.unpacker.cap = new_cap; header = *(uint16_t*)ptr;
} if (header == 0 || header == 0xffff || header == 0x0000) break; // End (pad) or error
memcpy(pn->u.unpacker.buf + pn->u.unpacker.len, cg + cg_pos, L); if (processed + 2 + header > pn->pkt_size) break;
pn->u.unpacker.len = new_len;
cg_pos += L; add_to_output(pn, ptr + 2, header);
/* Проверить, собрали ли весь пакет */
if (pn->u.unpacker.len >= pn->u.unpacker.total_len) { ptr += 2 + header;
if (pn->u.unpacker.len == pn->u.unpacker.total_len) { processed += 2 + header;
uint8_t* out = queue_data_new(pn->u.unpacker.total_len);
if (out) {
memcpy(out, pn->u.unpacker.buf, pn->u.unpacker.total_len);
queue_data_put(pn->output, out, 0);
}
} else {
/* Слишком много данных - ошибка */
pn->u.unpacker.error_count++;
}
reset_fragment_state(pn);
}
} else {
/* Обычная запись (не часть фрагмента) */
uint8_t* out = queue_data_new(L);
if (out) {
memcpy(out, cg + cg_pos, L);
queue_data_put(pn->output, out, 0);
}
cg_pos += L;
} }
} }
err:
queue_data_free(entry);
} }
queue_resume_callback(q);
} }

104
src/pkt_normalizer.h

@ -1,75 +1,59 @@
// pkt_normalizer.h // pkt_normalizer.h (упрощенная версия)
#ifndef PKT_NORMALIZER_H #ifndef PKT_NORMALIZER_H
#define PKT_NORMALIZER_H #define PKT_NORMALIZER_H
#include "ll_queue.h" #include "../lib/ll_queue.h"
#include "../lib/u_async.h" #include "../lib/u_async.h"
#include <stdint.h> #include <stdint.h>
/* Default fragment reassembly timeout in uasync timebase units (0.1 ms) */ // Структура для packer
#ifndef PKT_NORMALIZER_FRAGMENT_TIMEOUT struct PKTNORM {
#define PKT_NORMALIZER_FRAGMENT_TIMEOUT 5000 /* 500 ms */ // public:
#endif struct ll_queue* input; // Входная очередь в packer (через нее отправляем пакеты)
struct ll_queue* output; // Выходная очередь из unpacker (через нее принимаем пакеты)
// private:
struct ETCP_CONN* etcp;
uint16_t pkt_size;
// packer:
uasync_t* ua; // uasync instance
uint8_t* in_buf; // Буфер для упаковки
size_t len; // Текущая длина в буфере
size_t cap; // Емкость буфера (MAX_PACKET_SIZE)
void* flush_timer; // For timeout flush
struct ll_queue* pending; // For not filling input
// unpacker:
uint8_t* out_buf; // Буфер для сборки фрагментов
// size_t len; // Текущая накопленная длина
size_t total_len; // Ожидаемая общая длина (для фрагментов)
// size_t cap; // Емкость буфера
int in_fragment; // Флаг: идет сборка фрагмента (1) или нет (0)
struct ll_queue* send_pending; // For not filling etcp->input_queue
/* ETCP overhead for calculating fragment size from MTU */
#define ETCP_OVERHEAD 100 // Reserve 100 bytes for headers, crypto, etc.*/
/* Service packet callback type */
typedef void (*pkt_normalizer_service_callback_t)(void* user_data, uint8_t type, const uint8_t* data, size_t len);
struct pn_struct {
struct ll_queue* input;
struct ll_queue* output;
uasync_t* ua;
int is_packer;
union {
struct {
uint8_t* buf;
size_t len;
size_t cap;
int error_count;
} packer;
struct {
uint8_t* buf; /* буфер для сборки фрагментов */
size_t len; /* текущая накопленная длина */
size_t total_len; /* ожидаемая общая длина из первого фрагмента */
size_t cap; /* ёмкость буфера */
int error_count; /* счетчик ошибок сборки */
int in_fragment; /* флаг: идет сборка фрагментов (1) или нет (0) */
/* Service packet reassembly */
uint8_t* service_buf; /* буфер для сборки сервисных пакетов */
size_t service_len; /* текущая накопленная длина сервисного пакета */
size_t service_cap; /* ёмкость буфера сервисного пакета */
uint8_t service_type; /* тип сервисного пакета */
int in_service; /* флаг: идет сборка сервисного пакета (1) или нет (0) */
} unpacker;
} u;
/* Service packet callback */
pkt_normalizer_service_callback_t service_callback;
void* service_callback_user_data;
}; };
struct pkt_normalizer_pair { // Инициализация пары
struct pn_struct* packer; struct PKTNORM* pn_init(struct ETCP_CONN* etcp);// все что нужно (в т.ч. mtu и ua) берет из etcp
struct pn_struct* unpacker;
}; // Деинициализация пары
void pn_pair_deinit(struct PKTNORM* pn);
struct pn_struct* pkt_normalizer_init(uasync_t* ua, int is_packer, int mtu); // 1 for packer, 0 for unpacker, mtu for fragment size calculation // Сброс состояния (для unpacker)
void pkt_normalizer_deinit(struct pn_struct* pn); void pn_unpacker_reset_state(struct PKTNORM* pn);
struct pkt_normalizer_pair* pkt_normalizer_pair_init(uasync_t* ua, int mtu); // создаёт malloc data, копирует, помещает в input.
void pkt_normalizer_pair_deinit(struct pkt_normalizer_pair* pair); void pn_packer_send(struct PKTNORM* pn, uint8_t* data, uint16_t len);
/* Error handling */ /* Как работает:
int pkt_normalizer_get_error_count(const struct pn_struct* pn); Формат отправки в etcp: 2 байта размер, далее данные (порезанные на куски и отправленные через etcp)
void pkt_normalizer_reset_error_count(struct pn_struct* pn); собирает по возможности полные пакеты с размером pkt_size.
+ timeout: неполный пакет отправляется по таймауту 1ms (если очереди пустые)
+ входящая очередь etcp не должне наполняться для минимизации задержки - новый пакет отправляем только когда очередь пустая
/* Flush internal buffer (packer only) */ */
void pkt_normalizer_flush(struct pn_struct* pn);
int pkt_normalizer_send_service(struct pn_struct* pn, uint8_t type, const void* data, size_t len);
void pkt_normalizer_set_service_callback(struct pn_struct* pn, pkt_normalizer_service_callback_t callback, void* user_data);
void pkt_normalizer_reset_service_state(struct pn_struct* pn);
void pkt_normalizer_reset_state(struct pn_struct* pn);
#endif // PKT_NORMALIZER_H #endif // PKT_NORMALIZER_H

369
src/secure_channel.c.bak

@ -1,369 +0,0 @@
/* sc_lib.c - Secure Channel library implementation using TinyCrypt */
#include "secure_channel.h"
#include "../tinycrypt/lib/include/tinycrypt/ecc.h"
#include "../tinycrypt/lib/include/tinycrypt/ecc_dh.h"
#include "../tinycrypt/lib/include/tinycrypt/aes.h"
#include "../tinycrypt/lib/include/tinycrypt/ccm_mode.h"
#include "../tinycrypt/lib/include/tinycrypt/constants.h"
#include "../tinycrypt/lib/include/tinycrypt/ecc_platform_specific.h"
#include "../tinycrypt/lib/include/tinycrypt/sha256.h"
#include <string.h>
#include <stddef.h>
#include <sys/types.h>
#include <unistd.h>
#include <sys/time.h>
#include <stdio.h>
// Simple debug macros
#define DEBUG_CATEGORY_CRYPTO 1
#define DEBUG_ERROR(category, fmt, ...) fprintf(stderr, "ERROR: " fmt "\n", ##__VA_ARGS__)
#define DEBUG_INFO(category, fmt, ...) fprintf(stdout, "INFO: " fmt "\n", ##__VA_ARGS__)
#include <stdio.h>
#include <fcntl.h>
#include "crc32.h"
static const struct uECC_Curve_t *curve = NULL;
static uint8_t sc_urandom_seed[8] = {0};
static int sc_urandom_initialized = 0;
static void sc_init_random_seed(void)
{
int fd = open("/dev/urandom", O_RDONLY);
if (fd >= 0) {
ssize_t ret = read(fd, sc_urandom_seed, 8);
close(fd);
if (ret == 8) {
sc_urandom_initialized = 1;
}
}
}
static int sc_rng(uint8_t *dest, unsigned size)
{
int fd = open("/dev/urandom", O_RDONLY);
if (fd < 0) {
return 0;
}
ssize_t ret = read(fd, dest, size);
close(fd);
if (ret != size) {
return 0;
}
/* Mix in PID and microtime for additional entropy */
pid_t pid = getpid();
struct timeval tv;
gettimeofday(&tv, NULL);
for (unsigned i = 0; i < size; i++) {
dest[i] ^= ((pid >> (i % (sizeof(pid) * 8))) & 0xFF);
dest[i] ^= ((tv.tv_sec >> (i % (sizeof(tv.tv_sec) * 8))) & 0xFF);
dest[i] ^= ((tv.tv_usec >> (i % (sizeof(tv.tv_usec) * 8))) & 0xFF);
}
return 1;
}
static int sc_validate_key(const uint8_t *public_key)
{
if (!curve) {
curve = uECC_secp256r1();
}
int result = uECC_valid_public_key(public_key, curve);
DEBUG_INFO(DEBUG_CATEGORY_CRYPTO, "sc_validate_key: uECC_valid_public_key returned %d", result);
return result;
}
sc_status_t sc_generate_keypair(struct SC_MYKEYS *pk)
{
if (!pk) {
return SC_ERR_INVALID_ARG;
}
if (!curve) {
curve = uECC_secp256r1();
}
/* Set custom RNG function */
uECC_set_rng(sc_rng);
if (!uECC_make_key(pk->public_key, pk->private_key, curve)) {
return SC_ERR_CRYPTO;
}
return SC_OK;
}
// Конвертация hex строки в бинарный формат
static int hex_to_binary(const char *hex_str, uint8_t *binary, size_t binary_len) {
if (!hex_str || !binary || strlen(hex_str) != binary_len * 2) return -1;
for (size_t i = 0; i < binary_len; i++) {
unsigned int byte;
if (sscanf(hex_str + i * 2, "%2x", &byte) != 1) return -1;
binary[i] = (uint8_t)byte;
}
return 0;
}
sc_status_t sc_init_local_keys(struct SC_MYKEYS *mykeys, const char *public_key, const char *private_key) {
if (!mykeys || !public_key || !private_key) {
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "sc_init_local_keys: invalid arguments");
return SC_ERR_INVALID_ARG;
}
if (!curve) {
curve = uECC_secp256r1();
}
DEBUG_INFO(DEBUG_CATEGORY_CRYPTO, "sc_init_local_keys: public_key len=%zu, private_key len=%zu",
strlen(public_key), strlen(private_key));
/* Convert hex to binary first */
if (hex_to_binary(public_key, mykeys->public_key, SC_PUBKEY_SIZE)) {
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "sc_init_local_keys: failed to convert public key from hex");
return SC_ERR_INVALID_ARG;
}
if (hex_to_binary(private_key, mykeys->private_key, SC_PRIVKEY_SIZE)) {
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "sc_init_local_keys: failed to convert private key from hex");
return SC_ERR_INVALID_ARG;
}
/* Validate the converted binary public key */
if (sc_validate_key(mykeys->public_key) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "sc_init_local_keys: public key validation failed");
return SC_ERR_INVALID_ARG;
}
DEBUG_INFO(DEBUG_CATEGORY_CRYPTO, "sc_init_local_keys: keys initialized successfully");
return SC_OK;
}
sc_status_t sc_init_ctx(sc_context_t *ctx, struct SC_MYKEYS *mykeys) {
ctx->pk=mykeys;
ctx->initialized = 1;
ctx->peer_key_set = 0;
ctx->session_ready = 0;
ctx->tx_counter = 0;
ctx->rx_counter = 0;
return SC_OK;
}
sc_status_t sc_set_peer_public_key(sc_context_t *ctx, const char *peer_public_key_h, int mode) {
uint8_t shared_secret[SC_SHARED_SECRET_SIZE];
uint8_t peer_public_key[SC_PUBKEY_SIZE];
if (mode) {
if (hex_to_binary(peer_public_key_h, peer_public_key, SC_PUBKEY_SIZE)) {
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "sc_set_peer_public_key: invalid hex key format");
return SC_ERR_INVALID_ARG;
}
}
else memcpy(peer_public_key, peer_public_key_h, SC_PUBKEY_SIZE);
if (!ctx) {
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "sc_set_peer_public_key: invalid ctx");
return SC_ERR_INVALID_ARG;
}
if (!ctx->initialized) {
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "sc_set_peer_public_key: ctx not initialized");
return SC_ERR_NOT_INITIALIZED;
}
if (!curve) {
curve = uECC_secp256r1();
}
/* Validate peer public key */
if (sc_validate_key(peer_public_key) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "sc_set_peer_public_key: invalid key");
return SC_ERR_INVALID_ARG;
}
/* Compute shared secret using ECDH */
if (!ctx->pk) {
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "sc_set_peer_public_key: no private key");
return SC_ERR_NOT_INITIALIZED;
}
if (!uECC_shared_secret(peer_public_key, ctx->pk->private_key,
shared_secret, curve)) {
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "sc_set_peer_public_key: shared secret error");
return SC_ERR_CRYPTO;
}
/* Derive session key from shared secret (simple copy for demo) */
memcpy(ctx->session_key, shared_secret, SC_SESSION_KEY_SIZE);
/* Store peer public key */
memcpy(ctx->peer_public_key, peer_public_key, SC_PUBKEY_SIZE);
ctx->peer_key_set = 1;
ctx->session_ready = 1;
return SC_OK;
}
// Новая функция для генерации nonce с микросекундами
static void generate_nonce_with_timer(uint8_t *nonce) {
struct timeval tv;
gettimeofday(&tv, NULL);
uint32_t usec = (uint32_t)tv.tv_usec;
// Поместить usec в первые 4 байта (little-endian)
nonce[0] = usec & 0xFF;
nonce[1] = (usec >> 8) & 0xFF;
nonce[2] = (usec >> 16) & 0xFF;
nonce[3] = (usec >> 24) & 0xFF;
// Заполнить оставшиеся 9 байт случайными данными
sc_rng(nonce + 4, SC_NONCE_SIZE - 4);
}
sc_status_t sc_encrypt(sc_context_t *ctx, const uint8_t *plaintext, size_t plaintext_len, uint8_t *ciphertext, size_t *ciphertext_len) {
uint8_t nonce[SC_NONCE_SIZE];
uint8_t plaintext_with_crc[plaintext_len + SC_CRC32_SIZE];
size_t total_plaintext_len = plaintext_len + SC_CRC32_SIZE;
uint8_t combined_output[total_plaintext_len + SC_TAG_SIZE];
struct tc_aes_key_sched_struct sched;
struct tc_ccm_mode_struct ccm_state;
if (!ctx || !plaintext || !ciphertext || !ciphertext_len) {
return SC_ERR_INVALID_ARG;
}
if (!ctx->session_ready) {
return SC_ERR_NOT_INITIALIZED;
}
if (plaintext_len == 0) {
return SC_ERR_INVALID_ARG;
}
/* Добавляем CRC32 к данным */
memcpy(plaintext_with_crc, plaintext, plaintext_len);
uint32_t crc = crc32_calc(plaintext, plaintext_len);
plaintext_with_crc[plaintext_len] = (crc >> 0) & 0xFF;
plaintext_with_crc[plaintext_len + 1] = (crc >> 8) & 0xFF;
plaintext_with_crc[plaintext_len + 2] = (crc >> 16) & 0xFF;
plaintext_with_crc[plaintext_len + 3] = (crc >> 24) & 0xFF;
/* Генерируем nonce с таймером */
generate_nonce_with_timer(nonce);
/* Initialize AES key schedule */
if (tc_aes128_set_encrypt_key(&sched, ctx->session_key) != TC_CRYPTO_SUCCESS) {
return SC_ERR_CRYPTO;
}
/* Configure CCM mode */
if (tc_ccm_config(&ccm_state, &sched, nonce, SC_NONCE_SIZE, SC_TAG_SIZE) != TC_CRYPTO_SUCCESS) {
return SC_ERR_CRYPTO;
}
/* Encrypt and generate tag */
if (tc_ccm_generation_encryption(combined_output, sizeof(combined_output),
NULL, 0, /* no associated data */
plaintext_with_crc, total_plaintext_len,
&ccm_state) != TC_CRYPTO_SUCCESS) {
return SC_ERR_CRYPTO;
}
/* Copy nonce + ciphertext + tag to output buffer */
memcpy(ciphertext, nonce, SC_NONCE_SIZE);
memcpy(ciphertext + SC_NONCE_SIZE, combined_output, total_plaintext_len + SC_TAG_SIZE);
*ciphertext_len = SC_NONCE_SIZE + total_plaintext_len + SC_TAG_SIZE;
ctx->tx_counter++;
return SC_OK;
}
sc_status_t sc_decrypt(sc_context_t *ctx,
const uint8_t *ciphertext,
size_t ciphertext_len,
uint8_t *plaintext,
size_t *plaintext_len)
{
uint8_t nonce[SC_NONCE_SIZE];
struct tc_aes_key_sched_struct sched;
struct tc_ccm_mode_struct ccm_state;
size_t total_plaintext_len = ciphertext_len - SC_NONCE_SIZE - SC_TAG_SIZE;
uint8_t plaintext_with_crc[total_plaintext_len];
if (!ctx || !ciphertext || !plaintext || !plaintext_len) {
return SC_ERR_INVALID_ARG;
}
if (!ctx->session_ready) {
return SC_ERR_NOT_INITIALIZED;
}
if (ciphertext_len < SC_NONCE_SIZE + SC_TAG_SIZE + SC_CRC32_SIZE) {
return SC_ERR_INVALID_ARG;
}
/* Извлекаем nonce из начала ciphertext */
memcpy(nonce, ciphertext, SC_NONCE_SIZE);
/* Ciphertext для расшифровки начинается после nonce */
const uint8_t *encrypted_data = ciphertext + SC_NONCE_SIZE;
size_t encrypted_len = ciphertext_len - SC_NONCE_SIZE;
/* Initialize AES key schedule */
if (tc_aes128_set_encrypt_key(&sched, ctx->session_key) != TC_CRYPTO_SUCCESS) {
return SC_ERR_CRYPTO;
}
/* Configure CCM mode с извлечённым nonce */
if (tc_ccm_config(&ccm_state, &sched, nonce, SC_NONCE_SIZE, SC_TAG_SIZE) != TC_CRYPTO_SUCCESS) {
return SC_ERR_CRYPTO;
}
/* Decrypt and verify tag */
if (tc_ccm_decryption_verification(plaintext_with_crc, total_plaintext_len,
NULL, 0, /* no associated data */
encrypted_data, encrypted_len,
&ccm_state) != TC_CRYPTO_SUCCESS) {
return SC_ERR_AUTH_FAILED;
}
/* Проверяем CRC32 */
size_t data_len = total_plaintext_len - SC_CRC32_SIZE;
uint32_t expected_crc = crc32_calc(plaintext_with_crc, data_len);
uint32_t received_crc = (plaintext_with_crc[data_len] << 0) |
(plaintext_with_crc[data_len + 1] << 8) |
(plaintext_with_crc[data_len + 2] << 16) |
(plaintext_with_crc[data_len + 3] << 24);
if (expected_crc != received_crc) {
return SC_ERR_CRC_FAILED;
}
/* Копируем данные без CRC32 */
memcpy(plaintext, plaintext_with_crc, data_len);
*plaintext_len = data_len;
ctx->rx_counter++;
return SC_OK;
}
sc_status_t sc_compute_public_key_from_private(const uint8_t *private_key, uint8_t *public_key) {
if (!private_key || !public_key) {
return SC_ERR_INVALID_ARG;
}
if (!curve) {
curve = uECC_secp256r1();
}
if (!uECC_compute_public_key(private_key, public_key, curve)) {
return SC_ERR_CRYPTO;
}
return SC_OK;
}

388
src/secure_channel.c1

@ -1,388 +0,0 @@
/* sc_lib.c - Secure Channel library implementation using TinyCrypt */
#include "secure_channel.h"
#include "../tinycrypt/lib/include/tinycrypt/ecc.h"
#include "../tinycrypt/lib/include/tinycrypt/ecc_dh.h"
#include "../tinycrypt/lib/include/tinycrypt/aes.h"
#include "../tinycrypt/lib/include/tinycrypt/ccm_mode.h"
#include "../tinycrypt/lib/include/tinycrypt/constants.h"
#include "../tinycrypt/lib/include/tinycrypt/ecc_platform_specific.h"
#include "../tinycrypt/lib/include/tinycrypt/sha256.h"
#include <string.h>
#include <stddef.h>
#include <sys/types.h>
#include <unistd.h>
#include <sys/time.h>
#include <stdio.h>
// Simple debug macros
#define DEBUG_CATEGORY_CRYPTO 1
#define DEBUG_ERROR(category, fmt, ...) fprintf(stderr, "ERROR: " fmt "\n", ##__VA_ARGS__)
#define DEBUG_INFO(category, fmt, ...) fprintf(stdout, "INFO: " fmt "\n", ##__VA_ARGS__)
#include <stdio.h>
#include <fcntl.h>
#include "crc32.h"
static const struct uECC_Curve_t *curve = NULL;
static uint8_t sc_urandom_seed[8] = {0};
static int sc_urandom_initialized = 0;
static void sc_init_random_seed(void)
{
int fd = open("/dev/urandom", O_RDONLY);
if (fd >= 0) {
ssize_t ret = read(fd, sc_urandom_seed, 8);
close(fd);
if (ret == 8) {
sc_urandom_initialized = 1;
}
}
}
static int sc_rng(uint8_t *dest, unsigned size)
{
int fd = open("/dev/urandom", O_RDONLY);
if (fd < 0) {
return 0;
}
ssize_t ret = read(fd, dest, size);
close(fd);
if (ret != size) {
return 0;
}
/* Mix in PID and microtime for additional entropy */
pid_t pid = getpid();
struct timeval tv;
gettimeofday(&tv, NULL);
for (unsigned i = 0; i < size; i++) {
dest[i] ^= ((pid >> (i % (sizeof(pid) * 8))) & 0xFF);
dest[i] ^= ((tv.tv_sec >> (i % (sizeof(tv.tv_sec) * 8))) & 0xFF);
dest[i] ^= ((tv.tv_usec >> (i % (sizeof(tv.tv_usec) * 8))) & 0xFF);
}
return 1;
}
static int sc_validate_key(const uint8_t *public_key)
{
if (!curve) {
curve = uECC_secp256r1();
}
int result = uECC_valid_public_key(public_key, curve);
DEBUG_INFO(DEBUG_CATEGORY_CRYPTO, "sc_validate_key: uECC_valid_public_key returned %d", result);
return result;
}
sc_status_t sc_generate_keypair(struct SC_MYKEYS *pk)
{
if (!pk) {
return SC_ERR_INVALID_ARG;
}
if (!curve) {
curve = uECC_secp256r1();
}
/* Set custom RNG function */
uECC_set_rng(sc_rng);
if (!uECC_make_key(pk->public_key, pk->private_key, curve)) {
return SC_ERR_CRYPTO;
}
return SC_OK;
}
// Конвертация hex строки в бинарный формат
static int hex_to_binary(const char *hex_str, uint8_t *binary, size_t binary_len) {
if (!hex_str || !binary || strlen(hex_str) != binary_len * 2) return -1;
for (size_t i = 0; i < binary_len; i++) {
unsigned int byte;
if (sscanf(hex_str + i * 2, "%2x", &byte) != 1) return -1;
binary[i] = (uint8_t)byte;
}
return 0;
}
sc_status_t sc_init_local_keys(struct SC_MYKEYS *mykeys, const char *public_key, const char *private_key) {
if (!mykeys || !public_key || !private_key) {
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "sc_init_local_keys: invalid arguments");
return SC_ERR_INVALID_ARG;
}
if (!curve) {
curve = uECC_secp256r1();
}
DEBUG_INFO(DEBUG_CATEGORY_CRYPTO, "sc_init_local_keys: public_key len=%zu, private_key len=%zu",
strlen(public_key), strlen(private_key));
/* Convert hex to binary first */
if (hex_to_binary(public_key, mykeys->public_key, SC_PUBKEY_SIZE)) {
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "sc_init_local_keys: failed to convert public key from hex");
return SC_ERR_INVALID_ARG;
}
if (hex_to_binary(private_key, mykeys->private_key, SC_PRIVKEY_SIZE)) {
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "sc_init_local_keys: failed to convert private key from hex");
return SC_ERR_INVALID_ARG;
}
/* Validate the converted binary public key */
if (sc_validate_key(mykeys->public_key) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "sc_init_local_keys: public key validation failed");
return SC_ERR_INVALID_ARG;
}
DEBUG_INFO(DEBUG_CATEGORY_CRYPTO, "sc_init_local_keys: keys initialized successfully");
return SC_OK;
}
sc_status_t sc_init_ctx(sc_context_t *ctx, struct SC_MYKEYS *mykeys) {
ctx->pk=mykeys;
ctx->initialized = 1;
ctx->peer_key_set = 0;
ctx->session_ready = 0;
ctx->tx_counter = 0;
ctx->rx_counter = 0;
return SC_OK;
}
sc_status_t sc_set_peer_public_key(sc_context_t *ctx, const char *peer_public_key_h, int mode) {
uint8_t shared_secret[SC_SHARED_SECRET_SIZE];
uint8_t peer_public_key[SC_PUBKEY_SIZE];
if (mode) {
if (hex_to_binary(peer_public_key_h, peer_public_key, SC_PUBKEY_SIZE)) {
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "sc_set_peer_public_key: invalid hex key format");
return SC_ERR_INVALID_ARG;
}
}
else memcpy(peer_public_key, peer_public_key_h, SC_PUBKEY_SIZE);
if (!ctx) {
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "sc_set_peer_public_key: invalid ctx");
return SC_ERR_INVALID_ARG;
}
if (!ctx->initialized) {
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "sc_set_peer_public_key: ctx not initialized");
return SC_ERR_NOT_INITIALIZED;
}
if (!curve) {
curve = uECC_secp256r1();
}
/* Validate peer public key */
if (sc_validate_key(peer_public_key) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "sc_set_peer_public_key: invalid key");
return SC_ERR_INVALID_ARG;
}
/* Compute shared secret using ECDH */
if (!ctx->pk) {
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "sc_set_peer_public_key: no private key");
return SC_ERR_NOT_INITIALIZED;
}
if (!uECC_shared_secret(peer_public_key, ctx->pk->private_key,
shared_secret, curve)) {
DEBUG_ERROR(DEBUG_CATEGORY_CRYPTO, "sc_set_peer_public_key: shared secret error");
return SC_ERR_CRYPTO;
}
/* Derive session key from shared secret (simple copy for demo) */
memcpy(ctx->session_key, shared_secret, SC_SESSION_KEY_SIZE);
/* Store peer public key */
memcpy(ctx->peer_public_key, peer_public_key, SC_PUBKEY_SIZE);
ctx->peer_key_set = 1;
ctx->session_ready = 1;
return SC_OK;
}
static void sc_build_nonce(uint64_t counter, uint8_t *nonce_out)
{
struct tc_sha256_state_struct sha_ctx;
uint8_t hash[32];
struct timeval tv;
uint8_t data[8 + 8 + 4];
if (!sc_urandom_initialized) {
sc_init_random_seed();
}
gettimeofday(&tv, NULL);
memcpy(data, sc_urandom_seed, 8);
data[8] = (counter >> 0) & 0xFF;
data[9] = (counter >> 8) & 0xFF;
data[10] = (counter >> 16) & 0xFF;
data[11] = (counter >> 24) & 0xFF;
data[12] = (counter >> 32) & 0xFF;
data[13] = (counter >> 40) & 0xFF;
data[14] = (counter >> 48) & 0xFF;
data[15] = (counter >> 56) & 0xFF;
data[16] = (tv.tv_sec >> 0) & 0xFF;
data[17] = (tv.tv_sec >> 8) & 0xFF;
data[18] = (tv.tv_sec >> 16) & 0xFF;
data[19] = (tv.tv_sec >> 24) & 0xFF;
tc_sha256_init(&sha_ctx);
tc_sha256_update(&sha_ctx, data, 20);
tc_sha256_final(hash, &sha_ctx);
memcpy(nonce_out, hash, SC_NONCE_SIZE);
}
sc_status_t sc_encrypt(sc_context_t *ctx,
const uint8_t *plaintext,
size_t plaintext_len,
uint8_t *ciphertext,
size_t *ciphertext_len)
{
uint8_t nonce[SC_NONCE_SIZE];
struct tc_aes_key_sched_struct sched;
struct tc_ccm_mode_struct ccm_state;
size_t total_plaintext_len = plaintext_len + SC_CRC32_SIZE;
uint8_t plaintext_with_crc[total_plaintext_len];
uint8_t combined_output[total_plaintext_len + SC_TAG_SIZE];
if (!ctx || !plaintext || !ciphertext || !ciphertext_len) {
return SC_ERR_INVALID_ARG;
}
if (!ctx->session_ready) {
return SC_ERR_NOT_INITIALIZED;
}
if (plaintext_len == 0) {
return SC_ERR_INVALID_ARG;
}
/* Добавляем CRC32 к данным */
memcpy(plaintext_with_crc, plaintext, plaintext_len);
uint32_t crc = crc32_calc(plaintext, plaintext_len);
plaintext_with_crc[plaintext_len] = (crc >> 0) & 0xFF;
plaintext_with_crc[plaintext_len + 1] = (crc >> 8) & 0xFF;
plaintext_with_crc[plaintext_len + 2] = (crc >> 16) & 0xFF;
plaintext_with_crc[plaintext_len + 3] = (crc >> 24) & 0xFF;
/* Initialize AES key schedule */
if (tc_aes128_set_encrypt_key(&sched, ctx->session_key) != TC_CRYPTO_SUCCESS) {
return SC_ERR_CRYPTO;
}
/* Build nonce from counter */
sc_build_nonce(ctx->tx_counter, nonce);
/* Configure CCM mode */
if (tc_ccm_config(&ccm_state, &sched, nonce, SC_NONCE_SIZE, SC_TAG_SIZE) != TC_CRYPTO_SUCCESS) {
return SC_ERR_CRYPTO;
}
/* Encrypt and generate tag */
if (tc_ccm_generation_encryption(combined_output, sizeof(combined_output),
NULL, 0, /* no associated data */
plaintext_with_crc, total_plaintext_len,
&ccm_state) != TC_CRYPTO_SUCCESS) {
return SC_ERR_CRYPTO;
}
/* Copy ciphertext + tag to output buffer */
memcpy(ciphertext, combined_output, total_plaintext_len + SC_TAG_SIZE);
*ciphertext_len = total_plaintext_len + SC_TAG_SIZE;
ctx->tx_counter++;
return SC_OK;
}
sc_status_t sc_decrypt(sc_context_t *ctx,
const uint8_t *ciphertext,
size_t ciphertext_len,
uint8_t *plaintext,
size_t *plaintext_len)
{
uint8_t nonce[SC_NONCE_SIZE];
struct tc_aes_key_sched_struct sched;
struct tc_ccm_mode_struct ccm_state;
TCCcmMode_t c = &ccm_state;
size_t total_plaintext_len = ciphertext_len - SC_TAG_SIZE;
uint8_t plaintext_with_crc[total_plaintext_len];
if (!ctx || !ciphertext || !plaintext || !plaintext_len) {
return SC_ERR_INVALID_ARG;
}
if (!ctx->session_ready) {
return SC_ERR_NOT_INITIALIZED;
}
if (ciphertext_len < SC_TAG_SIZE + SC_CRC32_SIZE) {
return SC_ERR_INVALID_ARG;
}
/* Initialize AES key schedule */
if (tc_aes128_set_encrypt_key(&sched, ctx->session_key) != TC_CRYPTO_SUCCESS) {
return SC_ERR_CRYPTO;
}
/* Build nonce from counter */
sc_build_nonce(ctx->rx_counter, nonce);
/* Configure CCM mode */
if (tc_ccm_config(c, &sched, nonce, SC_NONCE_SIZE, SC_TAG_SIZE) != TC_CRYPTO_SUCCESS) {
return SC_ERR_CRYPTO;
}
/* Decrypt and verify tag */
if (tc_ccm_decryption_verification(plaintext_with_crc, total_plaintext_len,
NULL, 0, /* no associated data */
ciphertext, ciphertext_len,
c) != TC_CRYPTO_SUCCESS) {
return SC_ERR_AUTH_FAILED;
}
/* Проверяем CRC32 */
size_t data_len = total_plaintext_len - SC_CRC32_SIZE;
uint32_t expected_crc = crc32_calc(plaintext_with_crc, data_len);
uint32_t received_crc = (plaintext_with_crc[data_len] << 0) |
(plaintext_with_crc[data_len + 1] << 8) |
(plaintext_with_crc[data_len + 2] << 16) |
(plaintext_with_crc[data_len + 3] << 24);
if (expected_crc != received_crc) {
return SC_ERR_CRC_FAILED;
}
/* Копируем данные без CRC32 */
memcpy(plaintext, plaintext_with_crc, data_len);
*plaintext_len = data_len;
ctx->rx_counter++;
return SC_OK;
}
sc_status_t sc_compute_public_key_from_private(const uint8_t *private_key, uint8_t *public_key) {
if (!private_key || !public_key) {
return SC_ERR_INVALID_ARG;
}
if (!curve) {
curve = uECC_secp256r1();
}
if (!uECC_compute_public_key(private_key, public_key, curve)) {
return SC_ERR_CRYPTO;
}
return SC_OK;
}

68
src/secure_channel.h1

@ -1,68 +0,0 @@
// secure_channel.h
#ifndef SECURE_CHANNEL_H
#define SECURE_CHANNEL_H
#include <stdint.h>
#include <stddef.h>
// Размеры ключей
#define SC_PRIVKEY_SIZE 32
#define SC_PUBKEY_SIZE 64
#define SC_HASH_SIZE 32
#define SC_NONCE_SIZE 13 // CCM requires exactly 13 bytes
#define SC_SHARED_SECRET_SIZE SC_HASH_SIZE
#define SC_SESSION_KEY_SIZE 16
#define SC_TAG_SIZE 8
#define SC_CRC32_SIZE 4
// Коды возврата
#define SC_OK 0
#define SC_ERR_INVALID_ARG -1
#define SC_ERR_CRYPTO -2
#define SC_ERR_NOT_INITIALIZED -3
#define SC_ERR_AUTH_FAILED -4
#define SC_ERR_CRC_FAILED -5
#define SC_PEER_PUBKEY_BIN 0
#define SC_PEER_PUBKEY_HEX 1
// Типы
typedef int sc_status_t;
typedef struct secure_channel sc_context_t;
struct SC_MYKEYS {
/* Локальные ключи */
uint8_t private_key[SC_PRIVKEY_SIZE];
uint8_t public_key[SC_PUBKEY_SIZE];
};
// Контекст защищенного канала
struct secure_channel {
struct SC_MYKEYS* pk;
/* Ключи пира (после key exchange) */
uint8_t peer_public_key[SC_PUBKEY_SIZE];
uint8_t session_key[SC_SESSION_KEY_SIZE]; /* Derived session key */
/* Nonces для отправки и приема */
uint8_t send_nonce[SC_NONCE_SIZE];
uint8_t recv_nonce[SC_NONCE_SIZE];
uint8_t initialized;
uint8_t peer_key_set;
uint8_t session_ready;
uint64_t tx_counter;
uint64_t rx_counter;
};
// Функции инициализации
sc_status_t sc_init_ctx(sc_context_t *ctx, struct SC_MYKEYS *mykeys);
sc_status_t sc_generate_keypair(struct SC_MYKEYS *keys);
sc_status_t sc_init_local_keys(struct SC_MYKEYS *mykeys, const char *public_key, const char *private_key);
sc_status_t sc_set_peer_public_key(sc_context_t *ctx, const char *peer_public_key, int mode);// mode: 0-bin 1-hex key format
sc_status_t sc_compute_public_key_from_private(const uint8_t *private_key, uint8_t *public_key);
// Криптографические операции
sc_status_t sc_encrypt(sc_context_t *ctx, const uint8_t *plaintext, size_t plaintext_len, uint8_t *ciphertext, size_t *ciphertext_len);
sc_status_t sc_decrypt(sc_context_t *ctx, const uint8_t *ciphertext, size_t ciphertext_len, uint8_t *plaintext, size_t *plaintext_len);
#endif // SECURE_CHANNEL_H

202
tests/test_etcp_100_packets.c

@ -25,12 +25,24 @@ static struct UTUN_INSTANCE* client_instance = NULL;
static int test_completed = 0; static int test_completed = 0;
static void* packet_timeout_id = NULL; static void* packet_timeout_id = NULL;
// Test statistics // Test statistics - forward direction (client -> server)
static int packets_sent = 0; static int packets_sent_fwd = 0;
static int packets_received = 0; static int packets_received_fwd = 0;
static int current_packet_seq_fwd = 0;
static uint8_t received_packets_fwd[TOTAL_PACKETS];
// Test statistics - backward direction (server -> client)
static int packets_sent_back = 0;
static int packets_received_back = 0;
static int current_packet_seq_back = 0;
static uint8_t received_packets_back[TOTAL_PACKETS];
static uint8_t packet_buffer[PACKET_SIZE]; static uint8_t packet_buffer[PACKET_SIZE];
static int current_packet_seq = 0;
static uint8_t received_packets[TOTAL_PACKETS]; // Timing variables
static struct timespec start_time_fwd, end_time_fwd;
static struct timespec start_time_back, end_time_back;
static int phase = 0; // 0 = connecting, 1 = forward transfer, 2 = backward transfer
// Function to generate packet data // Function to generate packet data
static void generate_packet_data(int seq, uint8_t* buffer, int size) { static void generate_packet_data(int seq, uint8_t* buffer, int size) {
@ -55,39 +67,82 @@ static int is_connection_established(struct UTUN_INSTANCE* inst) {
return 0; return 0;
} }
// Send packets while queue has space // Send packets from client to server (forward direction)
static void send_packets(void) { static void send_packets_fwd(void) {
if (!client_instance || packets_sent >= TOTAL_PACKETS) return; if (!client_instance || packets_sent_fwd >= TOTAL_PACKETS) return;
struct ETCP_CONN* conn = client_instance->connections; struct ETCP_CONN* conn = client_instance->connections;
if (!conn || !conn->input_queue) return; if (!conn || !conn->input_queue) return;
// Start timing on first packet
if (packets_sent_fwd == 0) {
clock_gettime(CLOCK_MONOTONIC, &start_time_fwd);
phase = 1;
printf("Starting forward transfer (client -> server)...\n");
}
// Send while queue has space // Send while queue has space
while (packets_sent < TOTAL_PACKETS) { while (packets_sent_fwd < TOTAL_PACKETS) {
int queue_count = queue_entry_count(conn->input_queue); int queue_count = queue_entry_count(conn->input_queue);
if (queue_count >= MAX_QUEUE_SIZE) { if (queue_count >= MAX_QUEUE_SIZE) {
// Queue full, stop sending
break; break;
} }
generate_packet_data(current_packet_seq, packet_buffer, PACKET_SIZE); generate_packet_data(current_packet_seq_fwd, packet_buffer, PACKET_SIZE);
if (etcp_send(conn, packet_buffer, PACKET_SIZE) == 0) { if (etcp_send(conn, packet_buffer, PACKET_SIZE) == 0) {
packets_sent++; packets_sent_fwd++;
current_packet_seq++; current_packet_seq_fwd++;
} else { } else {
break; // Send failed break;
} }
} }
if (packets_sent >= TOTAL_PACKETS) { if (packets_sent_fwd >= TOTAL_PACKETS) {
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "All %d packets queued for sending", TOTAL_PACKETS); DEBUG_INFO(DEBUG_CATEGORY_ETCP, "All %d forward packets queued", TOTAL_PACKETS);
} }
} }
// Check received packets // Send packets from server to client (backward direction)
static void check_received_packets(void) { static void send_packets_back(void) {
if (!server_instance || packets_sent_back >= TOTAL_PACKETS) return;
struct ETCP_CONN* conn = server_instance->connections;
if (!conn || !conn->input_queue) return;
// Start timing on first packet
if (packets_sent_back == 0) {
clock_gettime(CLOCK_MONOTONIC, &start_time_back);
phase = 2;
printf("Starting backward transfer (server -> client)...\n");
}
// Send while queue has space
while (packets_sent_back < TOTAL_PACKETS) {
int queue_count = queue_entry_count(conn->input_queue);
if (queue_count >= MAX_QUEUE_SIZE) {
break;
}
generate_packet_data(current_packet_seq_back, packet_buffer, PACKET_SIZE);
if (etcp_send(conn, packet_buffer, PACKET_SIZE) == 0) {
packets_sent_back++;
current_packet_seq_back++;
} else {
break;
}
}
if (packets_sent_back >= TOTAL_PACKETS) {
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "All %d backward packets queued", TOTAL_PACKETS);
}
}
// Check packets received by server (forward direction)
static void check_received_packets_fwd(void) {
if (!server_instance) return; if (!server_instance) return;
struct ETCP_CONN* conn = server_instance->connections; struct ETCP_CONN* conn = server_instance->connections;
@ -96,26 +151,63 @@ static void check_received_packets(void) {
struct ETCP_FRAGMENT* pkt; struct ETCP_FRAGMENT* pkt;
while ((pkt = queue_data_get(conn->output_queue)) != NULL) { while ((pkt = queue_data_get(conn->output_queue)) != NULL) {
if (pkt->ll.size >= PACKET_SIZE) { if (pkt->ll.size >= PACKET_SIZE) {
int seq = pkt->pkt_data[0]; int seq = pkt->ll.dgram[0];
uint8_t expected[PACKET_SIZE];
generate_packet_data(seq, expected, PACKET_SIZE);
if (memcmp(pkt->ll.dgram, expected, PACKET_SIZE) == 0) {
if (seq >= 0 && seq < TOTAL_PACKETS) {
received_packets_fwd[seq] = 1;
}
packets_received_fwd++;
}
}
if (pkt->ll.dgram) {
memory_pool_free(conn->instance->data_pool, pkt->ll.dgram);
}
queue_data_free(pkt);
}
}
// Check packets received by client (backward direction)
static void check_received_packets_back(void) {
if (!client_instance) return;
struct ETCP_CONN* conn = client_instance->connections;
if (!conn || !conn->output_queue) return;
struct ETCP_FRAGMENT* pkt;
while ((pkt = queue_data_get(conn->output_queue)) != NULL) {
if (pkt->ll.size >= PACKET_SIZE) {
int seq = pkt->ll.dgram[0];
uint8_t expected[PACKET_SIZE]; uint8_t expected[PACKET_SIZE];
generate_packet_data(seq, expected, PACKET_SIZE); generate_packet_data(seq, expected, PACKET_SIZE);
if (memcmp(pkt->pkt_data, expected, PACKET_SIZE) == 0) { if (memcmp(pkt->ll.dgram, expected, PACKET_SIZE) == 0) {
if (seq >= 0 && seq < TOTAL_PACKETS) { if (seq >= 0 && seq < TOTAL_PACKETS) {
received_packets[seq] = 1; received_packets_back[seq] = 1;
} }
packets_received++; packets_received_back++;
} }
} }
if (pkt->pkt_data) { if (pkt->ll.dgram) {
memory_pool_free(conn->instance->data_pool, pkt->pkt_data); memory_pool_free(conn->instance->data_pool, pkt->ll.dgram);
} }
queue_data_free(pkt); queue_data_free(pkt);
} }
} }
// Calculate time difference in milliseconds
static double time_diff_ms(struct timespec* start, struct timespec* end) {
double seconds = end->tv_sec - start->tv_sec;
double nanoseconds = end->tv_nsec - start->tv_nsec;
return (seconds * 1000.0) + (nanoseconds / 1000000.0);
}
// Monitor function // Monitor function
static void monitor_and_send(void* arg) { static void monitor_and_send(void* arg) {
(void)arg; (void)arg;
@ -135,16 +227,43 @@ static void monitor_and_send(void* arg) {
} }
if (connection_checked) { if (connection_checked) {
// Send packets if queue has space // Phase 1: Forward transfer (client -> server)
send_packets(); if (packets_sent_fwd < TOTAL_PACKETS || packets_received_fwd < TOTAL_PACKETS) {
send_packets_fwd();
// Check received packets check_received_packets_fwd();
check_received_packets();
// Check if forward phase completed
if (packets_sent_fwd >= TOTAL_PACKETS && packets_received_fwd >= TOTAL_PACKETS) {
if (end_time_fwd.tv_sec == 0) {
clock_gettime(CLOCK_MONOTONIC, &end_time_fwd);
double duration = time_diff_ms(&start_time_fwd, &end_time_fwd);
printf("✅ Forward transfer completed: %d/%d packets in %.2f ms\n",
packets_received_fwd, TOTAL_PACKETS, duration);
}
}
}
// Phase 2: Backward transfer (server -> client)
else if (packets_sent_back < TOTAL_PACKETS || packets_received_back < TOTAL_PACKETS) {
send_packets_back();
check_received_packets_back();
}
// Check completion // Check completion
if (packets_sent >= TOTAL_PACKETS && packets_received >= TOTAL_PACKETS) { else {
clock_gettime(CLOCK_MONOTONIC, &end_time_back);
double duration_back = time_diff_ms(&start_time_back, &end_time_back);
double duration_total = time_diff_ms(&start_time_fwd, &end_time_back);
printf("✅ Backward transfer completed: %d/%d packets in %.2f ms\n",
packets_received_back, TOTAL_PACKETS, duration_back);
test_completed = 1; test_completed = 1;
printf("\n=== SUCCESS: All %d packets transmitted! ===\n", TOTAL_PACKETS); printf("\n=== SUCCESS: Bidirectional transfer completed! ===\n");
printf("Forward (client->server): %d/%d packets in %.2f ms\n",
packets_received_fwd, TOTAL_PACKETS, time_diff_ms(&start_time_fwd, &end_time_fwd));
printf("Backward (server->client): %d/%d packets in %.2f ms\n",
packets_received_back, TOTAL_PACKETS, duration_back);
printf("Total time: %.2f ms\n", duration_total);
if (packet_timeout_id) { if (packet_timeout_id) {
uasync_cancel_timeout(server_instance->ua, packet_timeout_id); uasync_cancel_timeout(server_instance->ua, packet_timeout_id);
packet_timeout_id = NULL; packet_timeout_id = NULL;
@ -163,7 +282,10 @@ static void test_timeout(void* arg) {
(void)arg; (void)arg;
if (!test_completed) { if (!test_completed) {
printf("\n=== TIMEOUT ===\n"); printf("\n=== TIMEOUT ===\n");
printf("Sent: %d/%d, Received: %d/%d\n", packets_sent, TOTAL_PACKETS, packets_received, TOTAL_PACKETS); printf("Forward: Sent: %d/%d, Received: %d/%d\n",
packets_sent_fwd, TOTAL_PACKETS, packets_received_fwd, TOTAL_PACKETS);
printf("Backward: Sent: %d/%d, Received: %d/%d\n",
packets_sent_back, TOTAL_PACKETS, packets_received_back, TOTAL_PACKETS);
test_completed = 2; test_completed = 2;
if (packet_timeout_id) { if (packet_timeout_id) {
uasync_cancel_timeout(server_instance->ua, packet_timeout_id); uasync_cancel_timeout(server_instance->ua, packet_timeout_id);
@ -173,9 +295,10 @@ static void test_timeout(void* arg) {
} }
int main() { int main() {
printf("=== ETCP 100 Packets Test ===\n\n"); printf("=== ETCP 100 Packets Bidirectional Test ===\n\n");
memset(received_packets, 0, sizeof(received_packets)); memset(received_packets_fwd, 0, sizeof(received_packets_fwd));
memset(received_packets_back, 0, sizeof(received_packets_back));
debug_config_init(); debug_config_init();
debug_set_level(DEBUG_LEVEL_DEBUG); debug_set_level(DEBUG_LEVEL_DEBUG);
@ -201,7 +324,7 @@ int main() {
} }
printf("✅ Client ready\n\n"); printf("✅ Client ready\n\n");
printf("Sending %d packets (max queue size: %d)...\n", TOTAL_PACKETS, MAX_QUEUE_SIZE); printf("Sending %d packets in each direction (max queue size: %d)...\n", TOTAL_PACKETS, MAX_QUEUE_SIZE);
packet_timeout_id = uasync_set_timeout(server_ua, 500, NULL, monitor_and_send); packet_timeout_id = uasync_set_timeout(server_ua, 500, NULL, monitor_and_send);
void* global_timeout_id = uasync_set_timeout(server_ua, TEST_TIMEOUT_MS, NULL, test_timeout); void* global_timeout_id = uasync_set_timeout(server_ua, TEST_TIMEOUT_MS, NULL, test_timeout);
@ -229,11 +352,14 @@ int main() {
if (test_completed == 1) { if (test_completed == 1) {
printf("\n=== TEST PASSED ===\n"); printf("\n=== TEST PASSED ===\n");
printf("✅ All %d packets transmitted\n", TOTAL_PACKETS); printf("✅ All %d packets transmitted in each direction\n", TOTAL_PACKETS);
return 0; return 0;
} else { } else {
printf("\n=== TEST FAILED ===\n"); printf("\n=== TEST FAILED ===\n");
printf("❌ Sent: %d/%d, Received: %d/%d\n", packets_sent, TOTAL_PACKETS, packets_received, TOTAL_PACKETS); printf("❌ Forward: Sent: %d/%d, Received: %d/%d\n",
packets_sent_fwd, TOTAL_PACKETS, packets_received_fwd, TOTAL_PACKETS);
printf("❌ Backward: Sent: %d/%d, Received: %d/%d\n",
packets_sent_back, TOTAL_PACKETS, packets_received_back, TOTAL_PACKETS);
return 1; return 1;
} }
} }

252
tests/test_etcp_simple_traffic.c

@ -17,7 +17,6 @@
#define TEST_TIMEOUT_MS 5000 // 5 seconds for packet transmission #define TEST_TIMEOUT_MS 5000 // 5 seconds for packet transmission
#define PACKET_SIZE 100 // Test packet size #define PACKET_SIZE 100 // Test packet size
#define MAX_QUEUE_ENTRIES 100 // Max entries in queue before considering it full
static struct UTUN_INSTANCE* server_instance = NULL; static struct UTUN_INSTANCE* server_instance = NULL;
static struct UTUN_INSTANCE* client_instance = NULL; static struct UTUN_INSTANCE* client_instance = NULL;
@ -29,15 +28,37 @@ static uint8_t test_packet_data[PACKET_SIZE];
static int packet_sent = 0; static int packet_sent = 0;
static int packet_received = 0; static int packet_received = 0;
// Function to check if connection is established // Function to check if connection is established - UPDATED WITH DEBUG
static int is_connection_established(struct UTUN_INSTANCE* inst) { static int is_connection_established(struct UTUN_INSTANCE* inst) {
if (!inst) return 0; printf("[DEBUG] is_connection_established: Checking instance %p\n", inst);
fflush(stdout);
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "is_connection_established: Checking instance %p", inst);
if (!inst) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "is_connection_established: Instance is NULL");
return 0;
}
printf("[DEBUG] is_connection_established: Instance has %d connections\n", inst->connections_count);
fflush(stdout);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "is_connection_established: Instance has %d connections", inst->connections_count);
struct ETCP_CONN* conn = inst->connections; struct ETCP_CONN* conn = inst->connections;
while (conn) { while (conn) {
printf("[DEBUG] is_connection_established: Checking connection %p\n", conn);
fflush(stdout);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "is_connection_established: Checking connection %p", conn);
struct ETCP_LINK* link = conn->links; struct ETCP_LINK* link = conn->links;
while (link) { while (link) {
printf("[DEBUG] is_connection_established: Link %p - initialized=%d, is_server=%d\n",
link, link->initialized, link->is_server);
fflush(stdout);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "is_connection_established: Link %p - initialized=%d, is_server=%d",
link, link->initialized, link->is_server);
if (link->initialized) { if (link->initialized) {
printf("[DEBUG] is_connection_established: FOUND initialized link!\n");
fflush(stdout);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "is_connection_established: FOUND initialized link!");
return 1; return 1;
} }
link = link->next; link = link->next;
@ -45,12 +66,18 @@ static int is_connection_established(struct UTUN_INSTANCE* inst) {
conn = conn->next; conn = conn->next;
} }
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "is_connection_established: No initialized links found");
return 0; return 0;
} }
// Function to send test packet via normalizer packer input queue // Function to send test packet to client's input queue
static void send_test_packet(void) { static void send_test_packet(void) {
printf("[DEBUG] send_test_packet: ENTERING\n");
fflush(stdout);
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "send_test_packet: ENTERING");
if (!client_instance || packet_sent) { if (!client_instance || packet_sent) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "send_test_packet: SKIPPING (client_instance=%p, packet_sent=%d)", client_instance, packet_sent);
return; return;
} }
@ -60,15 +87,21 @@ static void send_test_packet(void) {
return; return;
} }
if (!conn->normalizer || !conn->normalizer->packer || !conn->normalizer->packer->input) { printf("[DEBUG] send_test_packet: Found connection %p\n", conn);
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "send_test_packet: Client normalizer/packer/input not initialized"); fflush(stdout);
return; DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "send_test_packet: Found connection %p", conn);
} printf("[DEBUG] send_test_packet: connection input_queue=%p\n", conn->input_queue);
fflush(stdout);
// Check queue fullness before writing DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "send_test_packet: connection input_queue=%p", conn->input_queue);
size_t queue_count = queue_entry_count(conn->normalizer->packer->input); printf("[DEBUG] send_test_packet: connection output_queue=%p\n", conn->output_queue);
if (queue_count >= MAX_QUEUE_ENTRIES) { fflush(stdout);
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "send_test_packet: Packer input queue is full (%zu entries), waiting...", queue_count); DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "send_test_packet: connection output_queue=%p", conn->output_queue);
printf("[DEBUG] send_test_packet: connection next_tx_id=%u\n", conn->next_tx_id);
fflush(stdout);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "send_test_packet: connection next_tx_id=%u", conn->next_tx_id);
if (!conn->input_queue) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "send_test_packet: Client input queue is NULL");
return; return;
} }
@ -77,71 +110,86 @@ static void send_test_packet(void) {
test_packet_data[i] = (uint8_t)(i % 256); test_packet_data[i] = (uint8_t)(i % 256);
} }
// Allocate memory for packet using queue API DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "send_test_packet: Created test data, first byte=%02X, last byte=%02X",
void* entry_data = queue_data_new(PACKET_SIZE); test_packet_data[0], test_packet_data[PACKET_SIZE-1]);
if (!entry_data) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "send_test_packet: Failed to allocate queue data");
return;
}
memcpy(entry_data, test_packet_data, PACKET_SIZE); // Use new etcp_send function - much simpler API
printf("[DEBUG] send_test_packet: Calling etcp_send with len=%d\n", PACKET_SIZE);
fflush(stdout);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "send_test_packet: Calling etcp_send with len=%d", PACKET_SIZE);
int result = etcp_send(conn, test_packet_data, PACKET_SIZE);
printf("[DEBUG] send_test_packet: etcp_send returned %d\n", result);
fflush(stdout);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "send_test_packet: etcp_send returned %d", result);
// Set the length in ll_entry if (result != 0) {
struct ll_entry* entry = (struct ll_entry *)((char *)entry_data - offsetof(struct ll_entry, data)); DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "send_test_packet: Failed to send packet via etcp_send");
entry->len = PACKET_SIZE;
// Write to packer input queue
if (queue_data_put(conn->normalizer->packer->input, entry_data, 0) < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "send_test_packet: Failed to put data in packer input queue");
queue_data_free(entry_data);
return; return;
} }
packet_sent = 1; packet_sent = 1;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "send_test_packet: Test packet sent to packer input queue (%d bytes)", PACKET_SIZE); DEBUG_INFO(DEBUG_CATEGORY_ETCP, "send_test_packet: SUCCESS - Test packet sent via etcp_send (%d bytes)", PACKET_SIZE);
} }
// Function to check if packet received in server's unpacker output queue // Function to check if packet received in server's output queue
static void check_packet_received(void) { static void check_packet_received(void) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "check_packet_received: ENTERING");
if (!server_instance || packet_received) { if (!server_instance || packet_received) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "check_packet_received: SKIPPING (server_instance=%p, packet_received=%d)", server_instance, packet_received);
return; return;
} }
struct ETCP_CONN* conn = server_instance->connections; struct ETCP_CONN* conn = server_instance->connections;
if (!conn) { if (!conn) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "check_packet_received: No connection on server");
return; return;
} }
if (!conn->normalizer || !conn->normalizer->unpacker || !conn->normalizer->unpacker->output) { if (!conn->output_queue) {
return; DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "check_packet_received: No output_queue on server connection");
}
// Read from unpacker output queue
void* data = queue_data_get(conn->normalizer->unpacker->output);
if (!data) {
return; return;
} }
// Get the length from ll_entry DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "check_packet_received: Checking server connection %p", conn);
struct ll_entry* entry = (struct ll_entry *)((char *)data - offsetof(struct ll_entry, data)); DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "check_packet_received: Server output_queue count: %d", queue_entry_count(conn->output_queue));
size_t size = entry->len;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "check_packet_received: Found packet in unpacker output queue (size=%zu)", size); // Check if there's any packet in output queue
struct ETCP_FRAGMENT* pkt = queue_data_get(conn->output_queue);
// Verify packet data if (pkt) {
if (size >= PACKET_SIZE && memcmp(data, test_packet_data, PACKET_SIZE) == 0) { DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "check_packet_received: Found packet in output queue");
packet_received = 1;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "check_packet_received: SUCCESS - Packet received and data matches (%zu bytes)", size); // ETCP_FRAGMENT содержит ll_entry в начале, данные в pkt_data
size_t size = pkt->ll.size;
uint8_t* data = pkt->ll.dgram;
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "check_packet_received: Packet size=%zu, expected=%d", size, PACKET_SIZE);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "check_packet_received: First byte received=%02X, expected=%02X",
data[0], test_packet_data[0]);
// Проверяем только данные (первые PACKET_SIZE байт), размер может включать ETCP заголовки
if (size >= PACKET_SIZE && memcmp(data, test_packet_data, PACKET_SIZE) == 0) {
packet_received = 1;
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "check_packet_received: SUCCESS - Packet received in server output queue (%zu bytes), data matches", size);
} else {
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "check_packet_received: Packet received but data differs (expected %d bytes, got %zu)", PACKET_SIZE, size);
}
// Освобождаем ETCP_FRAGMENT и данные
if (pkt->ll.dgram) {
memory_pool_free(conn->instance->data_pool, pkt->ll.dgram);
}
queue_data_free(pkt);
} else { } else {
DEBUG_WARN(DEBUG_CATEGORY_ETCP, "check_packet_received: Packet received but data differs (expected %d bytes, got %zu)", PACKET_SIZE, size); DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "check_packet_received: No packet found in output queue");
} }
// Free the data
queue_data_free(data);
} }
// Monitor connection and send packet when ready // Monitor connection and send packet when ready
static void monitor_and_send(void* arg) { static void monitor_and_send(void* arg) {
printf("[DEBUG] monitor_and_send: ENTERING (arg=%p)\n", arg);
fflush(stdout);
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "monitor_and_send: ENTERING (arg=%p)", arg);
(void)arg; (void)arg;
if (test_completed) { if (test_completed) {
@ -157,9 +205,17 @@ static void monitor_and_send(void* arg) {
int server_ready = is_connection_established(server_instance); int server_ready = is_connection_established(server_instance);
int client_ready = is_connection_established(client_instance); int client_ready = is_connection_established(client_instance);
printf("[DEBUG] monitor_and_send: Connection check - server=%d, client=%d\n", server_ready, client_ready);
fflush(stdout);
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "monitor_and_send: Connection check - server=%d, client=%d", server_ready, client_ready);
// We only need client link initialized to send packet
// Server will accept packets via its socket
if (client_ready) { if (client_ready) {
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "monitor_and_send: Client link initialized, ready to send packet"); DEBUG_INFO(DEBUG_CATEGORY_ETCP, "monitor_and_send: Client link initialized, ready to send packet (server=%d, client=%d)", server_ready, client_ready);
connection_checked = 1; connection_checked = 1;
} else {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "monitor_and_send: Waiting for connection... (server=%d, client=%d)", server_ready, client_ready);
} }
} }
@ -203,12 +259,19 @@ static void test_timeout(void* arg) {
} }
int main() { int main() {
printf("=== ETCP Simple Traffic Test (via normalizer queues) ===\n\n"); printf("=== ETCP Simple Traffic Test (Queue-based) ===\n\n");
// Enable debug output // Enable debug output - MAXIMUM DEBUGGING
debug_config_init(); debug_config_init();
debug_set_level(DEBUG_LEVEL_INFO); debug_set_level(DEBUG_LEVEL_TRACE);
debug_set_categories(DEBUG_CATEGORY_ETCP | DEBUG_CATEGORY_CONNECTION); debug_set_categories(DEBUG_CATEGORY_ALL);
debug_enable_function_name(1);
DEBUG_TRACE(DEBUG_CATEGORY_MEMORY, "*************1");
printf("Function names enabled: YES\n");
printf("===========================\n\n");
// Explicitly disable TUN initialization for this test // Explicitly disable TUN initialization for this test
utun_instance_set_tun_init_enabled(0); utun_instance_set_tun_init_enabled(0);
@ -222,13 +285,24 @@ int main() {
return 1; return 1;
} }
// Initialize ETCP connections // Initialize connections and register sockets regardless of TUN state
if (server_instance->tun.fd < 0) {
printf("ℹ️ Server TUN disabled - initializing connections only\n");
}
// Initialize ETCP connections (creates sockets and links)
if (init_connections(server_instance) < 0) { if (init_connections(server_instance) < 0) {
printf("Failed to initialize server connections\n"); printf("Failed to initialize server connections\n");
utun_instance_destroy(server_instance); utun_instance_destroy(server_instance);
return 1; return 1;
} else {
printf("Server connections initialized: sockets=%p, connections=%p\n",
server_instance->etcp_sockets, server_instance->connections);
} }
printf("✅ Server instance initialized successfully (node_id=%llx)\n", (unsigned long long)server_instance->node_id);
printf("Server instance ready (node_id=%llx)\n\n", (unsigned long long)server_instance->node_id); printf("Server instance ready (node_id=%llx)\n\n", (unsigned long long)server_instance->node_id);
// Create client instance // Create client instance
@ -241,16 +315,61 @@ int main() {
return 1; return 1;
} }
// Initialize ETCP connections // Initialize connections and register sockets regardless of TUN state
if (init_connections(client_instance) < 0) { if (client_instance->tun.fd < 0) {
printf("Failed to initialize client connections\n"); printf("ℹ️ Client TUN disabled - initializing connections only\n");
utun_instance_destroy(server_instance);
utun_instance_destroy(client_instance);
return 1;
} }
// Initialize ETCP connections (creates sockets and links)
printf("About to call init_connections() for client instance\n");
fflush(stdout);
int conn_result = init_connections(client_instance);
printf("init_connections() returned: %d\n", conn_result);
fflush(stdout);
if (conn_result < 0) {
printf("Failed to initialize client connections (result=%d)\n", conn_result);
printf("But continuing test to analyze the issue...\n");
fflush(stdout);
// Don't return error - continue to analyze
} else {
printf("Client connections initialized: sockets=%p, connections=%p, count=%d\n",
client_instance->etcp_sockets, client_instance->connections,
client_instance ? client_instance->connections_count : -1);
fflush(stdout);
}
printf("✅ Client instance initialized successfully (node_id=%llx)\n", (unsigned long long)client_instance->node_id);
printf("Client instance ready (node_id=%llx)\n\n", (unsigned long long)client_instance->node_id); printf("Client instance ready (node_id=%llx)\n\n", (unsigned long long)client_instance->node_id);
// Debug: print connection and link info
printf("\n=== Connection Debug ===\n");
struct ETCP_CONN* conn = client_instance->connections;
while (conn) {
printf("Client connection %p: peer_node_id=%llx\n", conn, (unsigned long long)conn->peer_node_id);
struct ETCP_LINK* link = conn->links;
while (link) {
printf(" Link %p: initialized=%d, is_server=%d, remote_addr family=%d, conn->fd=%d\n",
link, link->initialized, link->is_server, link->remote_addr.ss_family, link->conn ? link->conn->fd : -1);
link = link->next;
}
conn = conn->next;
}
conn = server_instance->connections;
while (conn) {
printf("Server connection %p: peer_node_id=%llx\n", conn, (unsigned long long)conn->peer_node_id);
struct ETCP_LINK* link = conn->links;
while (link) {
printf(" Link %p: initialized=%d, is_server=%d\n", link, link->initialized, link->is_server);
link = link->next;
}
conn = conn->next;
}
printf("=== End Debug ===\n\n");
// Start monitoring and packet transmission // Start monitoring and packet transmission
printf("Starting packet transmission test...\n"); printf("Starting packet transmission test...\n");
packet_timeout_id = uasync_set_timeout(server_ua, 500, NULL, monitor_and_send); packet_timeout_id = uasync_set_timeout(server_ua, 500, NULL, monitor_and_send);
@ -309,15 +428,16 @@ int main() {
// Evaluate test result // Evaluate test result
if (test_completed == 1) { if (test_completed == 1) {
printf("\n=== TEST PASSED ===\n"); printf("\n=== TEST PASSED ===\n");
printf("Packet successfully transmitted from client packer input to server unpacker output\n"); printf("✅ Packet successfully transmitted from client input queue to server output queue\n");
printf("ETCP + normalizer integration verified\n"); printf("✅ ETCP connection and queue mechanisms verified\n");
return 0; return 0;
} else if (test_completed == 2) { } else if (test_completed == 2) {
printf("\n=== TEST FAILED: Packet not received within timeout ===\n"); printf("\n=== TEST FAILED: Packet not received within timeout ===\n");
printf("Connection may not have been established\n"); printf("❌ Connection may not have been established\n");
printf("❌ Check if UDP sockets were created and bound correctly\n");
return 1; return 1;
} else { } else {
printf("\n=== TEST FAILED: Unknown error ===\n"); printf("\n=== TEST FAILED: Unknown error ===\n");
return 1; return 1;
} }
} }
Loading…
Cancel
Save