diff --git a/lib/ll_queue.c b/lib/ll_queue.c index a727d524..b30a9e43 100644 --- a/lib/ll_queue.c +++ b/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) { - if (!q || q->hash_size == 0 || !entry) return; - - size_t slot = entry->id % q->hash_size; + 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; + } + + uint32_t slot = entry->index_hash % q->hash_size; entry->hash_next = q->hash_table[slot]; q->hash_table[slot] = entry; } -// Внутренняя функция удаления из хеш-таблицы static void remove_from_hash(struct ll_queue* q, struct ll_entry* entry) { 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]; while (*ptr) { 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) { 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; - #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 = 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 -// 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 + + 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); - entry->id = id; - - // Проверить лимит размера if (q->size_limit >= 0 && q->count >= q->size_limit) { queue_dgram_free(entry); - queue_entry_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; - } + if (q->tail) q->tail->next = entry; else q->head = entry; q->tail = entry; - + q->count++; - entry->int_len=entry->len; + entry->int_len = entry->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); - -// 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 - queue_check_consistency(q);// !!!! for debug - BEFORE callback + queue_check_consistency(q); +#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) { q->callback(q, q->callback_arg); } - - // Проверить ожидающие коллбэки (надо только при заборе из очереди) -// check_waiters(q); - + #ifdef QUEUE_DEBUG - queue_check_consistency(q);// !!!! for debug - AFTER callback + queue_check_consistency(q); #endif 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; - #ifdef QUEUE_THREAD_CHECK queue_check_thread(q); -#endif +#endif + + 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); - entry->id = id; - - // Проверить лимит размера if (q->size_limit >= 0 && q->count >= q->size_limit) { - queue_entry_free(entry); // Освободить элемент если превышен лимит + 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; - } + if (q->head) q->head->prev = entry; else q->tail = entry; q->head = entry; - + q->count++; - entry->int_len=entry->len; + entry->int_len = entry->len; q->total_bytes += entry->int_len; - + 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) { q->callback(q, q->callback_arg); } - - // Проверить ожидающие коллбэки -// check_waiters(q); - + #ifdef QUEUE_DEBUG - queue_check_consistency(q);// !!!! for debug + queue_check_consistency(q); #endif return 0; } + struct ll_entry* queue_data_get(struct ll_queue* q) { if (!q || !q->head) return NULL; @@ -496,19 +581,22 @@ void queue_cancel_wait(struct ll_queue* q, struct queue_waiter* waiter) { // ==================== Поиск и удаление по ID ==================== -struct ll_entry* queue_find_data_by_id(struct ll_queue* q, uint32_t id) { - if (!q || q->hash_size == 0 || !q->hash_table) return NULL; - - size_t slot = id % q->hash_size; +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 || index_size == 0 || !index_key) return NULL; + + uint32_t hash_val = make_hash(index_key, index_size); + + uint32_t slot = hash_val % q->hash_size; struct ll_entry* entry = q->hash_table[slot]; - + 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; } entry = entry->hash_next; } - return NULL; } diff --git a/lib/ll_queue.h b/lib/ll_queue.h index 9e9869a1..c9c93039 100644 --- a/lib/ll_queue.h +++ b/lib/ll_queue.h @@ -54,19 +54,21 @@ struct ll_queue; */ struct ll_entry { char* name; - struct ll_entry* next; ///< Следующий элемент в очереди - struct ll_entry* prev; ///< Предыдущий элемент в очереди - uint16_t size; ///< Размер пользовательского буфера data[] (байт) - uint16_t len; ///< Актуальная длина данных в dgram - uint16_t memlen; ///< Выделенный размер под dgram - uint16_t int_len; ///< Внутреннее (не использовать) - uint8_t* dgram; ///< Указатель на данные пакета - void (*dgram_free_fn)(uint8_t* data); ///< Кастомная функция освобождения dgram - struct memory_pool* dgram_pool; ///< Пул для dgram (если используется) - uint32_t id; ///< Идентификатор для поиска (задаётся при добавлении) - struct ll_entry* hash_next; ///< Следующий в хеш-цепочке - struct memory_pool* pool; ///< Пул, из которого выделен сам entry (NULL = malloc) - uint8_t data[0]; ///< Гибкий массив пользовательских данных (размер = size) + struct ll_entry* next; // Следующий элемент в очереди + struct ll_entry* prev; // Предыдущий элемент в очереди + uint16_t size; // Размер пользовательского буфера data[] (байт) + uint16_t len; // Актуальная длина данных в dgram + uint16_t memlen; // Выделенный размер под dgram + uint16_t int_len; // Внутреннее (не использовать) + uint8_t* dgram; // Указатель на данные пакета + void (*dgram_free_fn)(uint8_t* data); // Кастомная функция освобождения dgram + struct memory_pool* dgram_pool; // Пул для dgram (если используется) + uint16_t index_offset; // Смещение индекса в data[] + uint16_t index_size; // Длина индекса (0 = без индекса) + uint32_t index_hash; // Хеш по первым 4 байтам индекса + 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 { char* name; - struct ll_entry* head; ///< Голова очереди (отсюда извлекаем) - struct ll_entry* tail; ///< Хвост очереди (сюда добавляем) - int count; ///< Текущее количество элементов - size_t total_bytes; ///< Суммарный объём данных (сумма int_len) - int size_limit; ///< Максимальное количество элементов (-1 = без лимита) + struct ll_entry* head; // Голова очереди (отсюда извлекаем) + struct ll_entry* tail; // Хвост очереди (сюда добавляем) + int count; // Текущее количество элементов + size_t total_bytes; // Суммарный объём данных (сумма int_len) + int size_limit; // Максимальное количество элементов (-1 = без лимита) - queue_callback_fn callback; ///< Коллбэк автозабора - void* callback_arg; ///< Аргумент коллбэка - int callback_suspended; ///< 1 = коллбэки временно приостановлены + queue_callback_fn callback; // Коллбэк автозабора + void* callback_arg; // Аргумент коллбэка + int callback_suspended; // 1 = коллбэки временно приостановлены - void* resume_timeout_id; ///< ID таймера uasync для отложенного resume - struct UASYNC* ua; ///< Экземпляр uasync (обязателен для таймеров) + void* resume_timeout_id; // ID таймера uasync для отложенного resume + struct UASYNC* ua; // Экземпляр uasync (обязателен для таймеров) - struct queue_waiter waiter; ///< Встроенный waiter (только один) + struct queue_waiter waiter; // Встроенный waiter (только один) - struct ll_entry** hash_table; ///< Хеш-таблица для поиска по id (если hash_size > 0) - size_t hash_size; ///< Размер хеш-таблицы + struct ll_entry** hash_table; // Хеш-таблица для поиска по id (если hash_size > 0) + size_t hash_size; // Размер хеш-таблицы #ifdef QUEUE_THREAD_CHECK #ifdef _WIN32 @@ -215,7 +217,7 @@ void queue_cancel_wait(struct ll_queue* q, struct queue_waiter* waiter); * @param id идентификатор для поиска * @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, высокий приоритет). @@ -224,7 +226,12 @@ int queue_data_put(struct ll_queue* q, struct ll_entry* entry, uint32_t id); * @param id идентификатор * @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 Извлекает элемент из начала очереди. @@ -291,7 +298,9 @@ void queue_dgram_free(struct ll_entry* entry); * @param id искомый идентификатор * @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 Удаляет элемент из очереди (не освобождает память). diff --git a/src/config_parser.c b/src/config_parser.c index 52aef1c7..d852250b 100644 --- a/src/config_parser.c +++ b/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) { + if (strcmp(key, "my_node_name") == 0) { + return assign_string(global->name, sizeof(global->name), value); + } if (strcmp(key, "my_private_key") == 0) { 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 + cfg->global.name[0] = '\0'; cfg->global.keepalive_timeout = 2000; // Default 2 seconds cfg->global.firewall_rules = NULL; cfg->global.firewall_rule_count = 0; diff --git a/src/config_parser.h b/src/config_parser.h index 2024879c..68c79607 100644 --- a/src/config_parser.h +++ b/src/config_parser.h @@ -67,6 +67,7 @@ struct CFG_FIREWALL_RULE { }; struct global_config { + char name[16]; // Instance name char my_private_key_hex[MAX_KEY_LEN]; char my_public_key_hex[MAX_KEY_LEN]; uint64_t my_node_id; diff --git a/src/etcp.c b/src/etcp.c index f0877224..a26941aa 100644 --- a/src/etcp.c +++ b/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); // 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); queue_dgram_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) - memset(p, 0, sizeof(*p)); +// memset(p, 0, sizeof(*p)); p->seq = etcp->next_tx_id++; // Assign seq p->state = INFLIGHT_STATE_WAIT_SEND; 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 // 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); queue_dgram_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 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 etcp->retransmissions_count++; @@ -767,7 +767,7 @@ struct ETCP_DGRAM* etcp_request_pkt(struct ETCP_CONN* etcp) { inf_pkt->send_count++; 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); @@ -889,7 +889,7 @@ void etcp_output_try_assembly(struct ETCP_CONN* etcp) { // Look for contiguous packets starting from next_expected_id 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) { // No more contiguous packets found 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"); // 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_count++; 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, next_expected_id); // Put it back in recv_q if we can't add to output_queue - queue_data_put(etcp->recv_q, (struct ll_entry*)rx_pkt, next_expected_id); + queue_data_put_with_index(etcp->recv_q, (struct ll_entry*)rx_pkt, 0, 4); 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 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) { etcp->cnt_ack_hit_inf++; queue_remove_data(etcp->input_wait_ack, (struct ll_entry*)acked_pkt); } 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) { etcp->cnt_ack_hit_sndq++; 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->pkt_timestamp=pkt->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) { etcp->ack_resp_timer = uasync_set_timeout(etcp->instance->ua, ACK_DELAY_TB, etcp, ack_response_timer_cb); DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] set ack_timer for delayed ACK send", etcp->log_name); } - if ((int32_t)(etcp->last_delivered_id-seq)<0) if (queue_find_data_by_id(etcp->recv_q, seq)==NULL) {// проверяем есть ли пакет с этим seq + 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; 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; // Copy the actual payload data 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); if ((int32_t)(seq - etcp->last_delivered_id) == 1) etcp_output_try_assembly(etcp);// пробуем собрать выходную очередь из фрагментов } diff --git a/src/etcp.h b/src/etcp.h index f298f28f..0c4085cf 100644 --- a/src/etcp.h +++ b/src/etcp.h @@ -42,9 +42,9 @@ uint16_t get_current_timestamp(void); // пакет полностью удаляется когда приходит ACK (либо conn_reset/close) struct INFLIGHT_PACKET {// выделяется из etcp->inflight_pool struct ll_entry ll; + uint32_t seq; // packet seq (ID по документации) struct ETCP_LINK* last_link; // Last sent link uint64_t last_timestamp; // Last send timestamp - uint32_t seq; // packet seq (ID по документации) uint8_t send_count; // Number of sends uint8_t retrans_req_count; // Number of retrans requests 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_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*);// функцию которую должен вызвать драйвер линка при готовности линка принимать данные diff --git a/src/etcp_api.c b/src/etcp_api.c index 8fad3a08..405ace8b 100644 --- a/src/etcp_api.c +++ b/src/etcp_api.c @@ -126,7 +126,7 @@ int etcp_send(struct ETCP_CONN* conn, struct ll_entry* entry) { // Помещаем entry в очередь input normalizer // queue_data_put забирает ownership entry // 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"); if (result != 0) { diff --git a/src/etcp_connections.c b/src/etcp_connections.c index b527f5b3..0efd2bf1 100644 --- a/src/etcp_connections.c +++ b/src/etcp_connections.c @@ -1248,21 +1248,21 @@ ec_fr: 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, ""); if (!instance || !instance->config) return -1; if (instance->etcp_sockets) { - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Connections already initialized, skipping"); + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "Sockets already initialized, skipping"); return 0; } - + 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; while (server) { - // Create socket for this server - // Auto-detect local IP for public servers with 0.0.0.0 uint32_t default_ip = 0; if (server->type == CFG_SERVER_TYPE_PUBLIC) { @@ -1298,11 +1298,28 @@ int init_connections(struct UTUN_INSTANCE* instance) { snprintf(addr_str, sizeof(addr_str), "%s:%d", ip_to_str(&sin6->sin6_addr, AF_INET6).str, ntohs(sin6->sin6_port)); } - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Initialized server %s on %s (links: %zu)", + DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "Initialized server %s on %s (links: %zu)", server->name, addr_str, e_sock->num_channels); 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 struct CFG_CLIENT* client = config->clients; DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "init_connections called, instance=%p, config=%p, clients=%p, connections_count=%d", diff --git a/src/etcp_connections.h b/src/etcp_connections.h index eb097c5a..e220c3fc 100644 --- a/src/etcp_connections.h +++ b/src/etcp_connections.h @@ -159,7 +159,10 @@ struct ETCP_LINK { 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); // SOCKET FUNCTIONS diff --git a/src/pkt_normalizer.c b/src/pkt_normalizer.c index 3a086cd9..0921fb01 100644 --- a/src/pkt_normalizer.c +++ b/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"); - int ret = queue_data_put(pn->input, entry, 0); + int ret = queue_data_put(pn->input, entry); pn->in_total_pkts++; pn->in_total_bytes += len; DEBUG_TRACE(DEBUG_CATEGORY_NORMALIZER, "PUT to input end"); @@ -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); 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 передан во фрагмент, не освобождаем) pn->data = NULL; 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) { 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); - queue_data_put(pn->output, pn->recvpart, 0); + queue_data_put(pn->output, pn->recvpart); pn->out_total_pkts++; pn->out_total_bytes += pn->recvpart->len; pn->recvpart = NULL; diff --git a/src/route_bgp.c b/src/route_bgp.c index 7821cd4b..8c7bf52f 100644 --- a/src/route_bgp.c +++ b/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) { if (!instance) { 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->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"); if (!bgp->senders_list) { u_free(bgp); @@ -445,6 +549,11 @@ void route_bgp_destroy(struct UTUN_INSTANCE* instance) { } queue_free(instance->bgp->senders_list); + // Free my_nodeinfo + if (instance->bgp->my_nodeinfo) { + u_free(instance->bgp->my_nodeinfo); + } + u_free(instance->bgp); instance->bgp = NULL; } @@ -488,7 +597,7 @@ void route_bgp_new_conn(struct ETCP_CONN* conn) { if (!item_entry) return; ((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"); } diff --git a/src/route_bgp.h b/src/route_bgp.h index 2bc4590b..c059c5c3 100644 --- a/src/route_bgp.h +++ b/src/route_bgp.h @@ -2,8 +2,10 @@ #define ROUTE_BGP_H #include +#include #include "../lib/ll_queue.h" #include "route_lib.h" +#include "secure_channel.h" // ETCP ID для маршрутных пакетов #define ETCP_ID_ROUTE_ENTRY 0x01 @@ -31,6 +33,50 @@ struct BGP_ROUTE_PACKET { uint64_t hop_list[0]; // flexible array: next_hop → ... → destination } __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 (фиксированный) */ @@ -56,6 +102,8 @@ struct ROUTE_BGP_CONN_ITEM { struct ROUTE_BGP { struct UTUN_INSTANCE* instance; + struct BGP_NODEINFO_PACKET* my_nodeinfo; + uint16_t my_nodeinfo_size; struct ll_queue* senders_list; uint32_t table_version; // starts at 1, for future use }; diff --git a/src/route_bgp.txt b/src/route_bgp.txt index 910ef9fc..eed82156 100644 --- a/src/route_bgp.txt +++ b/src/route_bgp.txt @@ -39,5 +39,33 @@ subcmd: удаленный узел считает свои хеши и сравнивает. где не совпало смотрит сколько записей. если записей не много - передает эти записи. если записей много - добавляет 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. diff --git a/src/routing.c b/src/routing.c index 46659af4..fb37a68b 100644 --- a/src/routing.c +++ b/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) { // 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) { 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); diff --git a/src/tun_if.c b/src/tun_if.c index d93045c8..2225807b 100644 --- a/src/tun_if.c +++ b/src/tun_if.c @@ -59,7 +59,7 @@ static void tun_read_callback(int fd, void* user_arg) pkt->len = nread + 1; // 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"); u_free(packet_data); queue_entry_free(pkt); @@ -293,7 +293,7 @@ void tun_packet_handler(void* arg) { } // Положить в очередь (уже в main thread) - int ok = queue_data_put(tun->output_queue, entry, 0); + int ok = queue_data_put(tun->output_queue, entry); if (ok != 0) { DEBUG_ERROR(DEBUG_CATEGORY_TUN, "Put error"); 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->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) { 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->len > len) { - queue_data_put(tun->input_queue, pkt, 0); + queue_data_put(tun->input_queue, pkt); return -1; } diff --git a/src/utun_instance.c b/src/utun_instance.c index f62d94d7..d433ed64 100644 --- a/src/utun_instance.c +++ b/src/utun_instance.c @@ -45,6 +45,14 @@ static int instance_init_common(struct UTUN_INSTANCE* instance, struct UASYNC* u instance->ua = ua; 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 instance->node_id = config->global.my_node_id; @@ -107,7 +115,13 @@ static int instance_init_common(struct UTUN_INSTANCE* instance, struct UASYNC* u DEBUG_INFO(DEBUG_CATEGORY_TUN, "TUN initialization disabled - skipping TUN device setup"); 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 instance->bgp = route_bgp_init(instance); if (!instance->bgp) { @@ -177,6 +191,15 @@ struct UTUN_INSTANCE* utun_instance_create(struct UASYNC* ua, const char *config 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; } diff --git a/src/utun_instance.h b/src/utun_instance.h index 28c17533..f9197b45 100644 --- a/src/utun_instance.h +++ b/src/utun_instance.h @@ -24,6 +24,9 @@ struct control_server; // uTun instance configuration struct UTUN_INSTANCE { + // Identification + char name[16]; // Instance name from config + // Configuration (moved from utun_state) struct utun_config *config; diff --git a/tests/debug_simple.c b/tests/debug_simple.c index 6aaab2ea..1a77a753 100644 --- a/tests/debug_simple.c +++ b/tests/debug_simple.c @@ -26,7 +26,7 @@ int main() { printf("Data pointer: %p\n", data1); /* 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("After put: queue count = %d\n", queue_entry_count(q)); diff --git a/tests/test_intensive_memory_pool.c b/tests/test_intensive_memory_pool.c index 31e5a5d5..801fbe56 100644 --- a/tests/test_intensive_memory_pool.c +++ b/tests/test_intensive_memory_pool.c @@ -39,7 +39,7 @@ static double test_without_pools(int iterations) { // Добавить записи for (int i = 0; i < 10; i++) { void* data = queue_entry_new(64); - queue_data_put(queue, data, i); // Используем ID = i + queue_data_put(queue, data); // Используем ID = i } // Удалить записи (триггер waiters) @@ -91,7 +91,7 @@ static double test_with_pools(int iterations) { for (int i = 0; i < 10; i++) { void* data = queue_entry_new_from_pool(pool); if (data) { - queue_data_put(queue, data, i); // Используем ID = i + queue_data_put(queue, data); // Используем ID = i } } diff --git a/tests/test_intensive_memory_pool_new.c b/tests/test_intensive_memory_pool_new.c index 7d8d3e58..125923d4 100644 --- a/tests/test_intensive_memory_pool_new.c +++ b/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++) { void* data = queue_entry_new(64); - queue_data_put(queue, data, i); // Используем ID = i + queue_data_put(queue, data); // Используем ID = i } // Удалить записи (триггер waiters) @@ -91,7 +91,7 @@ static double test_with_pools(int iterations) { for (int i = 0; i < 10; i++) { void* data = queue_entry_new_from_pool(pool); if (data) { - queue_data_put(queue, data, i); // Используем ID = i + queue_data_put(queue, data); // Используем ID = i } } diff --git a/tests/test_ll_queue.c b/tests/test_ll_queue.c index 5de319d8..ff661043 100644 --- a/tests/test_ll_queue.c +++ b/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)); d->id = i; snprintf(d->name, sizeof(d->name), "item%d", 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); } ASSERT_EQ(queue_entry_count(q), 10, ""); @@ -119,12 +119,12 @@ static void test_lifo_priority(void) { d->id = i; d->value = i; 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)); 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); ASSERT(first && first->id == 999, "priority first"); @@ -150,7 +150,7 @@ static void test_callback(void) { for (int i = 0; i < 5; i++) { test_data_t *d = (test_data_t*)queue_entry_new(sizeof(*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); @@ -172,7 +172,7 @@ static void test_waiter(void) { for (int i = 0; i < 5; 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; @@ -197,15 +197,16 @@ static void test_limits_hash(void) { for (int i = 0; i < 3; i++) { 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)); - 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"); 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)); 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); 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); 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 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); ASSERT(d3 != NULL, "alloc after free failed"); 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 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)); d->id = rand() % 10000; 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++; } while (queue_entry_count(q)) queue_entry_free((struct ll_entry*)queue_data_get(q)); diff --git a/tests/test_memory_pool_and_config.c b/tests/test_memory_pool_and_config.c index 2f8ab7f3..d63b05e7 100644 --- a/tests/test_memory_pool_and_config.c +++ b/tests/test_memory_pool_and_config.c @@ -59,7 +59,7 @@ int main() { // Add some entries and trigger waiters for (int i = 0; i < 5; i++) { 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 diff --git a/tests/test_pkt_normalizer_standalone.c b/tests/test_pkt_normalizer_standalone.c index f324d2e6..a1e2e54f 100644 --- a/tests/test_pkt_normalizer_standalone.c +++ b/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) { fragments_sent++; // 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);