Browse Source

queue: add variable size hash

bgp: add to headers new node info
nodeinfo-routing-update
Evgeny 6 months ago
parent
commit
097e1acc6d
  1. 202
      lib/ll_queue.c
  2. 67
      lib/ll_queue.h
  3. 4
      src/config_parser.c
  4. 1
      src/config_parser.h
  5. 26
      src/etcp.c
  6. 6
      src/etcp.h
  7. 2
      src/etcp_api.c
  8. 27
      src/etcp_connections.c
  9. 5
      src/etcp_connections.h
  10. 6
      src/pkt_normalizer.c
  11. 111
      src/route_bgp.c
  12. 48
      src/route_bgp.h
  13. 28
      src/route_bgp.txt
  14. 2
      src/routing.c
  15. 8
      src/tun_if.c
  16. 23
      src/utun_instance.c
  17. 3
      src/utun_instance.h
  18. 2
      tests/debug_simple.c
  19. 4
      tests/test_intensive_memory_pool.c
  20. 4
      tests/test_intensive_memory_pool_new.c
  21. 27
      tests/test_ll_queue.c
  22. 2
      tests/test_memory_pool_and_config.c
  23. 2
      tests/test_pkt_normalizer_standalone.c

202
lib/ll_queue.c

@ -223,18 +223,22 @@ void queue_entry_free(struct ll_entry* entry) {
// Внутренняя функция добавления в хеш-таблицу // Внутренняя функция добавления в хеш-таблицу
static void add_to_hash(struct ll_queue* q, struct ll_entry* entry) { static void add_to_hash(struct ll_queue* q, struct ll_entry* entry) {
if (!q || q->hash_size == 0 || !entry) return; if (!q || !entry) return;
if (q->hash_size == 0 || entry->index_size == 0) {
DEBUG_ERROR(DEBUG_CATEGORY_LL_QUEUE, "[%s] Empty hash_size=%d or index_size=%d",q->name, q->hash_size, entry->index_size);
return;
}
size_t slot = entry->id % q->hash_size; uint32_t slot = entry->index_hash % q->hash_size;
entry->hash_next = q->hash_table[slot]; entry->hash_next = q->hash_table[slot];
q->hash_table[slot] = entry; q->hash_table[slot] = entry;
} }
// Внутренняя функция удаления из хеш-таблицы
static void remove_from_hash(struct ll_queue* q, struct ll_entry* entry) { static void remove_from_hash(struct ll_queue* q, struct ll_entry* entry) {
if (!q || q->hash_size == 0 || !entry) return; if (!q || q->hash_size == 0 || !entry) return;
size_t slot = entry->id % q->hash_size; uint32_t slot = entry->index_hash % q->hash_size; // теперь используем index_hash
struct ll_entry** ptr = &q->hash_table[slot]; struct ll_entry** ptr = &q->hash_table[slot];
while (*ptr) { while (*ptr) {
if (*ptr == entry) { if (*ptr == entry) {
@ -246,6 +250,20 @@ static void remove_from_hash(struct ll_queue* q, struct ll_entry* entry) {
} }
} }
static uint32_t make_hash(const void* data, uint16_t len) {// алгоритм FNV-1a
if (len == 0 || data == NULL) return 0;
const uint8_t* p = (const uint8_t*)data;
uint32_t hash = 0x811C9DC5u; // FNV-1a offset basis (32-bit)
for (uint16_t i = 0; i < len; ++i) {
hash ^= (uint32_t)p[i]; // XOR с байтом
hash *= 0x01000193u; // FNV-1a prime
}
return hash;
}
// Проверить и запустить ожидающие коллбэки // Проверить и запустить ожидающие коллбэки
static void check_waiters(struct ll_queue* q) { static void check_waiters(struct ll_queue* q) {
if (!q) return; if (!q) return;
@ -267,118 +285,185 @@ static void check_waiters(struct ll_queue* q) {
} }
} }
int queue_data_put(struct ll_queue* q, struct ll_entry* entry, uint32_t id) { int queue_data_put(struct ll_queue* q, struct ll_entry* entry) {
if (!q || !entry) return -1; if (!q || !entry) return -1;
#ifdef QUEUE_THREAD_CHECK #ifdef QUEUE_THREAD_CHECK
queue_check_thread(q); queue_check_thread(q);
#endif #endif
entry->index_offset = 0;
entry->index_size = 0;
entry->index_hash = 0;
if (q->hash_size > 0) {
DEBUG_ERROR(DEBUG_CATEGORY_LL_QUEUE, "[%s] Put no-hash data to hash queue",q->name);
queue_dgram_free(entry);
queue_entry_free(entry);
return -1;
}
if (q->size_limit >= 0 && q->count >= q->size_limit) {
queue_dgram_free(entry);
queue_entry_free(entry);
return -1;
}
// Добавляем в конец
entry->next = NULL;
entry->prev = q->tail;
if (q->tail) q->tail->next = entry; else q->head = entry;
q->tail = entry;
q->count++;
entry->int_len = entry->len;
q->total_bytes += entry->int_len;
// НЕ вызываем add_to_hash
if (q->count == 1 && !q->callback_suspended && q->callback) {
q->callback(q, q->callback_arg);
}
#ifdef QUEUE_DEBUG #ifdef QUEUE_DEBUG
// queue_check_consistency(q);// !!!! for debug - BEFORE callback queue_check_consistency(q);
#endif
return 0;
}
int queue_data_put_with_index(struct ll_queue* q, struct ll_entry* entry, uint16_t index_offset, uint16_t index_size) {
if (!q || !entry) return -1;
#ifdef QUEUE_THREAD_CHECK
queue_check_thread(q);
#endif #endif
entry->id = id; if (index_size == 0 || (size_t)index_offset + index_size > (size_t)entry->size) {
DEBUG_ERROR(DEBUG_CATEGORY_LL_QUEUE, "[%s] invalid index: offset=%u size=%u (data_size=%u)",q->name, index_offset, index_size, entry->size);
queue_dgram_free(entry);
queue_entry_free(entry);
return -1;
}
entry->index_offset = index_offset;
entry->index_size = index_size;
entry->index_hash = make_hash(entry->data + index_offset, index_size);
// Проверить лимит размера
if (q->size_limit >= 0 && q->count >= q->size_limit) { if (q->size_limit >= 0 && q->count >= q->size_limit) {
queue_dgram_free(entry); queue_dgram_free(entry);
queue_entry_free(entry); // Освободить элемент если превышен лимит queue_entry_free(entry);
return -1; return -1;
} }
// Добавить в конец // Добавляем в конец
entry->next = NULL; entry->next = NULL;
entry->prev = q->tail; entry->prev = q->tail;
if (q->tail) { if (q->tail) q->tail->next = entry; else q->head = entry;
q->tail->next = entry;
} else {
q->head = entry;
}
q->tail = entry; q->tail = entry;
q->count++; q->count++;
entry->int_len=entry->len; entry->int_len = entry->len;
q->total_bytes += entry->int_len; q->total_bytes += entry->int_len;
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);
add_to_hash(q, entry); add_to_hash(q, entry);
// DEBUG_DEBUG(DEBUG_CATEGORY_LL_QUEUE, "queue_data_put: added entry %p (id=%u), count=%d", entry, id, q->count); if (q->count == 1 && !q->callback_suspended && q->callback) {
q->callback(q, q->callback_arg);
}
// ВАЖНО: проверка консистентности ДО коллбэка, так как коллбэк может модифицировать очередь
#ifdef QUEUE_DEBUG #ifdef QUEUE_DEBUG
queue_check_consistency(q);// !!!! for debug - BEFORE callback queue_check_consistency(q);
#endif #endif
return 0;
}
int queue_data_put_first(struct ll_queue* q, struct ll_entry* entry) {
if (!q || !entry) return -1;
#ifdef QUEUE_THREAD_CHECK
queue_check_thread(q);
#endif
entry->index_offset = 0;
entry->index_size = 0;
entry->index_hash = 0;
if (q->hash_size > 0) {
DEBUG_ERROR(DEBUG_CATEGORY_LL_QUEUE, "[%s] Put no-hash data to hash queue",q->name);
queue_dgram_free(entry);
queue_entry_free(entry);
return -1;
}
if (q->size_limit >= 0 && q->count >= q->size_limit) {
queue_dgram_free(entry);
queue_entry_free(entry); // как было в оригинале
return -1;
}
entry->next = q->head;
entry->prev = NULL;
if (q->head) q->head->prev = entry; else q->tail = entry;
q->head = entry;
q->count++;
entry->int_len = entry->len;
q->total_bytes += entry->int_len;
// Если очередь была пуста и коллбэки разрешены - вызвать коллбэк
if (q->count == 1 && !q->callback_suspended && q->callback) { if (q->count == 1 && !q->callback_suspended && q->callback) {
q->callback(q, q->callback_arg); q->callback(q, q->callback_arg);
} }
// Проверить ожидающие коллбэки (надо только при заборе из очереди)
// check_waiters(q);
#ifdef QUEUE_DEBUG #ifdef QUEUE_DEBUG
queue_check_consistency(q);// !!!! for debug - AFTER callback queue_check_consistency(q);
#endif #endif
return 0; return 0;
} }
int queue_data_put_first(struct ll_queue* q, struct ll_entry* entry, uint32_t id) { // С хешем для put_first
int queue_data_put_first_with_index(struct ll_queue* q, struct ll_entry* entry, uint16_t index_offset, uint16_t index_size) {
if (!q || !entry) return -1; if (!q || !entry) return -1;
#ifdef QUEUE_THREAD_CHECK #ifdef QUEUE_THREAD_CHECK
queue_check_thread(q); queue_check_thread(q);
#endif #endif
entry->id = id; if (index_size == 0 || (size_t)index_offset + index_size > (size_t)entry->size) {
DEBUG_ERROR(DEBUG_CATEGORY_LL_QUEUE, "[%s] invalid index: offset=%u size=%u (data_size=%u)",q->name, index_offset, index_size, entry->size);
queue_dgram_free(entry);
queue_entry_free(entry);
return -1;
}
entry->index_offset = index_offset;
entry->index_size = index_size;
entry->index_hash = make_hash(entry->data + index_offset, index_size);
// Проверить лимит размера
if (q->size_limit >= 0 && q->count >= q->size_limit) { if (q->size_limit >= 0 && q->count >= q->size_limit) {
queue_entry_free(entry); // Освободить элемент если превышен лимит queue_dgram_free(entry);
queue_entry_free(entry);
return -1; return -1;
} }
// Добавить в начало
entry->next = q->head; entry->next = q->head;
entry->prev = NULL; entry->prev = NULL;
if (q->head) { if (q->head) q->head->prev = entry; else q->tail = entry;
q->head->prev = entry;
} else {
q->tail = entry;
}
q->head = entry; q->head = entry;
q->count++; q->count++;
entry->int_len=entry->len; entry->int_len = entry->len;
q->total_bytes += entry->int_len; q->total_bytes += entry->int_len;
add_to_hash(q, entry); add_to_hash(q, entry);
// DEBUG_DEBUG(DEBUG_CATEGORY_LL_QUEUE, "queue_data_put_first: added entry %p (id=%u), count=%d", entry, id, q->count);
// ВАЖНО: проверка консистентности ДО коллбэка
#ifdef QUEUE_DEBUG
queue_check_consistency(q);// !!!! for debug - BEFORE callback
#endif
// Если очередь была пуста и коллбэки разрешены - вызвать коллбэк
if (q->count == 1 && !q->callback_suspended && q->callback) { if (q->count == 1 && !q->callback_suspended && q->callback) {
q->callback(q, q->callback_arg); q->callback(q, q->callback_arg);
} }
// Проверить ожидающие коллбэки
// check_waiters(q);
#ifdef QUEUE_DEBUG #ifdef QUEUE_DEBUG
queue_check_consistency(q);// !!!! for debug queue_check_consistency(q);
#endif #endif
return 0; return 0;
} }
struct ll_entry* queue_data_get(struct ll_queue* q) { struct ll_entry* queue_data_get(struct ll_queue* q) {
if (!q || !q->head) return NULL; if (!q || !q->head) return NULL;
@ -496,19 +581,22 @@ void queue_cancel_wait(struct ll_queue* q, struct queue_waiter* waiter) {
// ==================== Поиск и удаление по ID ==================== // ==================== Поиск и удаление по ID ====================
struct ll_entry* queue_find_data_by_id(struct ll_queue* q, uint32_t id) { struct ll_entry* queue_find_data_by_index(struct ll_queue* q, const void* index_key, uint16_t index_size) {
if (!q || q->hash_size == 0 || !q->hash_table) return NULL; if (!q || q->hash_size == 0 || !q->hash_table || index_size == 0 || !index_key) return NULL;
uint32_t hash_val = make_hash(index_key, index_size);
size_t slot = id % q->hash_size; uint32_t slot = hash_val % q->hash_size;
struct ll_entry* entry = q->hash_table[slot]; struct ll_entry* entry = q->hash_table[slot];
while (entry) { while (entry) {
if (entry->id == id) { if (entry->index_size == index_size &&
(size_t)entry->index_offset + index_size <= (size_t)entry->size &&
memcmp(entry->data + entry->index_offset, index_key, index_size) == 0) {
return entry; return entry;
} }
entry = entry->hash_next; entry = entry->hash_next;
} }
return NULL; return NULL;
} }

67
lib/ll_queue.h

@ -54,19 +54,21 @@ struct ll_queue;
*/ */
struct ll_entry { struct ll_entry {
char* name; char* name;
struct ll_entry* next; ///< Следующий элемент в очереди struct ll_entry* next; // Следующий элемент в очереди
struct ll_entry* prev; ///< Предыдущий элемент в очереди struct ll_entry* prev; // Предыдущий элемент в очереди
uint16_t size; ///< Размер пользовательского буфера data[] (байт) uint16_t size; // Размер пользовательского буфера data[] (байт)
uint16_t len; ///< Актуальная длина данных в dgram uint16_t len; // Актуальная длина данных в dgram
uint16_t memlen; ///< Выделенный размер под dgram uint16_t memlen; // Выделенный размер под dgram
uint16_t int_len; ///< Внутреннее (не использовать) uint16_t int_len; // Внутреннее (не использовать)
uint8_t* dgram; ///< Указатель на данные пакета uint8_t* dgram; // Указатель на данные пакета
void (*dgram_free_fn)(uint8_t* data); ///< Кастомная функция освобождения dgram void (*dgram_free_fn)(uint8_t* data); // Кастомная функция освобождения dgram
struct memory_pool* dgram_pool; ///< Пул для dgram (если используется) struct memory_pool* dgram_pool; // Пул для dgram (если используется)
uint32_t id; ///< Идентификатор для поиска (задаётся при добавлении) uint16_t index_offset; // Смещение индекса в data[]
struct ll_entry* hash_next; ///< Следующий в хеш-цепочке uint16_t index_size; // Длина индекса (0 = без индекса)
struct memory_pool* pool; ///< Пул, из которого выделен сам entry (NULL = malloc) uint32_t index_hash; // Хеш по первым 4 байтам индекса
uint8_t data[0]; ///< Гибкий массив пользовательских данных (размер = size) struct ll_entry* hash_next; // Следующий в хеш-цепочке
struct memory_pool* pool; // Пул, из которого выделен сам entry (NULL = malloc)
uint8_t data[0]; // Гибкий массив пользовательских данных (размер = size)
}; };
/** /**
@ -107,23 +109,23 @@ struct queue_waiter {
*/ */
struct ll_queue { struct ll_queue {
char* name; char* name;
struct ll_entry* head; ///< Голова очереди (отсюда извлекаем) struct ll_entry* head; // Голова очереди (отсюда извлекаем)
struct ll_entry* tail; ///< Хвост очереди (сюда добавляем) struct ll_entry* tail; // Хвост очереди (сюда добавляем)
int count; ///< Текущее количество элементов int count; // Текущее количество элементов
size_t total_bytes; ///< Суммарный объём данных (сумма int_len) size_t total_bytes; // Суммарный объём данных (сумма int_len)
int size_limit; ///< Максимальное количество элементов (-1 = без лимита) int size_limit; // Максимальное количество элементов (-1 = без лимита)
queue_callback_fn callback; ///< Коллбэк автозабора queue_callback_fn callback; // Коллбэк автозабора
void* callback_arg; ///< Аргумент коллбэка void* callback_arg; // Аргумент коллбэка
int callback_suspended; ///< 1 = коллбэки временно приостановлены int callback_suspended; // 1 = коллбэки временно приостановлены
void* resume_timeout_id; ///< ID таймера uasync для отложенного resume void* resume_timeout_id; // ID таймера uasync для отложенного resume
struct UASYNC* ua; ///< Экземпляр uasync (обязателен для таймеров) struct UASYNC* ua; // Экземпляр uasync (обязателен для таймеров)
struct queue_waiter waiter; ///< Встроенный waiter (только один) struct queue_waiter waiter; // Встроенный waiter (только один)
struct ll_entry** hash_table; ///< Хеш-таблица для поиска по id (если hash_size > 0) struct ll_entry** hash_table; // Хеш-таблица для поиска по id (если hash_size > 0)
size_t hash_size; ///< Размер хеш-таблицы size_t hash_size; // Размер хеш-таблицы
#ifdef QUEUE_THREAD_CHECK #ifdef QUEUE_THREAD_CHECK
#ifdef _WIN32 #ifdef _WIN32
@ -215,7 +217,7 @@ void queue_cancel_wait(struct ll_queue* q, struct queue_waiter* waiter);
* @param id идентификатор для поиска * @param id идентификатор для поиска
* @return 0 — успех, -1 — превышен лимит (элемент освобождён) * @return 0 — успех, -1 — превышен лимит (элемент освобождён)
*/ */
int queue_data_put(struct ll_queue* q, struct ll_entry* entry, uint32_t id); int queue_data_put(struct ll_queue* q, struct ll_entry* entry);
/** /**
* @brief Добавляет элемент в начало очереди (LIFO, высокий приоритет). * @brief Добавляет элемент в начало очереди (LIFO, высокий приоритет).
@ -224,7 +226,12 @@ int queue_data_put(struct ll_queue* q, struct ll_entry* entry, uint32_t id);
* @param id идентификатор * @param id идентификатор
* @return 0 — успех, -1 — превышен лимит (элемент освобождён) * @return 0 — успех, -1 — превышен лимит (элемент освобождён)
*/ */
int queue_data_put_first(struct ll_queue* q, struct ll_entry* entry, uint32_t id); int queue_data_put_first(struct ll_queue* q, struct ll_entry* entry);
int queue_data_put_with_index(struct ll_queue* q, struct ll_entry* entry,
uint16_t index_offset, uint16_t index_size);
int queue_data_put_first_with_index(struct ll_queue* q, struct ll_entry* entry,
uint16_t index_offset, uint16_t index_size);
/** /**
* @brief Извлекает элемент из начала очереди. * @brief Извлекает элемент из начала очереди.
@ -291,7 +298,9 @@ void queue_dgram_free(struct ll_entry* entry);
* @param id искомый идентификатор * @param id искомый идентификатор
* @return элемент или NULL * @return элемент или NULL
*/ */
struct ll_entry* queue_find_data_by_id(struct ll_queue* q, uint32_t id); struct ll_entry* queue_find_data_by_index(struct ll_queue* q,
const void* index_key,
uint16_t index_size);
/** /**
* @brief Удаляет элемент из очереди (не освобождает память). * @brief Удаляет элемент из очереди (не освобождает память).

4
src/config_parser.c

@ -270,6 +270,9 @@ static struct CFG_SERVER* find_server_by_name(struct CFG_SERVER *servers, const
} }
static int parse_global(const char *key, const char *value, struct global_config *global) { static int parse_global(const char *key, const char *value, struct global_config *global) {
if (strcmp(key, "my_node_name") == 0) {
return assign_string(global->name, sizeof(global->name), value);
}
if (strcmp(key, "my_private_key") == 0) { if (strcmp(key, "my_private_key") == 0) {
return assign_string(global->my_private_key_hex, MAX_KEY_LEN, value); return assign_string(global->my_private_key_hex, MAX_KEY_LEN, value);
} }
@ -477,6 +480,7 @@ static struct utun_config* parse_config_internal(FILE *fp, const char *filename)
} }
// Set default values // Set default values
cfg->global.name[0] = '\0';
cfg->global.keepalive_timeout = 2000; // Default 2 seconds cfg->global.keepalive_timeout = 2000; // Default 2 seconds
cfg->global.firewall_rules = NULL; cfg->global.firewall_rules = NULL;
cfg->global.firewall_rule_count = 0; cfg->global.firewall_rule_count = 0;

1
src/config_parser.h

@ -67,6 +67,7 @@ struct CFG_FIREWALL_RULE {
}; };
struct global_config { struct global_config {
char name[16]; // Instance name
char my_private_key_hex[MAX_KEY_LEN]; char my_private_key_hex[MAX_KEY_LEN];
char my_public_key_hex[MAX_KEY_LEN]; char my_public_key_hex[MAX_KEY_LEN];
uint64_t my_node_id; uint64_t my_node_id;

26
src/etcp.c

@ -427,7 +427,7 @@ int etcp_int_send(struct ETCP_CONN* etcp, const void* data, uint16_t len) {
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] created PACKET %p with data %p (len=%zu)", etcp->log_name, pkt, packet_data, len); 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 // Add to input queue - input_queue_cb will process it
if (queue_data_put(etcp->input_queue, (struct ll_entry*)pkt, 0) != 0) { if (queue_data_put(etcp->input_queue, (struct ll_entry*)pkt) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to add to input queue", etcp->log_name); DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to add to input queue", etcp->log_name);
queue_dgram_free(&pkt->ll); queue_dgram_free(&pkt->ll);
queue_entry_free(&pkt->ll); queue_entry_free(&pkt->ll);
@ -536,7 +536,7 @@ static void input_queue_cb(struct ll_queue* q, void* arg) {
} }
// Setup inflight packet (based on protocol.txt) // Setup inflight packet (based on protocol.txt)
memset(p, 0, sizeof(*p)); // memset(p, 0, sizeof(*p));
p->seq = etcp->next_tx_id++; // Assign seq p->seq = etcp->next_tx_id++; // Assign seq
p->state = INFLIGHT_STATE_WAIT_SEND; p->state = INFLIGHT_STATE_WAIT_SEND;
p->last_timestamp = 0; p->last_timestamp = 0;
@ -551,7 +551,7 @@ static void input_queue_cb(struct ll_queue* q, void* arg) {
int len=p->ll.len;// сохраним len int len=p->ll.len;// сохраним len
// Add to send queue // Add to send queue
if (queue_data_put(etcp->input_send_q, &p->ll, p->seq) != 0) { if (queue_data_put_with_index(etcp->input_send_q, &p->ll, 0, 4) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to add packet seq=%u to input_send_q", etcp->log_name, p->seq); DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to add packet seq=%u to input_send_q", etcp->log_name, p->seq);
queue_dgram_free(&p->ll); queue_dgram_free(&p->ll);
queue_entry_free(&p->ll); queue_entry_free(&p->ll);
@ -599,7 +599,7 @@ static void ack_timeout_check(struct ETCP_CONN* etcp) {
// Change state and add to send_q for retransmission // Change state and add to send_q for retransmission
pkt->state = INFLIGHT_STATE_WAIT_SEND; pkt->state = INFLIGHT_STATE_WAIT_SEND;
queue_data_put(etcp->input_send_q, (struct ll_entry*)pkt, pkt->seq); queue_data_put_with_index(etcp->input_send_q, (struct ll_entry*)pkt, 0, 4);
// Update stats // Update stats
etcp->retransmissions_count++; etcp->retransmissions_count++;
@ -767,7 +767,7 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) {
inf_pkt->send_count++; inf_pkt->send_count++;
inf_pkt->state=INFLIGHT_STATE_WAIT_ACK; inf_pkt->state=INFLIGHT_STATE_WAIT_ACK;
queue_data_put(etcp->input_wait_ack, &inf_pkt->ll, inf_pkt->seq);// move dgram to wait_ack queue queue_data_put_with_index(etcp->input_wait_ack, &inf_pkt->ll, 0, 4);// move dgram to wait_ack queue
} }
size_t ack_q_size = queue_entry_count(etcp->ack_q); size_t ack_q_size = queue_entry_count(etcp->ack_q);
@ -889,7 +889,7 @@ void etcp_output_try_assembly(struct ETCP_CONN* etcp) {
// Look for contiguous packets starting from next_expected_id // Look for contiguous packets starting from next_expected_id
while (1) { while (1) {
struct ETCP_FRAGMENT* rx_pkt = (struct ETCP_FRAGMENT*)queue_find_data_by_id(etcp->recv_q, next_expected_id); struct ETCP_FRAGMENT* rx_pkt = (struct ETCP_FRAGMENT*)queue_find_data_by_index(etcp->recv_q, &next_expected_id, 4);
if (!rx_pkt) { if (!rx_pkt) {
// No more contiguous packets found // No more contiguous packets found
DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] no packet found for id=%u, stopping", etcp->log_name, next_expected_id); DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] no packet found for id=%u, stopping", etcp->log_name, next_expected_id);
@ -906,7 +906,7 @@ void etcp_output_try_assembly(struct ETCP_CONN* etcp) {
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "move: ETCP -> PN"); DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "move: ETCP -> PN");
// Add to output_queue using the same ETCP_FRAGMENT structure // Add to output_queue using the same ETCP_FRAGMENT structure
if (queue_data_put(etcp->output_queue, (struct ll_entry*)rx_pkt, next_expected_id) == 0) { if (queue_data_put(etcp->output_queue, (struct ll_entry*)rx_pkt) == 0) {
delivered_bytes += rx_pkt->ll.len; delivered_bytes += rx_pkt->ll.len;
delivered_count++; delivered_count++;
DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] moved packet id=%u to output_queue (qlen=%d)", etcp->log_name, DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "[%s] moved packet id=%u to output_queue (qlen=%d)", etcp->log_name,
@ -915,7 +915,7 @@ void etcp_output_try_assembly(struct ETCP_CONN* etcp) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to add packet id=%u to output_queue", etcp->log_name, DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "[%s] failed to add packet id=%u to output_queue", etcp->log_name,
next_expected_id); next_expected_id);
// Put it back in recv_q if we can't add to output_queue // 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); queue_data_put_with_index(etcp->recv_q, (struct ll_entry*)rx_pkt, 0, 4);
break; break;
} }
@ -937,13 +937,13 @@ void etcp_ack_recv(struct ETCP_CONN* etcp, uint32_t seq, uint16_t ts, uint16_t d
// Find the acknowledged packet in the wait_ack queue // Find the acknowledged packet in the wait_ack queue
struct INFLIGHT_PACKET* acked_pkt; struct INFLIGHT_PACKET* acked_pkt;
acked_pkt = (struct INFLIGHT_PACKET*)queue_find_data_by_id(etcp->input_wait_ack, seq); acked_pkt = (struct INFLIGHT_PACKET*)queue_find_data_by_index(etcp->input_wait_ack, &seq, 4);
if (acked_pkt) { if (acked_pkt) {
etcp->cnt_ack_hit_inf++; etcp->cnt_ack_hit_inf++;
queue_remove_data(etcp->input_wait_ack, (struct ll_entry*)acked_pkt); queue_remove_data(etcp->input_wait_ack, (struct ll_entry*)acked_pkt);
} }
else { else {
acked_pkt = (struct INFLIGHT_PACKET*)queue_find_data_by_id(etcp->input_send_q, seq); acked_pkt = (struct INFLIGHT_PACKET*)queue_find_data_by_index(etcp->input_send_q, &seq, 4);
if (acked_pkt) { if (acked_pkt) {
etcp->cnt_ack_hit_sndq++; etcp->cnt_ack_hit_sndq++;
queue_remove_data(etcp->input_send_q, (struct ll_entry*)acked_pkt); queue_remove_data(etcp->input_send_q, (struct ll_entry*)acked_pkt);
@ -1128,12 +1128,12 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) {
p->seq=seq; p->seq=seq;
p->pkt_timestamp=pkt->timestamp; p->pkt_timestamp=pkt->timestamp;
p->recv_timestamp=get_current_timestamp(); p->recv_timestamp=get_current_timestamp();
queue_data_put(etcp->ack_q, (struct ll_entry*)p, p->seq); queue_data_put_with_index(etcp->ack_q, (struct ll_entry*)p, 0, 4);
if (etcp->ack_resp_timer == NULL) { if (etcp->ack_resp_timer == NULL) {
etcp->ack_resp_timer = uasync_set_timeout(etcp->instance->ua, ACK_DELAY_TB, etcp, ack_response_timer_cb); 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); 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 if ((int32_t)(etcp->last_delivered_id-seq)<0) if (queue_find_data_by_index(etcp->recv_q, &seq, 4)==NULL) {// проверяем есть ли пакет с этим seq
uint32_t pkt_len=len-5; uint32_t pkt_len=len-5;
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); 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);
// отправляем пакет в очередь на сборку // отправляем пакет в очередь на сборку
@ -1158,7 +1158,7 @@ void etcp_conn_input(struct ETCP_DGRAM* pkt) {
rx_pkt->ll.memlen = etcp->instance->data_pool->object_size; rx_pkt->ll.memlen = etcp->instance->data_pool->object_size;
// Copy the actual payload data // Copy the actual payload data
memcpy(payload_data, data + 5, pkt_len); memcpy(payload_data, data + 5, pkt_len);
queue_data_put(etcp->recv_q, (struct ll_entry*)rx_pkt, seq); queue_data_put_with_index(etcp->recv_q, (struct ll_entry*)rx_pkt, 0, 4);
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); 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 ((int32_t)(seq - etcp->last_delivered_id) == 1) etcp_output_try_assembly(etcp);// пробуем собрать выходную очередь из фрагментов if ((int32_t)(seq - etcp->last_delivered_id) == 1) etcp_output_try_assembly(etcp);// пробуем собрать выходную очередь из фрагментов
} }

6
src/etcp.h

@ -42,9 +42,9 @@ uint16_t get_current_timestamp(void);
// пакет полностью удаляется когда приходит ACK (либо conn_reset/close) // пакет полностью удаляется когда приходит ACK (либо conn_reset/close)
struct INFLIGHT_PACKET {// выделяется из etcp->inflight_pool struct INFLIGHT_PACKET {// выделяется из etcp->inflight_pool
struct ll_entry ll; struct ll_entry ll;
uint32_t seq; // packet seq (ID по документации)
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
uint32_t seq; // packet seq (ID по документации)
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
@ -94,9 +94,9 @@ struct ETCP_CONN {
struct ll_queue* input_send_q; // очередь на отправку (inflight_pool -> INFLIGHT_PACKET) struct ll_queue* input_send_q; // очередь на отправку (inflight_pool -> INFLIGHT_PACKET)
struct ll_queue* input_wait_ack; // очередь ожидающих подтверждение (inflight_pool -> struct INFLIGHT_PACKET) struct ll_queue* input_wait_ack; // очередь ожидающих подтверждение (inflight_pool -> struct INFLIGHT_PACKET)
struct ll_queue* ack_q; // неотправленные подтверждения приема пакетов (instance.ack_pool -> struct ACK_PACKET) struct ll_queue* ack_q; // неотправленные подтверждения приема пакетов (instance.ack_pool -> struct ACK_PACKET) + index [+0 SEQ]
struct ll_queue* recv_q; // очередь на сборку фрагментированных пакетов(rx_pool -> struct ETCP_FRAGMENT) struct ll_queue* recv_q; // очередь на сборку фрагментированных пакетов(rx_pool -> struct ETCP_FRAGMENT) + index [+0 SEQ]
void (*link_ready_for_send_fn)(struct ETCP_CONN*);// функцию которую должен вызвать драйвер линка при готовности линка принимать данные void (*link_ready_for_send_fn)(struct ETCP_CONN*);// функцию которую должен вызвать драйвер линка при готовности линка принимать данные

2
src/etcp_api.c

@ -126,7 +126,7 @@ int etcp_send(struct ETCP_CONN* conn, struct ll_entry* entry) {
// Помещаем entry в очередь input normalizer // Помещаем entry в очередь input normalizer
// queue_data_put забирает ownership entry // queue_data_put забирает ownership entry
// DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, "Before put to input"); // DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, "Before put to input");
int result = queue_data_put(pn->input, entry, 0); int result = queue_data_put(pn->input, entry);
DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, "After put to input"); DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, "After put to input");
if (result != 0) { if (result != 0) {

27
src/etcp_connections.c

@ -1248,21 +1248,21 @@ ec_fr:
return; return;
} }
int init_connections(struct UTUN_INSTANCE* instance) { // Initialize only sockets (servers for incoming connections)
// Called before route_bgp_init() to populate etcp_sockets for nodeinfo
int init_sockets(struct UTUN_INSTANCE* instance) {
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, ""); DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "");
if (!instance || !instance->config) return -1; if (!instance || !instance->config) return -1;
if (instance->etcp_sockets) { if (instance->etcp_sockets) {
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Connections already initialized, skipping"); DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Sockets already initialized, skipping");
return 0; return 0;
} }
struct utun_config* config = instance->config; struct utun_config* config = instance->config;
// Initialize servers first - create sockets for incoming connections // Create sockets for servers (incoming connections)
struct CFG_SERVER* server = config->servers; struct CFG_SERVER* server = config->servers;
while (server) { while (server) {
// Create socket for this server
// Auto-detect local IP for public servers with 0.0.0.0 // Auto-detect local IP for public servers with 0.0.0.0
uint32_t default_ip = 0; uint32_t default_ip = 0;
if (server->type == CFG_SERVER_TYPE_PUBLIC) { if (server->type == CFG_SERVER_TYPE_PUBLIC) {
@ -1303,6 +1303,23 @@ int init_connections(struct UTUN_INSTANCE* instance) {
server = server->next; server = server->next;
} }
return 0;
}
int init_connections(struct UTUN_INSTANCE* instance) {
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "");
if (!instance || !instance->config) return -1;
struct utun_config* config = instance->config;
// If sockets already exist (created by init_sockets), skip server creation
if (!instance->etcp_sockets) {
DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Sockets not initialized, calling init_sockets()");
if (init_sockets(instance) < 0) {
return -1;
}
}
// Initialize clients - create outgoing connections // Initialize clients - create outgoing connections
struct CFG_CLIENT* client = config->clients; struct CFG_CLIENT* client = config->clients;
DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "init_connections called, instance=%p, config=%p, clients=%p, connections_count=%d", DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "init_connections called, instance=%p, config=%p, clients=%p, connections_count=%d",

5
src/etcp_connections.h

@ -159,7 +159,10 @@ struct ETCP_LINK {
uint16_t handshake_maxsize; // мax размер udp при handshake (выбирает рандом) uint16_t handshake_maxsize; // мax размер udp при handshake (выбирает рандом)
}; };
// INITIALIZATION (создаёт listen-сокеты и подключения из конфига) // INITIALIZATION
// Создаёт только listen-сокеты из конфига (серверы для incoming connections)
int init_sockets(struct UTUN_INSTANCE* instance);
// Создаёт listen-сокеты и client connections из конфига
int init_connections(struct UTUN_INSTANCE* instance); int init_connections(struct UTUN_INSTANCE* instance);
// SOCKET FUNCTIONS // SOCKET FUNCTIONS

6
src/pkt_normalizer.c

@ -204,7 +204,7 @@ void pn_packer_send(struct PKTNORM* pn, uint8_t* data, uint16_t len) {
} }
DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, "PUT to input"); DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, "PUT to input");
int ret = queue_data_put(pn->input, entry, 0); int ret = queue_data_put(pn->input, entry);
pn->in_total_pkts++; pn->in_total_pkts++;
pn->in_total_bytes += len; pn->in_total_bytes += len;
DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, "PUT to input end"); DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, "PUT to input end");
@ -246,7 +246,7 @@ static void pn_send_to_etcp(struct PKTNORM* pn) {
DEBUG_INFO(DEBUG_CATEGORY_NORMALIZER, "pn->etcp: size=%d memlen=%d frag_size=%d", frag->ll.len, frag->ll.memlen, pn->frag_size); DEBUG_INFO(DEBUG_CATEGORY_NORMALIZER, "pn->etcp: size=%d memlen=%d frag_size=%d", frag->ll.len, frag->ll.memlen, pn->frag_size);
if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP)) log_dump("NORM->ETCP", pn->data, frag->ll.len); if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP)) log_dump("NORM->ETCP", pn->data, frag->ll.len);
queue_data_put(pn->etcp->input_queue, (struct ll_entry*)frag, 0); queue_data_put(pn->etcp->input_queue, (struct ll_entry*)frag);
// Сбросить структуру (dgram передан во фрагмент, не освобождаем) // Сбросить структуру (dgram передан во фрагмент, не освобождаем)
pn->data = NULL; pn->data = NULL;
pn->data_ptr = 0; pn->data_ptr = 0;
@ -398,7 +398,7 @@ static void pn_unpacker_cb(struct ll_queue* q, void* arg) {
if (pn->recvpart->len == pn->recvpart->memlen) { if (pn->recvpart->len == pn->recvpart->memlen) {
DEBUG_DEBUG(DEBUG_CATEGORY_NORMALIZER, "unpacked dgram (size=%d)", pn->recvpart->len); DEBUG_DEBUG(DEBUG_CATEGORY_NORMALIZER, "unpacked dgram (size=%d)", pn->recvpart->len);
if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP)) log_dump("NORM->", pn->recvpart->dgram, pn->recvpart->len); if (debug_should_output(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DUMP)) log_dump("NORM->", pn->recvpart->dgram, pn->recvpart->len);
queue_data_put(pn->output, pn->recvpart, 0); queue_data_put(pn->output, pn->recvpart);
pn->out_total_pkts++; pn->out_total_pkts++;
pn->out_total_bytes += pn->recvpart->len; pn->out_total_bytes += pn->recvpart->len;
pn->recvpart = NULL; pn->recvpart = NULL;

111
src/route_bgp.c

@ -381,6 +381,102 @@ static void route_bgp_etcp_conn_cbk(struct ETCP_CONN* conn, void* arg) {
} }
} }
// Build MY_NODEINFO packet for this instance
static int route_bgp_build_my_nodeinfo(struct UTUN_INSTANCE* instance, struct ROUTE_BGP* bgp) {
if (!instance || !bgp) return -1;
// Count IPv4 sockets
int v4_socket_count = 0;
struct ETCP_SOCKET* sock = instance->etcp_sockets;
while (sock) {
if (sock->local_addr.ss_family == AF_INET) {
v4_socket_count++;
}
sock = sock->next;
}
// Count IPv4 routes (my_subnets)
int v4_route_count = 0;
struct CFG_ROUTE_ENTRY* subnet = instance->config->my_subnets;
while (subnet) {
if (subnet->ip.family == AF_INET) {
v4_route_count++;
}
subnet = subnet->next;
}
// Get name length
size_t name_len = strlen(instance->name);
// Calculate total size
size_t pkt_size = sizeof(struct BGP_NODEINFO_PACKET)
+ name_len
+ v4_socket_count * sizeof(struct BGP_NODEINFO_IPV4_SOCKET)
+ v4_route_count * sizeof(struct BGP_NODEINFO_IPV4_ROUTE);
// Allocate
struct BGP_NODEINFO_PACKET* pkt = u_malloc(pkt_size);
if (!pkt) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "route_bgp_build_my_nodeinfo: malloc failed");
return -1;
}
// Fill fixed fields
pkt->cmd = ETCP_ID_ROUTE_ENTRY;
pkt->subcmd = ROUTE_SUBCMD_NODEINFO;
pkt->node_id = htobe64(instance->node_id);
pkt->ver = 1;
memcpy(pkt->public_key, instance->my_keys.public_key, SC_PUBKEY_SIZE);
pkt->node_name_len = (uint8_t)name_len;
pkt->local_v4_sockets = (uint8_t)v4_socket_count;
pkt->local_v6_sockets = 0;
pkt->local_v4_routes = (uint8_t)v4_route_count;
pkt->local_v6_routes = 0;
// Fill variable data
uint8_t* ptr = (uint8_t*)pkt + sizeof(struct BGP_NODEINFO_PACKET);
// Name
if (name_len > 0) {
memcpy(ptr, instance->name, name_len);
ptr += name_len;
}
// IPv4 sockets
struct BGP_NODEINFO_IPV4_SOCKET* sock_arr = (struct BGP_NODEINFO_IPV4_SOCKET*)ptr;
sock = instance->etcp_sockets;
while (sock) {
if (sock->local_addr.ss_family == AF_INET) {
struct sockaddr_in* in = (struct sockaddr_in*)&sock->local_addr;
memcpy(sock_arr->addr, &in->sin_addr, 4);
sock_arr->port = in->sin_port;
sock_arr++;
}
sock = sock->next;
}
ptr = (uint8_t*)sock_arr;
// IPv4 routes
struct BGP_NODEINFO_IPV4_ROUTE* route_arr = (struct BGP_NODEINFO_IPV4_ROUTE*)ptr;
subnet = instance->config->my_subnets;
while (subnet) {
if (subnet->ip.family == AF_INET) {
memcpy(route_arr->addr, &subnet->ip.addr.v4, 4);
route_arr->prefix_length = subnet->netmask;
route_arr++;
}
subnet = subnet->next;
}
bgp->my_nodeinfo = pkt;
bgp->my_nodeinfo_size = pkt_size;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "Built my_nodeinfo: name='%s', v4_sockets=%d, v4_routes=%d",
instance->name, v4_socket_count, v4_route_count);
return 0;
}
struct ROUTE_BGP* route_bgp_init(struct UTUN_INSTANCE* instance) { struct ROUTE_BGP* route_bgp_init(struct UTUN_INSTANCE* instance) {
if (!instance) { if (!instance) {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "route_bgp_init: instance is NULL"); DEBUG_ERROR(DEBUG_CATEGORY_BGP, "route_bgp_init: instance is NULL");
@ -398,6 +494,14 @@ struct ROUTE_BGP* route_bgp_init(struct UTUN_INSTANCE* instance) {
bgp->instance = instance; bgp->instance = instance;
bgp->table_version = 1; bgp->table_version = 1;
// Build my_nodeinfo packet
if (route_bgp_build_my_nodeinfo(instance, bgp) != 0) {
queue_free(bgp->senders_list);
u_free(bgp);
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "route_bgp_init: build_my_nodeinfo failed");
return NULL;
}
bgp->senders_list = queue_new(instance->ua, 0, "BGP_senders"); bgp->senders_list = queue_new(instance->ua, 0, "BGP_senders");
if (!bgp->senders_list) { if (!bgp->senders_list) {
u_free(bgp); u_free(bgp);
@ -445,6 +549,11 @@ void route_bgp_destroy(struct UTUN_INSTANCE* instance) {
} }
queue_free(instance->bgp->senders_list); queue_free(instance->bgp->senders_list);
// Free my_nodeinfo
if (instance->bgp->my_nodeinfo) {
u_free(instance->bgp->my_nodeinfo);
}
u_free(instance->bgp); u_free(instance->bgp);
instance->bgp = NULL; instance->bgp = NULL;
} }
@ -488,7 +597,7 @@ void route_bgp_new_conn(struct ETCP_CONN* conn) {
if (!item_entry) return; if (!item_entry) return;
((struct ROUTE_BGP_CONN_ITEM*)item_entry->data)->conn = conn; ((struct ROUTE_BGP_CONN_ITEM*)item_entry->data)->conn = conn;
queue_data_put(bgp->senders_list, item_entry, 0); queue_data_put(bgp->senders_list, item_entry);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "New connection added to senders_list"); DEBUG_INFO(DEBUG_CATEGORY_BGP, "New connection added to senders_list");
} }

48
src/route_bgp.h

@ -2,8 +2,10 @@
#define ROUTE_BGP_H #define ROUTE_BGP_H
#include <stdint.h> #include <stdint.h>
#include <stddef.h>
#include "../lib/ll_queue.h" #include "../lib/ll_queue.h"
#include "route_lib.h" #include "route_lib.h"
#include "secure_channel.h"
// ETCP ID для маршрутных пакетов // ETCP ID для маршрутных пакетов
#define ETCP_ID_ROUTE_ENTRY 0x01 #define ETCP_ID_ROUTE_ENTRY 0x01
@ -31,6 +33,50 @@ struct BGP_ROUTE_PACKET {
uint64_t hop_list[0]; // flexible array: next_hop → ... → destination uint64_t hop_list[0]; // flexible array: next_hop → ... → destination
} __attribute__((packed)); } __attribute__((packed));
/**
* @brief Основной пакет маршрута (переменной длины)
*/
struct BGP_NODEINFO_PACKET {
uint8_t cmd; // ETCP_ID_ROUTE_ENTRY
uint8_t subcmd; // ROUTE_SUBCMD_NODEINFO
uint64_t node_id; // (big-endian)
uint8_t ver; // версия пакета (сейчас 1)
uint8_t public_key[SC_PUBKEY_SIZE]; // node pubkey
uint8_t node_name_len; // размер в байтах (без null терминации)
uint8_t local_v4_sockets; // BGP_NODEINFO_IPV4_SOCKET число локальных сокетов ipv4 (для incoming connections)
uint8_t local_v6_sockets; // BGP_NODEINFO_IPV6_SOCKET число локальных сокетов ipv6 (для incoming connections) (пока 0)
uint8_t local_v4_routes; // BGP_NODEINFO_IPV4_ROUTE число локальных маршрутов ipv4
uint8_t local_v6_routes; // BGP_NODEINFO_IPV6_ROUTE число локальных маршрутов ipv6 (пока 0)
uint8_t tranzit_nodes; // BGP_NODEINFO_TRANZIT_NODE лучшие транзитные узлы для этой ноды (минимальный пинг / лучшее качество каналов)
// далее идут по порядку следования полей в структуре: node_name, сокеты, роуты
} __attribute__((packed));
struct BGP_NODEINFO_IPV4_SOCKET {
uint8_t addr[4];// network byte order
uint16_t port;
} __attribute__((packed));
struct BGP_NODEINFO_IPV6_SOCKET {
uint8_t addr[16];
uint16_t port;
} __attribute__((packed));
struct BGP_NODEINFO_IPV4_ROUTE {
uint8_t addr[4];// network byte order
uint8_t prefix_length;
} __attribute__((packed));
struct BGP_NODEINFO_IPV6_ROUTE {
uint8_t addr[16];
uint8_t prefix_length;
} __attribute__((packed));
struct BGP_NODEINFO_TRANZIT_NODE {
uint64_t node_id; // (big-endian)
uint16_t rtt; // x0.1 ms
uint16_t link_q; // меньше - лучше (потери + 1/BW)
} __attribute__((packed));
/** /**
* @brief Пакет WITHDRAW (фиксированный) * @brief Пакет WITHDRAW (фиксированный)
*/ */
@ -56,6 +102,8 @@ struct ROUTE_BGP_CONN_ITEM {
struct ROUTE_BGP { struct ROUTE_BGP {
struct UTUN_INSTANCE* instance; struct UTUN_INSTANCE* instance;
struct BGP_NODEINFO_PACKET* my_nodeinfo;
uint16_t my_nodeinfo_size;
struct ll_queue* senders_list; struct ll_queue* senders_list;
uint32_t table_version; // starts at 1, for future use uint32_t table_version; // starts at 1, for future use
}; };

28
src/route_bgp.txt

@ -39,5 +39,33 @@ subcmd:
удаленный узел считает свои хеши и сравнивает. где не совпало смотрит сколько записей. удаленный узел считает свои хеши и сравнивает. где не совпало смотрит сколько записей.
если записей не много - передает эти записи. если записей не много - передает эти записи.
если записей много - добавляет 4 бита к хеш таблице и строит субтаблицу для если записей много - добавляет 4 бита к хеш таблице и строит субтаблицу для
==========================================
Формат роутинга:
Таблица узлов состоит из записей:
- uid
- name
- links
- 3 транзитных узла с метриками (RTT)
- маршруты узла
- текущая загрузка линка (за последние 10 сек) можно частоту адаптировать под размер сети
- bandidth limit
- transit bandwidth limit
две группы узлов:
- узлы за nat. подключаются через транзитные узлы. измеряют пинги до транзитных и выбирают N (3 default) лучшие линки. 3 лучших используем для распространения маршрутов
- транзитные узлы. имеют линки с загрузкой.
добавить кодограмму - отменить распространение маршрутов по линку (+ сделать важным линком)
карта маршрутизации:
План:
- сделать передачу роутинга в
- сделать фоновый probe для узлов (условно 1 нода в секунду). выигравшие по качеству соатновятся основными
- сделать etcp дизконнект:
- отправить disconnect request + дождаться ack дальше master удаляет, slave удаляет по down.

2
src/routing.c

@ -178,7 +178,7 @@ static void route_pkt(struct UTUN_INSTANCE* instance, struct ll_entry* entry, ui
if (route->conn_list == NULL) { if (route->conn_list == NULL) {
// Local route - send to TUN (entry has [cmd=0][IP data], TUN skips cmd byte) // Local route - send to TUN (entry has [cmd=0][IP data], TUN skips cmd byte)
int put_err = queue_data_put(instance->tun->input_queue, entry, 0); int put_err = queue_data_put(instance->tun->input_queue, entry);
if (put_err != 0) { if (put_err != 0) {
DEBUG_WARN(DEBUG_CATEGORY_ROUTING, "route_pkt: failed to put to TUN: dst=%s err=%d", DEBUG_WARN(DEBUG_CATEGORY_ROUTING, "route_pkt: failed to put to TUN: dst=%s err=%d",
ip_to_str(&addr, AF_INET).str, put_err); ip_to_str(&addr, AF_INET).str, put_err);

8
src/tun_if.c

@ -59,7 +59,7 @@ static void tun_read_callback(int fd, void* user_arg)
pkt->len = nread + 1; pkt->len = nread + 1;
// Add to output queue (TUN → routing) // Add to output queue (TUN → routing)
if (queue_data_put(tun->output_queue, pkt, 0) != 0) { if (queue_data_put(tun->output_queue, pkt) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_TUN, "Failed to add packet to output queue"); DEBUG_ERROR(DEBUG_CATEGORY_TUN, "Failed to add packet to output queue");
u_free(packet_data); u_free(packet_data);
queue_entry_free(pkt); queue_entry_free(pkt);
@ -293,7 +293,7 @@ void tun_packet_handler(void* arg) {
} }
// Положить в очередь (уже в main thread) // Положить в очередь (уже в main thread)
int ok = queue_data_put(tun->output_queue, entry, 0); int ok = queue_data_put(tun->output_queue, entry);
if (ok != 0) { if (ok != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_TUN, "Put error"); DEBUG_ERROR(DEBUG_CATEGORY_TUN, "Put error");
queue_entry_free(entry); queue_entry_free(entry);
@ -318,7 +318,7 @@ int tun_inject_packet(struct tun_if* tun, const uint8_t* buf, size_t len)
pkt->dgram = data; pkt->dgram = data;
pkt->len = len + 1; pkt->len = len + 1;
int ret = queue_data_put(tun->output_queue, pkt, 0); int ret = queue_data_put(tun->output_queue, pkt);
if (ret != 0) { if (ret != 0) {
u_free(data); u_free(data);
@ -340,7 +340,7 @@ ssize_t tun_read_packet(struct tun_if* tun, uint8_t* buf, size_t len)
if (!pkt) return 0; if (!pkt) return 0;
if (pkt->len > len) { if (pkt->len > len) {
queue_data_put(tun->input_queue, pkt, 0); queue_data_put(tun->input_queue, pkt);
return -1; return -1;
} }

23
src/utun_instance.c

@ -45,6 +45,14 @@ static int instance_init_common(struct UTUN_INSTANCE* instance, struct UASYNC* u
instance->ua = ua; instance->ua = ua;
instance->config = config; instance->config = config;
// Set name from config
if (config->global.name[0] != '\0') {
strncpy(instance->name, config->global.name, sizeof(instance->name) - 1);
instance->name[sizeof(instance->name) - 1] = '\0';
} else {
instance->name[0] = '\0';
}
// Set node_id from config // Set node_id from config
instance->node_id = config->global.my_node_id; instance->node_id = config->global.my_node_id;
@ -108,6 +116,12 @@ static int instance_init_common(struct UTUN_INSTANCE* instance, struct UASYNC* u
instance->tun = NULL; instance->tun = NULL;
} }
// Initialize sockets first (needed for BGP nodeinfo)
if (init_sockets(instance) < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "Failed to initialize sockets");
return -1;
}
// Initialize BGP module for route exchange // Initialize BGP module for route exchange
instance->bgp = route_bgp_init(instance); instance->bgp = route_bgp_init(instance);
if (!instance->bgp) { if (!instance->bgp) {
@ -177,6 +191,15 @@ struct UTUN_INSTANCE* utun_instance_create(struct UASYNC* ua, const char *config
return NULL; return NULL;
} }
// Log instance info
if (instance->name[0] != '\0') {
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "uTun instance '%s' created, node_id=0x%llx",
instance->name, (unsigned long long)instance->node_id);
} else {
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "uTun instance created, node_id=0x%llx",
(unsigned long long)instance->node_id);
}
return instance; return instance;
} }

3
src/utun_instance.h

@ -24,6 +24,9 @@ struct control_server;
// uTun instance configuration // uTun instance configuration
struct UTUN_INSTANCE { struct UTUN_INSTANCE {
// Identification
char name[16]; // Instance name from config
// Configuration (moved from utun_state) // Configuration (moved from utun_state)
struct utun_config *config; struct utun_config *config;

2
tests/debug_simple.c

@ -26,7 +26,7 @@ int main() {
printf("Data pointer: %p\n", data1); printf("Data pointer: %p\n", data1);
/* Put data into queue */ /* Put data into queue */
int put_result = queue_data_put(q, data1, data1->id); int put_result = queue_data_put(q, data1);
printf("queue_data_put returned: %d\n", put_result); printf("queue_data_put returned: %d\n", put_result);
printf("After put: queue count = %d\n", queue_entry_count(q)); printf("After put: queue count = %d\n", queue_entry_count(q));

4
tests/test_intensive_memory_pool.c

@ -39,7 +39,7 @@ static double test_without_pools(int iterations) {
// Добавить записи // Добавить записи
for (int i = 0; i < 10; i++) { for (int i = 0; i < 10; i++) {
void* data = queue_entry_new(64); void* data = queue_entry_new(64);
queue_data_put(queue, data, i); // Используем ID = i queue_data_put(queue, data); // Используем ID = i
} }
// Удалить записи (триггер waiters) // Удалить записи (триггер waiters)
@ -91,7 +91,7 @@ static double test_with_pools(int iterations) {
for (int i = 0; i < 10; i++) { for (int i = 0; i < 10; i++) {
void* data = queue_entry_new_from_pool(pool); void* data = queue_entry_new_from_pool(pool);
if (data) { if (data) {
queue_data_put(queue, data, i); // Используем ID = i queue_data_put(queue, data); // Используем ID = i
} }
} }

4
tests/test_intensive_memory_pool_new.c

@ -39,7 +39,7 @@ static double test_without_pools(int iterations) {
// Добавить записи // Добавить записи
for (int i = 0; i < 10; i++) { for (int i = 0; i < 10; i++) {
void* data = queue_entry_new(64); void* data = queue_entry_new(64);
queue_data_put(queue, data, i); // Используем ID = i queue_data_put(queue, data); // Используем ID = i
} }
// Удалить записи (триггер waiters) // Удалить записи (триггер waiters)
@ -91,7 +91,7 @@ static double test_with_pools(int iterations) {
for (int i = 0; i < 10; i++) { for (int i = 0; i < 10; i++) {
void* data = queue_entry_new_from_pool(pool); void* data = queue_entry_new_from_pool(pool);
if (data) { if (data) {
queue_data_put(queue, data, i); // Используем ID = i queue_data_put(queue, data); // Используем ID = i
} }
} }

27
tests/test_ll_queue.c

@ -95,7 +95,7 @@ static void test_fifo(void) {
test_data_t *d = (test_data_t*)queue_entry_new(sizeof(test_data_t)); test_data_t *d = (test_data_t*)queue_entry_new(sizeof(test_data_t));
d->id = i; snprintf(d->name, sizeof(d->name), "item%d", i); d->id = i; snprintf(d->name, sizeof(d->name), "item%d", i);
d->value = i*10; d->checksum = checksum(d); d->value = i*10; d->checksum = checksum(d);
queue_data_put(q, (struct ll_entry*)d, d->id); queue_data_put_with_index(q, (struct ll_entry*)d, 0, 4);
} }
ASSERT_EQ(queue_entry_count(q), 10, ""); ASSERT_EQ(queue_entry_count(q), 10, "");
@ -119,12 +119,12 @@ static void test_lifo_priority(void) {
d->id = i; d->id = i;
d->value = i; d->value = i;
d->checksum = checksum(d); d->checksum = checksum(d);
queue_data_put(q, (struct ll_entry*)d, i); queue_data_put_with_index(q, (struct ll_entry*)d, 0, 4);
} }
test_data_t *pri = (test_data_t*)queue_entry_new(sizeof(*pri)); test_data_t *pri = (test_data_t*)queue_entry_new(sizeof(*pri));
pri->id = 999; pri->value = 999; pri->checksum = checksum(pri); pri->id = 999; pri->value = 999; pri->checksum = checksum(pri);
queue_data_put_first(q, (struct ll_entry*)pri, 999); queue_data_put_first_with_index(q, (struct ll_entry*)pri, 0, 4);
test_data_t *first = (test_data_t*)queue_data_get(q); test_data_t *first = (test_data_t*)queue_data_get(q);
ASSERT(first && first->id == 999, "priority first"); ASSERT(first && first->id == 999, "priority first");
@ -150,7 +150,7 @@ static void test_callback(void) {
for (int i = 0; i < 5; i++) { for (int i = 0; i < 5; i++) {
test_data_t *d = (test_data_t*)queue_entry_new(sizeof(*d)); test_data_t *d = (test_data_t*)queue_entry_new(sizeof(*d));
d->id = i; d->value = i*10; d->checksum = checksum(d); d->id = i; d->value = i*10; d->checksum = checksum(d);
queue_data_put(q, (struct ll_entry*)d, d->id); queue_data_put_with_index(q, (struct ll_entry*)d, 0, 4);
} }
for (int i = 0; i < 30 && cnt < 5; i++) uasync_poll(ua, 5); for (int i = 0; i < 30 && cnt < 5; i++) uasync_poll(ua, 5);
@ -172,7 +172,7 @@ static void test_waiter(void) {
for (int i = 0; i < 5; i++) { for (int i = 0; i < 5; i++) {
test_data_t *d = (test_data_t*)queue_entry_new(sizeof(*d)); d->id = i; test_data_t *d = (test_data_t*)queue_entry_new(sizeof(*d)); d->id = i;
queue_data_put(q, (struct ll_entry*)d, i); queue_data_put_with_index(q, (struct ll_entry*)d, 0, 4);
} }
called = 0; called = 0;
@ -197,15 +197,16 @@ static void test_limits_hash(void) {
for (int i = 0; i < 3; i++) { for (int i = 0; i < 3; i++) {
test_data_t *d = (test_data_t*)queue_entry_new(sizeof(*d)); d->id = i*10+1; test_data_t *d = (test_data_t*)queue_entry_new(sizeof(*d)); d->id = i*10+1;
queue_data_put(q, (struct ll_entry*)d, d->id); queue_data_put_with_index(q, (struct ll_entry*)d, 0, 4);
} }
test_data_t *ex = (test_data_t*)queue_entry_new(sizeof(*ex)); test_data_t *ex = (test_data_t*)queue_entry_new(sizeof(*ex));
ASSERT_EQ(queue_data_put(q, (struct ll_entry*)ex, 999), -1, "limit reject"); ASSERT_EQ(queue_data_put_with_index(q, (struct ll_entry*)ex, 0, 4), -1, "limit reject");
test_data_t *found = (test_data_t*)queue_find_data_by_id(q, 21); uint32_t hash=21;
test_data_t *found = (test_data_t*)queue_find_data_by_index(q, &hash, 4);
ASSERT(found && found->id == 21, "hash find"); ASSERT(found && found->id == 21, "hash find");
queue_remove_data(q, (struct ll_entry*)found); queue_remove_data(q, (struct ll_entry*)found);
ASSERT(queue_find_data_by_id(q, 21) == NULL, "removed"); ASSERT(queue_find_data_by_index(q, &hash, 4) == NULL, "removed");
while (queue_entry_count(q)) queue_entry_free((struct ll_entry*)queue_data_get(q)); while (queue_entry_count(q)) queue_entry_free((struct ll_entry*)queue_data_get(q));
queue_free(q); uasync_destroy(ua, 0); queue_free(q); uasync_destroy(ua, 0);
@ -223,11 +224,11 @@ static void test_pool(void) {
test_data_t *d1 = (test_data_t*)queue_entry_new_from_pool(pool); test_data_t *d1 = (test_data_t*)queue_entry_new_from_pool(pool);
d1->id = 1; d1->checksum = checksum(d1); d1->id = 1; d1->checksum = checksum(d1);
queue_data_put(q, (struct ll_entry*)d1, 1); queue_data_put_with_index(q, (struct ll_entry*)d1, 0, 4);
test_data_t *d2 = (test_data_t*)queue_entry_new_from_pool(pool); test_data_t *d2 = (test_data_t*)queue_entry_new_from_pool(pool);
d2->id = 2; d2->checksum = checksum(d2); d2->id = 2; d2->checksum = checksum(d2);
queue_data_put(q, (struct ll_entry*)d2, 2); queue_data_put_with_index(q, (struct ll_entry*)d2, 0, 4);
queue_entry_free((struct ll_entry*)queue_data_get(q)); // free d1 back to pool queue_entry_free((struct ll_entry*)queue_data_get(q)); // free d1 back to pool
queue_entry_free((struct ll_entry*)queue_data_get(q)); // free d2 back to pool queue_entry_free((struct ll_entry*)queue_data_get(q)); // free d2 back to pool
@ -236,7 +237,7 @@ static void test_pool(void) {
test_data_t *d3 = (test_data_t*)queue_entry_new_from_pool(pool); test_data_t *d3 = (test_data_t*)queue_entry_new_from_pool(pool);
ASSERT(d3 != NULL, "alloc after free failed"); ASSERT(d3 != NULL, "alloc after free failed");
d3->id = 3; d3->checksum = checksum(d3); d3->id = 3; d3->checksum = checksum(d3);
queue_data_put(q, (struct ll_entry*)d3, 3); queue_data_put_with_index(q, (struct ll_entry*)d3, 0, 4);
queue_entry_free((struct ll_entry*)queue_data_get(q)); // free d3 back queue_entry_free((struct ll_entry*)queue_data_get(q)); // free d3 back
size_t alloc2 = 0, reuse2 = 0; size_t alloc2 = 0, reuse2 = 0;
@ -260,7 +261,7 @@ static void test_stress(void) {
test_data_t *d = (test_data_t*)queue_entry_new(sizeof(*d)); test_data_t *d = (test_data_t*)queue_entry_new(sizeof(*d));
d->id = rand() % 10000; d->id = rand() % 10000;
d->checksum = checksum(d); d->checksum = checksum(d);
queue_data_put(q, (struct ll_entry*)d, d->id); queue_data_put_with_index(q, (struct ll_entry*)d, 0, 4);
stats.ops++; stats.ops++;
} }
while (queue_entry_count(q)) queue_entry_free((struct ll_entry*)queue_data_get(q)); while (queue_entry_count(q)) queue_entry_free((struct ll_entry*)queue_data_get(q));

2
tests/test_memory_pool_and_config.c

@ -59,7 +59,7 @@ int main() {
// Add some entries and trigger waiters // Add some entries and trigger waiters
for (int i = 0; i < 5; i++) { for (int i = 0; i < 5; i++) {
void* data = queue_entry_new(10); void* data = queue_entry_new(10);
queue_data_put(queue, data, i); // Используем ID = i queue_data_put(queue, data); // Используем ID = i
} }
// Remove entries to trigger waiter callbacks // Remove entries to trigger waiter callbacks

2
tests/test_pkt_normalizer_standalone.c

@ -82,7 +82,7 @@ static void loopback_callback(struct ll_queue* q, void* arg) {
while ((frag = (struct ETCP_FRAGMENT*)queue_data_get(q)) != NULL) { while ((frag = (struct ETCP_FRAGMENT*)queue_data_get(q)) != NULL) {
fragments_sent++; fragments_sent++;
// Move fragment from input to output (loopback) // Move fragment from input to output (loopback)
queue_data_put(mock_etcp.output_queue, (struct ll_entry*)frag, 0); queue_data_put(mock_etcp.output_queue, (struct ll_entry*)frag);
} }
queue_resume_callback(q); queue_resume_callback(q);

Loading…
Cancel
Save