Browse Source

Add async queue wait functionality and enhance packet normalizer tests

- Implement queue_wait_threshold() with automatic waiter checking in ll_queue
- Add pkt_normalizer_flush() to force buffer sending
- Enhance test suite with async wait tests, fragmentation verification, and buffer flush tests
- Fix test edge cases by flushing packer buffer between subtests
- Improve test reliability with increased iteration limits and while loop processing
- Add connection module for secure UDP communication with ECC/AES-CCM cryptography
v2_dev
jek 9 months ago
parent
commit
d57ab133c5
  1. 12
      Makefile
  2. 724
      connection.c
  3. 118
      connection.h
  4. 528
      etcp.c
  5. 202
      etcp.h
  6. 179
      etcp.h1
  7. 337
      etcp.txt
  8. 109
      ll_queue.c
  9. 38
      ll_queue.h
  10. 195
      pkt_normalizer.c
  11. 4
      pkt_normalizer.h
  12. 2842
      stress.log
  13. 2834
      stress.out
  14. 1355
      stress_50.log
  15. 1098
      stress_new.log
  16. 1
      tests/simple_uasync.h
  17. 1
      tests/test_etcp_simple.c
  18. 12
      tests/test_etcp_stress.c
  19. 373
      tests/test_pkt_normalizer.c
  20. 87
      u_async.c
  21. 3
      u_async.h

12
Makefile

@ -1,5 +1,5 @@
CC := gcc
CFLAGS := -Os -std=c99 -Wall -Wextra -D_ISOC99_SOURCE -DENABLE_TESTS
CFLAGS := -Os -std=c99 -Wall -Wextra -D_ISOC99_SOURCE -DENABLE_TESTS -DETCP_DEBUG -DETCP_DEBUG_EXT
INCLUDES := -Itinycrypt/lib/include/ -Itinycrypt/lib/source/ -Itinycrypt/tests/include/ -I.
TEST_DIR := tests
@ -21,7 +21,7 @@ PN_OBJS := pkt_normalizer.o settings.o
LL_QUEUE_OBJS := ll_queue.o
ETCP_OBJS := etcp.o
all: $(TEST_DIR)/test_ecc_encrypt $(TEST_DIR)/test_sc_lib $(TEST_DIR)/test_udp_secure $(TEST_DIR)/test_pkt_normalizer $(TEST_DIR)/test_etcp $(TEST_DIR)/test_etcp_stress $(TEST_DIR)/test_etcp_simple
all: $(TEST_DIR)/test_ecc_encrypt $(TEST_DIR)/test_sc_lib $(TEST_DIR)/test_udp_secure $(TEST_DIR)/test_pkt_normalizer $(TEST_DIR)/test_etcp $(TEST_DIR)/test_etcp_stress $(TEST_DIR)/test_etcp_simple $(TEST_DIR)/test_connection $(TEST_DIR)/test_connection_stress
$(TEST_DIR)/test_pkt_normalizer: $(TEST_DIR)/test_pkt_normalizer.o $(PN_OBJS) $(LL_QUEUE_OBJS) $(UASYNC_OBJS)
$(CC) $(CFLAGS) $(INCLUDES) -o $@ $^
@ -35,6 +35,12 @@ $(TEST_DIR)/test_etcp_stress: $(TEST_DIR)/test_etcp_stress.o $(ETCP_OBJS) $(LL_Q
$(TEST_DIR)/test_etcp_simple: $(TEST_DIR)/test_etcp_simple.o $(ETCP_OBJS) $(LL_QUEUE_OBJS) $(TEST_DIR)/simple_uasync.o
$(CC) $(CFLAGS) $(INCLUDES) -o $@ $^
$(TEST_DIR)/test_connection: $(TEST_DIR)/test_connection.o connection.o $(ETCP_OBJS) $(PN_OBJS) $(LL_QUEUE_OBJS) $(UASYNC_OBJS) $(SC_LIB_OBJS) $(TINYCRYPT_OBJS)
$(CC) $(CFLAGS) $(INCLUDES) -o $@ $^
$(TEST_DIR)/test_connection_stress: $(TEST_DIR)/test_connection_stress.o connection.o $(ETCP_OBJS) $(PN_OBJS) $(LL_QUEUE_OBJS) $(UASYNC_OBJS) $(SC_LIB_OBJS) $(TINYCRYPT_OBJS)
$(CC) $(CFLAGS) $(INCLUDES) -o $@ $^
$(TEST_DIR)/test_sc_lib: $(TEST_DIR)/test_sc_lib.o $(SC_LIB_OBJS) $(TINYCRYPT_OBJS)
$(CC) $(CFLAGS) $(INCLUDES) -o $@ $^
@ -51,7 +57,7 @@ $(TEST_DIR)/test_etcp: $(TEST_DIR)/test_etcp.o $(ETCP_OBJS) $(LL_QUEUE_OBJS) $(U
$(CC) $(CFLAGS) $(INCLUDES) -c $< -o $@
clean:
rm -f $(TEST_DIR)/test_ecc_encrypt $(TEST_DIR)/test_sc_lib $(TEST_DIR)/test_udp_secure $(TEST_DIR)/test_pkt_normalizer $(TEST_DIR)/test_etcp $(TEST_DIR)/test_etcp_stress $(TEST_DIR)/test_etcp_simple \
rm -f $(TEST_DIR)/test_ecc_encrypt $(TEST_DIR)/test_sc_lib $(TEST_DIR)/test_udp_secure $(TEST_DIR)/test_pkt_normalizer $(TEST_DIR)/test_etcp $(TEST_DIR)/test_etcp_stress $(TEST_DIR)/test_etcp_simple $(TEST_DIR)/test_connection $(TEST_DIR)/test_connection_stress \
*.o tinycrypt/lib/source/*.o $(TEST_DIR)/*.o
.PHONY: all clean

724
connection.c

@ -0,0 +1,724 @@
// connection.c - Реализация минималистичного API для защищенных UDP подключений
#include "connection.h"
#include "sc_lib.h"
#include "pkt_normalizer.h"
#include "etcp.h"
#include "u_async.h"
#include "ll_queue.h"
#include "settings.h"
#include <stdlib.h>
#include <string.h>
#include <stdio.h>
#include <errno.h>
/* Сетевые заголовки */
#include <unistd.h>
#include <sys/socket.h>
#include <netinet/in.h>
#include <arpa/inet.h>
#include <fcntl.h>
/* Внутренняя структура подключения */
struct conn_handle {
/* Сетевая часть */
int sockfd;
struct sockaddr_in local_addr;
struct sockaddr_in remote_addr;
conn_mode_t mode;
uint8_t remote_defined; /* 1 если удаленный адрес определен */
uint8_t socket_connected; /* 1 если сокет создан и настроен */
/* Криптография */
sc_context_t crypto_ctx;
uint8_t keys_set; /* 1 если ключи установлены */
uint8_t crypto_ready; /* 1 если сессионный ключ вычислен */
/* Обработка данных */
epkt_t* etcp;
pkt_normalizer_pair* normalizer;
ll_queue_t* app_input_queue; /* Данные от приложения для отправки */
ll_queue_t* app_output_queue; /* Собранные данные для приложения */
/* Callback'и */
conn_recv_callback_t recv_callback;
void* recv_callback_user;
/* Состояние */
uint8_t is_closing;
uint8_t is_destroying;
/* Таймеры */
void* socket_read_timer;
/* Статистика */
conn_stats_t stats;
};
/* Внутренние функции */
static int create_udp_socket(const char* ip, uint16_t port, struct sockaddr_in* addr);
static void socket_read_callback(void* arg);
static void etcp_tx_callback(epkt_t* epkt, uint8_t* data, uint16_t len, void* arg);
static void app_output_callback(ll_queue_t* q, ll_entry_t* entry, void* arg);
static void process_received_packet(conn_handle_t* conn, uint8_t* data, size_t len);
static void update_stats_from_etcp(conn_handle_t* conn);
/* Bridge functions for data flow */
static void app_input_bridge(ll_queue_t* q, ll_entry_t* entry, void* arg);
static void packer_output_bridge(ll_queue_t* q, ll_entry_t* entry, void* arg);
static void etcp_output_bridge(ll_queue_t* q, ll_entry_t* entry, void* arg);
/* Глобальная инициализация uasync (вызывается один раз) */
static int uasync_initialized = 0;
/* ==================== Публичные функции ==================== */
conn_handle_t* conn_create(void)
{
/* Инициализация uasync глобально */
if (!uasync_initialized) {
uasync_init();
uasync_initialized = 1;
}
/* Выделение памяти */
conn_handle_t* conn = calloc(1, sizeof(conn_handle_t));
if (!conn) {
return NULL;
}
/* Инициализация очередей приложения */
conn->app_input_queue = queue_new();
conn->app_output_queue = queue_new();
if (!conn->app_input_queue || !conn->app_output_queue) {
if (conn->app_input_queue) queue_free(conn->app_input_queue);
if (conn->app_output_queue) queue_free(conn->app_output_queue);
free(conn);
return NULL;
}
/* Настройка callback для выходной очереди */
queue_set_callback(conn->app_output_queue, app_output_callback, conn);
/* Инициализация статистики */
memset(&conn->stats, 0, sizeof(conn->stats));
/* Начальное состояние */
conn->sockfd = -1;
conn->mode = CONN_MODE_CLIENT; /* по умолчанию клиент */
conn->is_closing = 0;
conn->is_destroying = 0;
return conn;
}
int conn_set_keys(conn_handle_t* conn,
const uint8_t* my_pub_key,
const uint8_t* my_priv_key,
const uint8_t* peer_pub_key)
{
if (!conn || conn->socket_connected) {
return -1; /* Ключи нужно устанавливать до подключения */
}
sc_status_t status;
if (my_pub_key && my_priv_key) {
/* Использование предоставленных ключей */
status = sc_init_local_keys(&conn->crypto_ctx, my_pub_key, my_priv_key);
if (status != SC_OK) {
return -1;
}
} else {
/* Автоматическая генерация ключей */
status = sc_generate_keypair(&conn->crypto_ctx);
if (status != SC_OK) {
return -1;
}
}
/* Установка ключа пира, если предоставлен */
if (peer_pub_key) {
status = sc_set_peer_public_key(&conn->crypto_ctx, peer_pub_key);
if (status != SC_OK) {
return -1;
}
conn->crypto_ready = 1;
} else {
/* Для сервера ключ пира будет получен из первого пакета */
conn->crypto_ready = 0;
}
conn->keys_set = 1;
return 0;
}
int conn_connect(conn_handle_t* conn,
const char* local_ip,
uint16_t local_port,
const char* remote_ip,
uint16_t remote_port,
conn_mode_t mode)
{
if (!conn || conn->socket_connected) {
return -1; /* Уже подключен */
}
/* Проверка ключей */
if (!conn->keys_set) {
/* Автоматическая установка ключей по умолчанию */
if (conn_set_keys(conn, NULL, NULL, NULL) != 0) {
return -1;
}
}
/* Создание UDP сокета */
conn->sockfd = create_udp_socket(local_ip, local_port, &conn->local_addr);
if (conn->sockfd < 0) {
return -1;
}
/* Сохранение режима */
conn->mode = mode;
/* Настройка удаленного адреса для клиента */
if (mode == CONN_MODE_CLIENT && remote_ip) {
memset(&conn->remote_addr, 0, sizeof(conn->remote_addr));
conn->remote_addr.sin_family = AF_INET;
conn->remote_addr.sin_port = htons(remote_port);
if (inet_pton(AF_INET, remote_ip, &conn->remote_addr.sin_addr) != 1) {
close(conn->sockfd);
conn->sockfd = -1;
return -1;
}
conn->remote_defined = 1;
} else {
/* Сервер или клиент без указанного адреса */
conn->remote_defined = 0;
}
/* Инициализация нормализатора пакетов */
conn->normalizer = pkt_normalizer_pair_init();
if (!conn->normalizer) {
close(conn->sockfd);
conn->sockfd = -1;
return -1;
}
/* Инициализация ETCP */
conn->etcp = etcp_init();
if (!conn->etcp) {
pkt_normalizer_pair_deinit(conn->normalizer);
close(conn->sockfd);
conn->sockfd = -1;
return -1;
}
/* Настройка callback'ов ETCP */
etcp_set_callback(conn->etcp, etcp_tx_callback, conn);
/* Начальная полоса пропускания: максимальная для uint16_t ≈ 6553 байт/мс */
etcp_set_bandwidth(conn->etcp, 65535); /* в единицах 0.1мс: 6553.5 * 10 */
/* Связывание нормализатора с ETCP */
/* Packer: app_input_queue -> normalizer -> etcp_tx_put */
/* Unpacker: etcp_output_queue -> normalizer -> app_output_queue */
/* Set up bridge callbacks */
/* 1. App input -> packer input */
queue_set_callback(conn->app_input_queue, app_input_bridge, conn);
/* 2. Packer output -> ETCP tx */
if (conn->normalizer->packer && conn->normalizer->packer->output) {
queue_set_callback(conn->normalizer->packer->output, packer_output_bridge, conn);
}
/* 3. ETCP output -> unpacker input */
ll_queue_t* etcp_output_q = etcp_get_output_queue(conn->etcp);
if (etcp_output_q) {
queue_set_callback(etcp_output_q, etcp_output_bridge, conn);
}
/* 4. Unpacker output -> app output queue */
if (conn->normalizer->unpacker && conn->normalizer->unpacker->output) {
queue_set_callback(conn->normalizer->unpacker->output, app_output_callback, conn);
}
/* Запуск таймера чтения сокета */
conn->socket_read_timer = uasync_set_timeout(1, conn, socket_read_callback);
conn->socket_connected = 1;
return 0;
}
void conn_set_recv_callback(conn_handle_t* conn,
conn_recv_callback_t callback,
void* user_data)
{
if (!conn) return;
conn->recv_callback = callback;
conn->recv_callback_user = user_data;
}
int conn_send(conn_handle_t* conn, const uint8_t* data, size_t len)
{
if (!conn || !data || len == 0 || !conn->socket_connected || conn->is_closing) {
return -1;
}
/* Проверка готовности криптографии (для клиента) */
/* TODO: implement proper key exchange */
/* if (conn->mode == CONN_MODE_CLIENT && !conn->crypto_ready) {
return -1;
} */
/* Создание записи очереди */
ll_entry_t* entry = queue_entry_new(len);
if (!entry) {
return -1;
}
/* Копирование данных */
memcpy(ll_entry_data(entry), data, len);
/* Добавление в очередь ввода приложения */
int result = queue_entry_put(conn->app_input_queue, entry);
if (result != 0) {
queue_entry_free(entry);
return -1;
}
/* Обновление статистики */
conn->stats.bytes_sent += len;
/* Активация обработки (через callback очереди) */
queue_resume_callback(conn->app_input_queue);
return 0;
}
void conn_close(conn_handle_t* conn)
{
if (!conn || conn->is_closing) {
return;
}
conn->is_closing = 1;
/* Отмена таймеров */
if (conn->socket_read_timer) {
uasync_cancel_timeout(conn->socket_read_timer);
conn->socket_read_timer = NULL;
}
/* Закрытие сокета */
if (conn->sockfd >= 0) {
close(conn->sockfd);
conn->sockfd = -1;
}
conn->socket_connected = 0;
}
void conn_destroy(conn_handle_t* conn)
{
if (!conn) return;
conn->is_destroying = 1;
/* Закрытие подключения если активно */
conn_close(conn);
/* Освобождение ETCP */
if (conn->etcp) {
etcp_free(conn->etcp);
conn->etcp = NULL;
}
/* Освобождение нормализатора */
if (conn->normalizer) {
pkt_normalizer_pair_deinit(conn->normalizer);
conn->normalizer = NULL;
}
/* Освобождение очередей */
if (conn->app_input_queue) {
queue_free(conn->app_input_queue);
conn->app_input_queue = NULL;
}
if (conn->app_output_queue) {
queue_free(conn->app_output_queue);
conn->app_output_queue = NULL;
}
/* Очистка криптографического контекста */
memset(&conn->crypto_ctx, 0, sizeof(conn->crypto_ctx));
/* Освобождение дескриптора */
free(conn);
}
int conn_get_stats(conn_handle_t* conn, conn_stats_t* stats)
{
if (!conn || !stats) {
return -1;
}
/* Обновление статистики из ETCP */
update_stats_from_etcp(conn);
/* Копирование статистики */
memcpy(stats, &conn->stats, sizeof(conn_stats_t));
return 0;
}
/* ==================== Внутренние функции ==================== */
/* Bridge functions implementation */
static void app_input_bridge(ll_queue_t* q, ll_entry_t* entry, void* arg)
{
(void)entry; /* unused parameter - we'll get entry from queue */
conn_handle_t* conn = (conn_handle_t*)arg;
if (!conn || !q || !conn->normalizer || !conn->normalizer->packer) {
return;
}
/* Remove entry from app_input_queue */
ll_entry_t* entry_to_forward = queue_entry_get(q);
if (!entry_to_forward) {
queue_resume_callback(q);
return;
}
size_t len = ll_entry_size(entry_to_forward);
printf("[CONN DEBUG] app_input_bridge: forwarding packet len=%zu to packer\n", len);
/* Forward entry from app_input_queue to packer's input queue */
queue_entry_put(conn->normalizer->packer->input, entry_to_forward);
queue_resume_callback(q);
}
static void packer_output_bridge(ll_queue_t* q, ll_entry_t* entry, void* arg)
{
(void)entry; /* unused parameter - we'll get entry from queue */
printf("[CONN DEBUG] packer_output_bridge ENTER, queue count=%d\n", queue_entry_count(q));
conn_handle_t* conn = (conn_handle_t*)arg;
if (!conn || !q || !conn->etcp) {
queue_resume_callback(q);
return;
}
/* Remove entry from queue */
ll_entry_t* entry_to_process = queue_entry_get(q);
if (!entry_to_process) {
queue_resume_callback(q);
return;
}
/* Extract data from entry */
uint8_t* data = ll_entry_data(entry_to_process);
size_t len = ll_entry_size(entry_to_process);
printf("[CONN DEBUG] packer_output_bridge: processing packet len=%zu\n", len);
if (len > 65535) {
/* Too large for ETCP (max 64KB) */
queue_entry_free(entry_to_process);
queue_resume_callback(q);
return;
}
uint8_t* data_to_send = data;
uint16_t data_len = (uint16_t)len;
/* Encrypt if crypto is ready - TEMPORARILY DISABLED */
if (conn->crypto_ready) {
/* Allocate buffer for ciphertext + tag */
uint8_t* ciphertext = malloc(len);
uint8_t tag[SC_TAG_SIZE];
sc_status_t status = sc_encrypt(&conn->crypto_ctx, data, len, ciphertext, tag);
if (status == SC_OK) {
/* Combine ciphertext and tag into single buffer */
uint8_t* encrypted_data = malloc(len + SC_TAG_SIZE);
if (encrypted_data) {
memcpy(encrypted_data, ciphertext, len);
memcpy(encrypted_data + len, tag, SC_TAG_SIZE);
data_to_send = encrypted_data;
data_len = (uint16_t)(len + SC_TAG_SIZE);
}
/* Free ciphertext buffer */
free(ciphertext);
}
/* If encryption failed, fall back to plaintext */
}
/* Note: etcp_tx_put copies data, we can free entry after call */
int result = etcp_tx_put(conn->etcp, data_to_send, data_len);
/* Free encrypted buffer if it was allocated */
if (conn->crypto_ready && data_to_send != data) {
free(data_to_send);
}
queue_entry_free(entry_to_process);
(void)result; /* Ignore result for now */
queue_resume_callback(q);
}
static void etcp_output_bridge(ll_queue_t* q, ll_entry_t* entry, void* arg)
{
(void)entry; /* unused parameter - we'll get entry from queue */
conn_handle_t* conn = (conn_handle_t*)arg;
if (!conn || !q || !conn->normalizer || !conn->normalizer->unpacker) {
return;
}
/* Remove entry from ETCP output queue */
ll_entry_t* entry_to_forward = queue_entry_get(q);
if (!entry_to_forward) {
queue_resume_callback(q);
return;
}
/* Extract data from entry */
uint8_t* data = ll_entry_data(entry_to_forward);
size_t len = ll_entry_size(entry_to_forward);
ll_entry_t* new_entry = NULL;
/* Decrypt if crypto is ready - TEMPORARILY DISABLED */
if (conn->crypto_ready) {
size_t ciphertext_len = len - SC_TAG_SIZE;
uint8_t* ciphertext = data;
uint8_t* tag = data + ciphertext_len;
/* Allocate buffer for plaintext */
uint8_t* plaintext = malloc(ciphertext_len);
if (plaintext) {
sc_status_t status = sc_decrypt(&conn->crypto_ctx, ciphertext, ciphertext_len, tag, plaintext);
if (status == SC_OK) {
/* Create new entry with plaintext */
new_entry = queue_entry_new(ciphertext_len);
if (new_entry) {
memcpy(ll_entry_data(new_entry), plaintext, ciphertext_len);
}
}
free(plaintext);
}
}
/* If we created a new entry, replace the old one */
if (new_entry) {
queue_entry_free(entry_to_forward);
entry_to_forward = new_entry;
}
/* Forward entry to unpacker's input queue */
queue_entry_put(conn->normalizer->unpacker->input, entry_to_forward);
queue_resume_callback(q);
}
static int create_udp_socket(const char* ip, uint16_t port, struct sockaddr_in* addr)
{
int sock = socket(AF_INET, SOCK_DGRAM, 0);
if (sock < 0) {
return -1;
}
/* Установка non-blocking режима */
int flags = fcntl(sock, F_GETFL, 0);
if (flags < 0) {
close(sock);
return -1;
}
if (fcntl(sock, F_SETFL, flags | O_NONBLOCK) < 0) {
close(sock);
return -1;
}
/* Настройка адреса */
memset(addr, 0, sizeof(*addr));
addr->sin_family = AF_INET;
addr->sin_port = htons(port);
if (ip && ip[0] != '\0') {
if (inet_pton(AF_INET, ip, &addr->sin_addr) != 1) {
close(sock);
return -1;
}
} else {
addr->sin_addr.s_addr = INADDR_ANY;
}
/* Bind к адресу */
if (bind(sock, (struct sockaddr*)addr, sizeof(*addr)) < 0) {
close(sock);
return -1;
}
return sock;
}
static void socket_read_callback(void* arg)
{
conn_handle_t* conn = (conn_handle_t*)arg;
if (!conn || conn->is_closing || conn->sockfd < 0) {
return;
}
uint8_t buffer[2048]; /* Максимальный размер UDP пакета */
struct sockaddr_in from_addr;
socklen_t from_len = sizeof(from_addr);
/* Чтение всех доступных пакетов */
while (1) {
ssize_t received = recvfrom(conn->sockfd, buffer, sizeof(buffer), 0,
(struct sockaddr*)&from_addr, &from_len);
if (received <= 0) {
if (errno == EAGAIN || errno == EWOULDBLOCK) {
/* Нет больше данных */
break;
}
/* Ошибка чтения */
break;
}
/* Обновление статистики */
conn->stats.bytes_received += received;
conn->stats.packets_received++;
printf("[CONN DEBUG] socket_read_callback: received packet len=%zd, total_received=%u\n", received, conn->stats.packets_received);
/* Для сервера: установка удаленного адреса при первом пакете */
if (!conn->remote_defined && conn->mode == CONN_MODE_SERVER) {
memcpy(&conn->remote_addr, &from_addr, sizeof(conn->remote_addr));
conn->remote_defined = 1;
}
/* Обработка полученного пакета */
process_received_packet(conn, buffer, received);
}
/* Перезапуск таймера чтения */
if (!conn->is_closing) {
conn->socket_read_timer = uasync_set_timeout(1, conn, socket_read_callback);
}
}
static void etcp_tx_callback(epkt_t* epkt, uint8_t* data, uint16_t len, void* arg)
{
(void)epkt; /* unused parameter */
conn_handle_t* conn = (conn_handle_t*)arg;
if (!conn || conn->is_closing || conn->sockfd < 0) {
free(data); /* Данные были выделены в etcp.c */
return;
}
/* Шифрование данных если криптография готова */
uint8_t* packet_to_send = data;
uint16_t packet_len = len;
if (conn->crypto_ready) {
/* TODO: реализовать шифрование через sc_encrypt */
/* Временно пропускаем шифрование для тестирования */
}
/* Отправка через UDP сокет */
if (conn->remote_defined) {
printf("[CONN DEBUG] etcp_tx_callback: sending packet len=%u, stats.packets_sent=%u\n", packet_len, conn->stats.packets_sent);
ssize_t sent = sendto(conn->sockfd, packet_to_send, packet_len, 0,
(struct sockaddr*)&conn->remote_addr,
sizeof(conn->remote_addr));
if (sent == (ssize_t)packet_len) {
conn->stats.packets_sent++;
printf("[CONN DEBUG] etcp_tx_callback: sent successfully, new count=%u\n", conn->stats.packets_sent);
} else {
printf("[CONN DEBUG] etcp_tx_callback: send failed, sent=%zd, errno=%d\n", sent, errno);
}
}
/* Освобождение данных (выделены в etcp.c) - packet is stored in ETCP sent_list for retransmission */
/* DO NOT free here */
(void)packet_to_send; /* unused */
}
static void app_output_callback(ll_queue_t* q, ll_entry_t* entry, void* arg)
{
(void)entry; /* unused parameter - we'll get entry from queue */
conn_handle_t* conn = (conn_handle_t*)arg;
if (!conn || !q || conn->is_closing) {
return;
}
/* Remove entry from unpacker output queue */
ll_entry_t* entry_to_process = queue_entry_get(q);
if (!entry_to_process) {
queue_resume_callback(q);
return;
}
/* Получение данных из записи */
uint8_t* data = ll_entry_data(entry_to_process);
size_t len = ll_entry_size(entry_to_process);
/* Вызов пользовательского callback'а если установлен */
if (conn->recv_callback) {
conn->recv_callback(conn, data, len, conn->recv_callback_user);
}
/* Освобождение записи */
queue_entry_free(entry_to_process);
/* Обновление статистики */
conn->stats.fragments_assembled++;
queue_resume_callback(q);
}
static void process_received_packet(conn_handle_t* conn, uint8_t* data, size_t len)
{
if (!conn || !data || len == 0) {
return;
}
/* Расшифрование если криптография готова */
uint8_t* packet_to_process = data;
uint16_t packet_len = len;
if (conn->crypto_ready) {
/* TODO: реализовать расшифрование через sc_decrypt */
/* Временно пропускаем для тестирования */
} else {
/* Криптография не готова, но ключи должны быть установлены заранее */
printf("[CONN] WARNING: Received packet but crypto not ready (keys should be pre-shared)\n");
/* Все равно передаем пакет в ETCP для обработки */
/* В реальном использовании это ошибка конфигурации */
}
/* Передача пакета в ETCP для обработки */
printf("[CONN DEBUG] process_received_packet: forwarding to ETCP, len=%u\n", packet_len);
etcp_rx_input(conn->etcp, packet_to_process, packet_len);
}
static void update_stats_from_etcp(conn_handle_t* conn)
{
if (!conn || !conn->etcp) {
return;
}
/* Получение метрик из ETCP */
conn->stats.current_rtt_ms = etcp_get_rtt(conn->etcp) / 10; /* 0.1ms -> ms */
conn->stats.jitter_ms = etcp_get_jitter(conn->etcp) / 10;
/* TODO: получение счетчика ретрансмиссий из ETCP */
}

118
connection.h

@ -0,0 +1,118 @@
// connection.h - Минималистичный API для защищенных UDP подключений
#ifndef CONNECTION_H
#define CONNECTION_H
#include <stdint.h>
#include <stddef.h>
#ifdef __cplusplus
extern "C" {
#endif
/* Непрозрачный дескриптор подключения */
typedef struct conn_handle conn_handle_t;
/* Режим подключения */
typedef enum {
CONN_MODE_CLIENT, /* Инициируем подключение к указанному удаленному адресу */
CONN_MODE_SERVER /* Ожидаем входящие подключения */
} conn_mode_t;
/* Callback для входящих данных */
typedef void (*conn_recv_callback_t)(conn_handle_t* conn,
const uint8_t* data,
size_t len,
void* user_data);
/*
* Создание дескриптора подключения (только выделение памяти).
* Возвращает NULL при ошибке.
*/
conn_handle_t* conn_create(void);
/*
* Установка криптографических ключей.
* Должна быть вызвана до conn_connect().
*
* @param conn Дескриптор подключения
* @param my_pub_key Публичный ключ (64 байта для secp256r1), NULL для авто-генерации
* @param my_priv_key Приватный ключ (32 байта), NULL для авто-генерации
* @param peer_pub_key Публичный ключ пира (64 байта), NULL для сервера (получит из первого пакета)
*
* @return 0 при успехе, -1 при ошибке
*/
int conn_set_keys(conn_handle_t* conn,
const uint8_t* my_pub_key,
const uint8_t* my_priv_key,
const uint8_t* peer_pub_key);
/*
* Подключение к удаленному узлу или начало ожидания входящих подключений.
*
* @param conn Дескриптор подключения
* @param local_ip Локальный IP для bind (NULL для "0.0.0.0")
* @param local_port Локальный порт (0 для авто-выбора)
* @param remote_ip Удаленный IP (NULL для серверного режима)
* @param remote_port Удаленный порт (игнорируется если remote_ip NULL)
* @param mode Режим подключения (CONN_MODE_CLIENT/SERVER)
*
* @return 0 при успехе, -1 при ошибке
*/
int conn_connect(conn_handle_t* conn,
const char* local_ip,
uint16_t local_port,
const char* remote_ip,
uint16_t remote_port,
conn_mode_t mode);
/*
* Установка callback'а для входящих данных.
* Callback будет вызываться при получении полных собранных пакетов.
*/
void conn_set_recv_callback(conn_handle_t* conn,
conn_recv_callback_t callback,
void* user_data);
/*
* Отправка данных.
* Данные будут автоматически фрагментированы, зашифрованы и отправлены.
*
* @return 0 при успехе, -1 при ошибке
*/
int conn_send(conn_handle_t* conn, const uint8_t* data, size_t len);
/*
* Закрытие подключения (немедленное, без протокола завершения).
* После вызова дескриптор можно уничтожить через conn_destroy().
*/
void conn_close(conn_handle_t* conn);
/*
* Полное уничтожение дескриптора подключения и освобождение всех ресурсов.
* Автоматически вызывает conn_close() если подключение активно.
*/
void conn_destroy(conn_handle_t* conn);
/*
* Получение статистики подключения (опционально).
*
* @return 0 при успехе, -1 при ошибке
*/
typedef struct {
uint64_t bytes_sent;
uint64_t bytes_received;
uint32_t packets_sent;
uint32_t packets_received;
uint32_t retransmissions;
uint32_t fragments_assembled;
uint16_t current_rtt_ms; /* Текущее RTT в миллисекундах */
uint16_t jitter_ms; /* Джиттер в миллисекундах */
} conn_stats_t;
int conn_get_stats(conn_handle_t* conn, conn_stats_t* stats);
#ifdef __cplusplus
}
#endif
#endif /* CONNECTION_H */

528
etcp.c

@ -4,6 +4,9 @@
#include <stdlib.h>
#include <string.h>
#include <stdio.h>
#include <sys/time.h>
#include <time.h>
// Internal structures
typedef struct rx_packet {
@ -24,6 +27,7 @@ typedef struct sent_packet {
uint16_t payload_len; // Payload length (for window accounting)
uint16_t send_time;
uint8_t need_ack;
uint8_t need_retransmit;
} sent_packet_t;
// Forward declarations of internal functions
@ -33,6 +37,7 @@ static void update_metrics(epkt_t* epkt, uint16_t rtt);
static uint16_t get_current_timestamp(void);
static uint16_t timestamp_diff(uint16_t t1, uint16_t t2);
static int id_compare(uint16_t id1, uint16_t id2);
static void retransmit_packet(epkt_t* epkt, uint16_t id);
// Timer callbacks
static void tx_timer_callback(void* arg);
@ -78,8 +83,17 @@ epkt_t* etcp_init(void) {
epkt->jitter = 0;
epkt->bytes_sent_total = 0;
// Initialize statistics
epkt->retransmissions_count = 0;
epkt->ack_packets_count = 0;
epkt->control_packets_count = 0;
epkt->total_packets_sent = 0;
epkt->unique_packets_sent = 0;
epkt->bytes_received_total = 0;
// Initialize IDs
epkt->next_tx_id = 1;
epkt->last_sent_id = 0;
epkt->last_rx_id = 0;
epkt->last_delivered_id = 0;
@ -93,13 +107,16 @@ epkt_t* etcp_init(void) {
// Initialize window management
epkt->unacked_bytes = 0;
epkt->window_size = 0;
epkt->last_acked_id = 0;
epkt->last_rx_ack_id = 0;
epkt->retrans_timer_period = 20; // Default 2ms (20 timebase units)
epkt->next_retrans_time = 0;
epkt->window_blocked = 0;
// Forward progress tracking
epkt->oldest_missing_id = 0;
epkt->missing_since_time = 0;
// No timers yet
epkt->next_tx_timer = NULL;
epkt->retransmit_timer = NULL;
@ -174,10 +191,15 @@ void etcp_update_window(epkt_t* epkt) {
}
// window = RTT * bandwidth * 2
// RTT in timebase (0.1ms), bandwidth in bytes per timebase
// Multiply using 32-bit to avoid overflow
uint32_t rtt32 = rtt;
uint32_t bw32 = epkt->bandwidth;
epkt->window_size = rtt32 * bw32 * 2;
// Multiply using 64-bit to avoid overflow, cap at UINT32_MAX
uint64_t rtt64 = rtt;
uint64_t bw64 = epkt->bandwidth;
uint64_t window64 = rtt64 * bw64 * 2;
if (window64 > UINT32_MAX) {
epkt->window_size = UINT32_MAX;
} else {
epkt->window_size = (uint32_t)window64;
}
// Update retransmission timer period: max(RTT/2, 2ms)
uint16_t rtt_half = rtt / 2;
@ -187,6 +209,94 @@ void etcp_update_window(epkt_t* epkt) {
epkt->retrans_timer_period = rtt_half;
}
// Reset connection state
void etcp_reset(epkt_t* epkt) {
if (!epkt) return;
// Cancel timers
if (epkt->next_tx_timer) {
uasync_cancel_timeout(epkt->next_tx_timer);
epkt->next_tx_timer = NULL;
}
if (epkt->retransmit_timer) {
uasync_cancel_timeout(epkt->retransmit_timer);
epkt->retransmit_timer = NULL;
}
// Clear tx queue
ll_entry_t* entry;
while ((entry = queue_entry_get(epkt->tx_queue)) != NULL) {
queue_entry_free(entry);
}
// Clear output queue
while ((entry = queue_entry_get(epkt->output_queue)) != NULL) {
queue_entry_free(entry);
}
// Free rx_list
rx_packet_t* rx = epkt->rx_list;
while (rx) {
rx_packet_t* next = rx->next;
if (rx->data) free(rx->data);
free(rx);
rx = next;
}
epkt->rx_list = NULL;
// Free sent_list
sent_packet_t* sent = epkt->sent_list;
while (sent) {
sent_packet_t* next = sent->next;
if (sent->data) free(sent->data);
free(sent);
sent = next;
}
epkt->sent_list = NULL;
// Reset state
epkt->last_sent_timestamp = get_current_timestamp();
epkt->bytes_allowed = 0;
etcp_update_window(epkt);
// Reset metrics
epkt->rtt_last = 0;
epkt->rtt_avg_10 = 0;
epkt->rtt_avg_100 = 0;
epkt->jitter = 0;
epkt->bytes_sent_total = 0;
// Reset statistics
epkt->retransmissions_count = 0;
epkt->ack_packets_count = 0;
epkt->control_packets_count = 0;
epkt->total_packets_sent = 0;
epkt->unique_packets_sent = 0;
epkt->bytes_received_total = 0;
// Reset IDs
epkt->next_tx_id = 1;
epkt->last_rx_id = 0;
epkt->last_delivered_id = 0;
// Reset history
epkt->rtt_history_idx = 0;
epkt->rtt_history_count = 0;
// Reset pending arrays
epkt->pending_ack_count = 0;
epkt->pending_retransmit_count = 0;
// Reset window management
epkt->unacked_bytes = 0;
epkt->window_size = 0;
epkt->last_acked_id = 0;
epkt->last_rx_ack_id = 0;
epkt->retrans_timer_period = 20;
epkt->next_retrans_time = 0;
epkt->window_blocked = 0;
}
// Get RTT
uint16_t etcp_get_rtt(epkt_t* epkt) {
return epkt ? epkt->rtt_last : 0;
@ -209,9 +319,13 @@ int etcp_tx_put(epkt_t* epkt, uint8_t* data, uint16_t len) {
memcpy(ll_entry_data(entry), data, len);
// Add to queue
ETCP_DEBUG_LOG("etcp_tx_put: adding packet len=%u, tx_queue count before=%d\n", len, queue_entry_count(epkt->tx_queue));
int result = queue_entry_put(epkt->tx_queue, entry);
if (result != 0) {
ETCP_DEBUG_LOG("etcp_tx_put: queue_entry_put failed, result=%d\n", result);
queue_entry_free(entry);
} else {
ETCP_DEBUG_LOG("etcp_tx_put: queued successfully, tx_queue count after=%d\n", queue_entry_count(epkt->tx_queue));
}
return result;
@ -236,10 +350,30 @@ int etcp_tx_queue_size(epkt_t* epkt) {
// Get current timestamp (0.1us timebase)
static uint16_t get_current_timestamp(void) {
// TODO: Implement using u_async time or system time
// For now, return increasing counter
#ifdef ENABLE_TESTS
// For tests, use a simple counter to ensure deterministic behavior
static uint16_t counter = 0;
return counter++;
#else
// Production: use monotonic clock if available
#ifdef CLOCK_MONOTONIC
struct timespec ts;
if (clock_gettime(CLOCK_MONOTONIC, &ts) == 0) {
// Convert to 0.1ms units (100us)
// 1 second = 10,000,000 timebase units (0.1us each)
// But we need modulo 65536 for 16-bit timestamp
uint64_t ns = (uint64_t)ts.tv_sec * 1000000000ULL + (uint64_t)ts.tv_nsec;
uint64_t timebase_units = ns / 100000; // 0.1ms = 100,000ns
return (uint16_t)(timebase_units & 0xFFFF);
}
#endif
// Fallback to gettimeofday (not monotonic but available everywhere)
struct timeval tv;
gettimeofday(&tv, NULL);
uint64_t us = (uint64_t)tv.tv_sec * 1000000ULL + (uint64_t)tv.tv_usec;
uint64_t timebase_units = us / 100; // 0.1ms = 100us
return (uint16_t)(timebase_units & 0xFFFF);
#endif
}
// Calculate positive difference between timestamps (considering wrap-around)
@ -340,9 +474,8 @@ static void tx_process(epkt_t* epkt) {
if (epkt->unacked_bytes + data_len > epkt->window_size) {
// Window full, block transmission
epkt->window_blocked = 1;
// DEBUG
// printf("Window blocked: unacked=%u, data_len=%u, window=%u\n",
// epkt->unacked_bytes, data_len, epkt->window_size);
ETCP_LOG("Window blocked: unacked=%u, data_len=%u, window=%u\n",
epkt->unacked_bytes, data_len, epkt->window_size);
// Put entry back to queue
queue_entry_put_first(epkt->tx_queue, entry);
// Schedule check when window might open (after retransmission timer)
@ -361,13 +494,13 @@ static void tx_process(epkt_t* epkt) {
// Add space for ACKs
if (epkt->pending_ack_count > 0) {
packet_size += 1 + epkt->pending_ack_count * 4; // hdr=0x01 + ids+timestamps
packet_size += 1 + 1 + epkt->pending_ack_count * 4 + 4; // hdr=0x01 + count + ids+timestamps + 2 IDs
}
// Add space for retransmission requests
if (epkt->pending_retransmit_count > 0) {
uint8_t count = epkt->pending_retransmit_count;
if (count > 32) count = 32;
packet_size += 1 + count * 2 + 2; // hdr + IDs + last delivered ID
packet_size += 1 + count * 2 + 4; // hdr + IDs + last delivered ID + last received ID
}
// Check if we have enough bandwidth
@ -399,6 +532,10 @@ static void tx_process(epkt_t* epkt) {
uint16_t id;
if (data_packet) {
id = epkt->next_tx_id++;
// Update last sent ID if newer
if (id != 0 && id_compare(id, epkt->last_sent_id) > 0) {
epkt->last_sent_id = id;
}
} else {
id = 0; // metrics-only packet
}
@ -410,13 +547,21 @@ static void tx_process(epkt_t* epkt) {
// Add ACKs if any
if (epkt->pending_ack_count > 0) {
uint8_t count = epkt->pending_ack_count;
if (count > 32) count = 32;
*ptr++ = 0x01; // hdr for timestamp report
for (int i = 0; i < epkt->pending_ack_count; i++) {
*ptr++ = count; // number of timestamp pairs
for (int i = 0; i < count; i++) {
*ptr++ = epkt->pending_ack_ids[i] >> 8;
*ptr++ = epkt->pending_ack_ids[i] & 0xFF;
*ptr++ = epkt->pending_ack_timestamps[i] >> 8;
*ptr++ = epkt->pending_ack_timestamps[i] & 0xFF;
}
// Add last delivered ID and last received ID
*ptr++ = epkt->last_delivered_id >> 8;
*ptr++ = epkt->last_delivered_id & 0xFF;
*ptr++ = epkt->last_rx_id >> 8;
*ptr++ = epkt->last_rx_id & 0xFF;
epkt->pending_ack_count = 0;
}
@ -424,14 +569,22 @@ static void tx_process(epkt_t* epkt) {
if (epkt->pending_retransmit_count > 0) {
uint8_t count = epkt->pending_retransmit_count;
if (count > 32) count = 32;
ETCP_LOG("Sending retransmit requests: count=%u, IDs: ", count);
for (int i = 0; i < count; i++) {
ETCP_LOG("%u ", epkt->pending_retransmit_ids[i]);
}
ETCP_LOG("last_delivered=%u, last_rx=%u\n", epkt->last_delivered_id, epkt->last_rx_id);
*ptr++ = 0x10 + (count - 1); // hdr with count
for (int i = 0; i < count; i++) {
*ptr++ = epkt->pending_retransmit_ids[i] >> 8;
*ptr++ = epkt->pending_retransmit_ids[i] & 0xFF;
}
// Add last delivered ID
// Add last delivered ID and last received ID
*ptr++ = epkt->last_delivered_id >> 8;
*ptr++ = epkt->last_delivered_id & 0xFF;
*ptr++ = epkt->last_rx_id >> 8;
*ptr++ = epkt->last_rx_id & 0xFF;
epkt->pending_retransmit_count = 0;
}
@ -442,6 +595,22 @@ static void tx_process(epkt_t* epkt) {
ptr += data_len;
}
// Update statistics before sending
epkt->total_packets_sent++;
if (data_packet) {
epkt->unique_packets_sent++;
}
if (epkt->pending_ack_count > 0) {
epkt->ack_packets_count++;
epkt->control_packets_count++;
}
if (epkt->pending_retransmit_count > 0) {
epkt->control_packets_count++;
}
// Send packet
epkt->tx_callback(epkt, packet, packet_size, epkt->tx_callback_arg);
@ -506,6 +675,33 @@ static void retransmit_check(epkt_t* epkt) {
}
sent = sent->next;
}
// Special retransmission for the newest unacked packet after 2×RTT
if (epkt->last_sent_id != 0 && id_compare(epkt->last_sent_id, epkt->last_acked_id) > 0) {
// Find the newest packet in sent_list
sent_packet_t* sent = epkt->sent_list;
while (sent) {
if (sent->id == epkt->last_sent_id && sent->need_ack) {
uint16_t age = timestamp_diff(current_time, sent->send_time);
uint16_t threshold_new = epkt->rtt_avg_10 * 2;
if (age > threshold_new) {
// Check if already in retransmit queue
int already = 0;
for (int i = 0; i < epkt->pending_retransmit_count; i++) {
if (epkt->pending_retransmit_ids[i] == sent->id) {
already = 1;
break;
}
}
if (!already && epkt->pending_retransmit_count < 32) {
epkt->pending_retransmit_ids[epkt->pending_retransmit_count++] = sent->id;
}
}
break;
}
sent = sent->next;
}
}
// Reschedule check with updated period
if (epkt->retrans_timer_period > 0) {
@ -519,6 +715,9 @@ static void retransmit_check(epkt_t* epkt) {
int etcp_rx_input(epkt_t* epkt, uint8_t* pkt, uint16_t len) {
if (!epkt || !pkt || len < 4) return -1;
// Update received bytes statistics
epkt->bytes_received_total += len;
// Parse header
uint8_t* ptr = pkt;
uint16_t id = (ptr[0] << 8) | ptr[1];
@ -547,89 +746,135 @@ int etcp_rx_input(epkt_t* epkt, uint8_t* pkt, uint16_t len) {
payload_len = len;
break; // Payload is the rest of the packet
} else if (hdr == 0x01) {
// Timestamp report
if (len >= 4) {
uint16_t ack_id = (ptr[0] << 8) | ptr[1];
uint16_t ack_timestamp = (ptr[2] << 8) | ptr[3];
ptr += 4;
len -= 4;
// Calculate RTT
uint16_t current_time = get_current_timestamp();
uint16_t rtt_raw = timestamp_diff(current_time, ack_timestamp);
if (rtt_raw > 0) {
update_metrics(epkt, rtt_raw);
}
// Remove acknowledged packet from sent_list
sent_packet_t* sent = epkt->sent_list;
sent_packet_t* prev = NULL;
while (sent) {
if (sent->id == ack_id) {
if (prev) {
prev->next = sent->next;
} else {
epkt->sent_list = sent->next;
// Timestamp report with count
if (len >= 1) {
uint8_t count = *ptr++;
len--;
if (len >= count * 4 + 4) {
// Process each timestamp pair
for (int i = 0; i < count; i++) {
uint16_t ack_id = (ptr[0] << 8) | ptr[1];
uint16_t ack_timestamp = (ptr[2] << 8) | ptr[3];
ptr += 4;
len -= 4;
// Calculate RTT
uint16_t current_time = get_current_timestamp();
uint16_t rtt_raw = timestamp_diff(current_time, ack_timestamp);
if (rtt_raw > 0) {
update_metrics(epkt, rtt_raw);
}
// Update unacked bytes
epkt->unacked_bytes -= sent->payload_len;
// DEBUG
// printf("ACK received for id=%u, unacked_bytes now=%u, payload_len=%u\n",
// ack_id, epkt->unacked_bytes, sent->payload_len);
// Update last acknowledged ID
if (id_compare(ack_id, epkt->last_acked_id) > 0) {
epkt->last_acked_id = ack_id;
// Remove acknowledged packet from sent_list
sent_packet_t* sent = epkt->sent_list;
sent_packet_t* prev = NULL;
while (sent) {
if (sent->id == ack_id) {
if (prev) {
prev->next = sent->next;
} else {
epkt->sent_list = sent->next;
}
// Update unacked bytes (with underflow protection)
if (sent->payload_len > epkt->unacked_bytes) {
epkt->unacked_bytes = 0;
} else {
epkt->unacked_bytes -= sent->payload_len;
}
ETCP_LOG("ACK received for id=%u, unacked_bytes now=%u, payload_len=%u\n",
ack_id, epkt->unacked_bytes, sent->payload_len);
// Update last acknowledged ID
if (id_compare(ack_id, epkt->last_acked_id) > 0) {
epkt->last_acked_id = ack_id;
}
// Update last received ACK ID
if (id_compare(ack_id, epkt->last_rx_ack_id) > 0) {
epkt->last_rx_ack_id = ack_id;
}
// Window may have opened
epkt->window_blocked = 0;
// Forward progress tracking
epkt->oldest_missing_id = 0;
epkt->missing_since_time = 0;
free(sent->data);
free(sent);
break;
}
prev = sent;
sent = sent->next;
}
// Update last received ACK ID
if (id_compare(ack_id, epkt->last_rx_ack_id) > 0) {
epkt->last_rx_ack_id = ack_id;
}
// Read last delivered ID and last received ID
uint16_t last_delivered = (ptr[0] << 8) | ptr[1];
uint16_t last_received = (ptr[2] << 8) | ptr[3];
ptr += 4;
len -= 4;
// Update our last delivered if newer
if (id_compare(last_delivered, epkt->last_delivered_id) > 0) {
epkt->last_delivered_id = last_delivered;
// Clear oldest missing if it's now delivered or skipped
if (epkt->oldest_missing_id != 0 && id_compare(epkt->oldest_missing_id, epkt->last_delivered_id) <= 0) {
epkt->oldest_missing_id = 0;
epkt->missing_since_time = 0;
}
// Window may have opened
epkt->window_blocked = 0;
free(sent->data);
free(sent);
break;
}
prev = sent;
sent = sent->next;
// Update last_rx_ack_id (latest known received packet)
if (id_compare(last_received, epkt->last_rx_ack_id) > 0) {
epkt->last_rx_ack_id = last_received;
}
}
}
} else if (hdr >= 0x10 && hdr <= 0x2F) {
// Retransmission request
// Retransmission request with two IDs
uint8_t count = (hdr & 0x0F) + 1;
if (len >= count * 2 + 2) {
if (len >= count * 2 + 4) {
ETCP_LOG("Retransmit request received: count=%u, IDs: ", count);
// Read IDs to retransmit
for (int i = 0; i < count; i++) {
uint16_t retransmit_id = (ptr[0] << 8) | ptr[1];
ptr += 2;
len -= 2;
// Add to retransmit queue
if (epkt->pending_retransmit_count < 32) {
epkt->pending_retransmit_ids[epkt->pending_retransmit_count++] = retransmit_id;
}
// Retransmit the requested packet
retransmit_packet(epkt, retransmit_id);
ETCP_LOG("%u ", retransmit_id);
}
ETCP_LOG("\n");
// Read last delivered ID
// Read last delivered ID and last received ID
uint16_t last_delivered = (ptr[0] << 8) | ptr[1];
ptr += 2;
len -= 2;
uint16_t last_received = (ptr[2] << 8) | ptr[3];
ptr += 4;
len -= 4;
// Update our last delivered if newer
ETCP_LOG("Comparing last_delivered: remote=%u, local=%u, cmp=%d\n",
last_delivered, epkt->last_delivered_id, id_compare(last_delivered, epkt->last_delivered_id));
if (id_compare(last_delivered, epkt->last_delivered_id) > 0) {
epkt->last_delivered_id = last_delivered;
ETCP_LOG("Updated last_delivered_id to %u\n", epkt->last_delivered_id);
// Clear oldest missing if it's now delivered or skipped
if (epkt->oldest_missing_id != 0 && id_compare(epkt->oldest_missing_id, epkt->last_delivered_id) <= 0) {
epkt->oldest_missing_id = 0;
epkt->missing_since_time = 0;
}
}
// Also update last_rx_ack_id (latest known received packet)
if (id_compare(last_delivered, epkt->last_rx_ack_id) > 0) {
epkt->last_rx_ack_id = last_delivered;
// Update last_rx_ack_id (latest known received packet)
if (id_compare(last_received, epkt->last_rx_ack_id) > 0) {
epkt->last_rx_ack_id = last_received;
}
ETCP_LOG("Updated from retransmit request: last_delivered=%u, last_received=%u, our_last_delivered=%u\n",
last_delivered, last_received, epkt->last_delivered_id);
}
}
// Unknown hdr - skip?
}
// Add to rx_list if has payload
if (has_payload && payload_len > 0) {
// Add to rx_list if has payload and id != 0 (id=0 is for metrics-only)
if (has_payload && payload_len > 0 && id != 0) {
// Check for duplicate
rx_packet_t* current = epkt->rx_list;
rx_packet_t* prev = NULL;
@ -670,6 +915,13 @@ int etcp_rx_input(epkt_t* epkt, uint8_t* pkt, uint16_t len) {
epkt->rx_list = new_pkt;
}
ETCP_LOG("Received packet id=%u, last_delivered=%u\n", id, epkt->last_delivered_id);
// If this packet was the oldest missing, clear the tracking
if (epkt->oldest_missing_id != 0 && epkt->oldest_missing_id == id) {
epkt->oldest_missing_id = 0;
epkt->missing_since_time = 0;
}
// Add to pending ACKs
if (epkt->pending_ack_count < 32) {
epkt->pending_ack_ids[epkt->pending_ack_count] = id;
@ -683,28 +935,51 @@ int etcp_rx_input(epkt_t* epkt, uint8_t* pkt, uint16_t len) {
// Move continuous sequence to output queue
uint16_t next_expected = epkt->last_delivered_id + 1;
rx_packet_t* rx = epkt->rx_list;
ETCP_LOG("Delivery check: next_expected=%u\n", next_expected);
while (rx && rx->id == next_expected) {
// DEBUG
// printf("Delivering packet id=%u to output queue\n", rx->id);
// Add to output queue
ll_entry_t* entry = queue_entry_new(rx->data_len);
if (entry) {
memcpy(ll_entry_data(entry), rx->data, rx->data_len);
queue_entry_put(epkt->output_queue, entry);
}
// Update last delivered
epkt->last_delivered_id = next_expected;
next_expected++;
// Continue delivering as long as we find the next expected packet
int delivered;
do {
delivered = 0;
rx_packet_t* rx = epkt->rx_list;
rx_packet_t* prev = NULL;
// Remove from rx_list
epkt->rx_list = rx->next;
free(rx->data);
free(rx);
rx = epkt->rx_list;
}
// Search for packet with id == next_expected
while (rx) {
if (rx->id == next_expected) {
ETCP_LOG("Delivering packet id=%u to output queue\n", rx->id);
// Add to output queue
ll_entry_t* entry = queue_entry_new(rx->data_len);
if (entry) {
memcpy(ll_entry_data(entry), rx->data, rx->data_len);
queue_entry_put(epkt->output_queue, entry);
}
// Update last delivered
epkt->last_delivered_id = next_expected;
// Clear oldest missing if it's now delivered or skipped
if (epkt->oldest_missing_id != 0 && id_compare(epkt->oldest_missing_id, epkt->last_delivered_id) <= 0) {
epkt->oldest_missing_id = 0;
epkt->missing_since_time = 0;
}
next_expected++;
// Remove from rx_list
if (prev) {
prev->next = rx->next;
} else {
epkt->rx_list = rx->next;
}
free(rx->data);
free(rx);
delivered = 1;
break;
}
prev = rx;
rx = rx->next;
}
} while (delivered);
}
return 0;
@ -733,6 +1008,25 @@ static void retransmit_timer_callback(void* arg) {
static void request_retransmission_for_gaps(epkt_t* epkt) {
if (!epkt || !epkt->rx_list) return;
// Forward progress: if oldest missing packet has been missing for > 3×RTT, skip it
if (epkt->oldest_missing_id != 0) {
uint16_t current_time = get_current_timestamp();
uint16_t missing_duration = timestamp_diff(current_time, epkt->missing_since_time);
uint16_t threshold = epkt->rtt_avg_10 * 3;
if (threshold < 60) threshold = 60; // Minimum 6ms
if (missing_duration > threshold) {
ETCP_LOG("Forward progress: skipping missing packet id=%u (missing for %u > threshold %u)\n",
epkt->oldest_missing_id, missing_duration, threshold);
// Advance last_delivered_id past the missing packet
epkt->last_delivered_id = epkt->oldest_missing_id;
epkt->oldest_missing_id = 0;
epkt->missing_since_time = 0;
// Continue to deliver any now-contiguous packets
// The function will be called again after this update
}
}
// Find gaps in rx_list
uint16_t expected = epkt->last_delivered_id + 1;
rx_packet_t* current = epkt->rx_list;
@ -744,6 +1038,13 @@ static void request_retransmission_for_gaps(epkt_t* epkt) {
while (id_compare(missing, current->id) < 0) {
if (epkt->pending_retransmit_count < 32) {
epkt->pending_retransmit_ids[epkt->pending_retransmit_count++] = missing;
ETCP_LOG("Gap: requesting retransmit for id=%u (last_delivered=%u)\n",
missing, epkt->last_delivered_id);
// Track oldest missing packet for forward progress
if (epkt->oldest_missing_id == 0 || id_compare(missing, epkt->oldest_missing_id) < 0) {
epkt->oldest_missing_id = missing;
epkt->missing_since_time = get_current_timestamp();
}
}
missing++;
}
@ -763,6 +1064,54 @@ static void schedule_ack_timer(epkt_t* epkt) {
}
}
// Retransmit a specific packet
static void retransmit_packet(epkt_t* epkt, uint16_t id) {
if (!epkt || !epkt->tx_callback) return;
// Find packet in sent_list
sent_packet_t* sent = epkt->sent_list;
while (sent) {
if (sent->id == id) {
// Update statistics
epkt->retransmissions_count++;
epkt->total_packets_sent++;
// Resend the packet (with same data)
epkt->tx_callback(epkt, sent->data, sent->data_len, epkt->tx_callback_arg);
// Update send time for retransmission timeout
sent->send_time = get_current_timestamp();
ETCP_LOG("Retransmitted packet id=%u\n", id);
return;
}
sent = sent->next;
}
ETCP_LOG("Cannot retransmit packet id=%u: not found in sent_list\n", id);
}
// ==================== Statistics API ====================
void etcp_get_stats(epkt_t* epkt,
uint32_t* retransmissions,
uint32_t* total_packets_sent,
uint32_t* unique_packets_sent,
uint32_t* bytes_sent_total,
uint32_t* bytes_received_total,
uint32_t* ack_packets_count,
uint32_t* control_packets_count) {
if (!epkt) return;
if (retransmissions) *retransmissions = epkt->retransmissions_count;
if (total_packets_sent) *total_packets_sent = epkt->total_packets_sent;
if (unique_packets_sent) *unique_packets_sent = epkt->unique_packets_sent;
if (bytes_sent_total) *bytes_sent_total = epkt->bytes_sent_total;
if (bytes_received_total) *bytes_received_total = epkt->bytes_received_total;
if (ack_packets_count) *ack_packets_count = epkt->ack_packets_count;
if (control_packets_count) *control_packets_count = epkt->control_packets_count;
}
// ==================== Queue Callbacks ====================
static void tx_queue_callback(ll_queue_t* q, ll_entry_t* entry, void* arg) {
@ -771,6 +1120,7 @@ static void tx_queue_callback(ll_queue_t* q, ll_entry_t* entry, void* arg) {
epkt_t* epkt = (epkt_t*)arg;
if (!epkt) return;
ETCP_LOG("tx_queue_callback triggered\n");
// Start transmission process
tx_process(epkt);
}

202
etcp.h

@ -1,4 +1,4 @@
// etcp.h - Extended Transmission Control Protocol
// etcp.h - Расширенный протокол управления передачей (Extended Transmission Control Protocol)
#ifndef ETCP_H
#define ETCP_H
@ -6,157 +6,211 @@
#include <stddef.h>
#include "ll_queue.h"
// Отладочное логирование
#ifdef ETCP_DEBUG
#include <stdio.h>
#define ETCP_LOG(fmt, ...) printf("[ETCP] " fmt, ##__VA_ARGS__)
#ifdef ETCP_DEBUG_EXT
#define ETCP_DEBUG_LOG(fmt, ...) printf("[ETCP_DEBUG] " fmt, ##__VA_ARGS__)
#else
#define ETCP_DEBUG_LOG(fmt, ...) ((void)0)
#endif
#else
#define ETCP_LOG(fmt, ...) ((void)0)
#define ETCP_DEBUG_LOG(fmt, ...) ((void)0)
#endif
#ifdef __cplusplus
extern "C" {
#endif
// Forward declarations
// Предварительные объявления
typedef struct epkt epkt_t;
// Callback type for sending packets via UDP
// Тип обратного вызова для отправки пакетов через UDP
typedef void (*etcp_tx_callback_t)(epkt_t* epkt, uint8_t* pkt, uint16_t len, void* arg);
// Main ETCP structure
// Основная структура ETCP
struct epkt {
// Queues
ll_queue_t* tx_queue; // Queue of data to send
ll_queue_t* output_queue; // Output queue (reassembled data)
// Очереди
ll_queue_t* tx_queue; // Очередь данных для отправки
ll_queue_t* output_queue; // Выходная очередь (собранные данные)
// Received packets sorted linked list
// Список полученных пакетов (отсортированный связанный список)
struct rx_packet* rx_list;
// Sent packets (for retransmission)
// Отправленные пакеты (для повторной передачи)
struct sent_packet* sent_list;
// Metrics
uint16_t rtt_last; // Last RTT (timebase 0.1us)
uint16_t rtt_avg_10; // Average RTT last 10 packets
uint16_t rtt_avg_100; // Average RTT last 100 packets
uint16_t jitter; // Jitter (averaged)
uint16_t bandwidth; // Current bandwidth (bytes per timebase)
uint32_t bytes_sent_total; // Total bytes sent
uint16_t last_sent_timestamp; // Timestamp of last sent packet
uint32_t bytes_allowed; // Calculated bytes allowed to send
// Метрики
uint16_t rtt_last; // Последнее RTT (в единицах времени 0.1 мкс)
uint16_t rtt_avg_10; // Среднее RTT за последние 10 пакетов
uint16_t rtt_avg_100; // Среднее RTT за последние 100 пакетов
uint16_t jitter; // Джиттер (усредненный)
uint16_t bandwidth; // Текущая пропускная способность (байты за единицу времени)
uint32_t bytes_sent_total; // Общее количество отправленных байт
uint16_t last_sent_timestamp; // Временная метка последнего отправленного пакета
uint32_t bytes_allowed; // Рассчитанное количество разрешенных к отправке байт
// State
uint16_t next_tx_id; // Next ID for transmission
uint16_t last_rx_id; // Last received ID (for ACK)
uint16_t last_delivered_id; // Last delivered to output_queue ID
// Статистика
uint32_t retransmissions_count; // Количество ретрансмиссий
uint32_t ack_packets_count; // Количество отправленных пакетов подтверждения
uint32_t control_packets_count; // Количество отправленных управляющих пакетов (ACK + запросы ретрансмиссии)
uint32_t total_packets_sent; // Общее количество отправленных пакетов (включая ретрансмиссии)
uint32_t unique_packets_sent; // Количество уникальных отправленных пакетов (без ретрансмиссий)
uint32_t bytes_received_total; // Общее количество полученных байт
// Timers
void* next_tx_timer; // Timer for next transmission
void* retransmit_timer; // Timer for retransmissions
// Состояние
uint16_t next_tx_id; // Следующий ID для передачи
uint16_t last_sent_id; // Последний отправленный ID (для ретрансмиссии самого нового пакета)
uint16_t last_rx_id; // Последний полученный ID (для подтверждения)
uint16_t last_delivered_id; // Последний ID, переданный в output_queue
// Callback
// Таймеры
void* next_tx_timer; // Таймер для следующей передачи
void* retransmit_timer; // Таймер для повторных передач
// Обратный вызов
etcp_tx_callback_t tx_callback;
void* tx_callback_arg;
// RTT history for averaging
// История RTT для усреднения
uint16_t rtt_history[100];
uint8_t rtt_history_idx;
uint8_t rtt_history_count;
// Pending ACKs
// Ожидающие подтверждения
uint16_t pending_ack_ids[32];
uint16_t pending_ack_timestamps[32];
uint8_t pending_ack_count;
// Pending retransmission requests
// Ожидающие запросы на повторную передачу
uint16_t pending_retransmit_ids[32];
uint8_t pending_retransmit_count;
// Window management
uint32_t unacked_bytes; // Number of bytes sent but not yet acknowledged
uint32_t window_size; // Current window size in bytes (calculated)
uint16_t last_acked_id; // Last acknowledged packet ID
uint16_t last_rx_ack_id; // Latest received ACK ID from receiver
uint16_t retrans_timer_period; // Current retransmission timer period (timebase)
uint16_t next_retrans_time; // Time of next retransmission check
uint8_t window_blocked; // Flag: transmission blocked by window limit
// Управление окном
uint32_t unacked_bytes; // Количество байт, отправленных но еще не подтвержденных
uint32_t window_size; // Текущий размер окна в байтах (рассчитывается)
uint16_t last_acked_id; // Последний подтвержденный ID пакета
uint16_t last_rx_ack_id; // Последний полученный ID подтверждения от получателя
uint16_t retrans_timer_period; // Текущий период таймера повторной передачи (в единицах времени)
uint16_t next_retrans_time; // Время следующей проверки повторной передачи
uint8_t window_blocked; // Флаг: передача заблокирована из-за ограничения окна
// Forward progress tracking
uint16_t oldest_missing_id; // Oldest missing packet ID
uint16_t missing_since_time; // Time when oldest missing packet was first detected
};
// API Functions
// Функции API
/**
* @brief Initialize new ETCP instance
* @return Pointer to new instance or NULL on error
* @brief Инициализировать новый экземпляр ETCP
* @return Указатель на новый экземпляр или NULL в случае ошибки
*/
epkt_t* etcp_init(void);
/**
* @brief Free ETCP instance and all associated resources
* @param epkt Instance to free
* @brief Освободить экземпляр ETCP и все связанные ресурсы
* @param epkt Экземпляр для освобождения
*/
void etcp_free(epkt_t* epkt);
/**
* @brief Set callback for sending packets via UDP
* @param epkt ETCP instance
* @param cb Callback function
* @param arg User argument passed to callback
* @brief Установить обратный вызов для отправки пакетов через UDP
* @param epkt Экземпляр ETCP
* @param cb Функция обратного вызова
* @param arg Пользовательский аргумент, передаваемый в обратный вызов
*/
void etcp_set_callback(epkt_t* epkt, etcp_tx_callback_t cb, void* arg);
/**
* @brief Process received UDP packet
* @param epkt ETCP instance
* @param pkt Packet data
* @param len Packet length
* @return 0 on success, -1 on error
* @brief Обработать полученный UDP пакет
* @param epkt Экземпляр ETCP
* @param pkt Данные пакета
* @param len Длина пакета
* @return 0 при успехе, -1 при ошибке
*/
int etcp_rx_input(epkt_t* epkt, uint8_t* pkt, uint16_t len);
/**
* @brief Get total number of packets waiting in transmission queues
* @param epkt ETCP instance
* @return Number of packets
* @brief Получить общее количество пакетов, ожидающих в очередях передачи
* @param epkt Экземпляр ETCP
* @return Количество пакетов
*/
int etcp_tx_queue_size(epkt_t* epkt);
/**
* @brief Put data into transmission queue
* @param epkt ETCP instance
* @param data Data to send
* @param len Data length
* @return 0 on success, -1 on error
* @brief Поместить данные в очередь передачи
* @param epkt Экземпляр ETCP
* @param data Данные для отправки
* @param len Длина данных
* @return 0 при успехе, -1 при ошибке
*/
int etcp_tx_put(epkt_t* epkt, uint8_t* data, uint16_t len);
/**
* @brief Get output queue for reading received data
* @param epkt ETCP instance
* @return Pointer to output queue (ll_queue_t*)
* @brief Получить выходную очередь для чтения полученных данных
* @param epkt Экземпляр ETCP
* @return Указатель на выходную очередь (ll_queue_t*)
*/
ll_queue_t* etcp_get_output_queue(epkt_t* epkt);
/**
* @brief Set bandwidth limit
* @param epkt ETCP instance
* @param bandwidth Bytes per timebase (0.1us)
* @brief Установить ограничение пропускной способности
* @param epkt Экземпляр ETCP
* @param bandwidth Байты за единицу времени (0.1 мкс)
*/
void etcp_set_bandwidth(epkt_t* epkt, uint16_t bandwidth);
/**
* @brief Update window size based on current RTT and bandwidth
* @param epkt ETCP instance
* Window size = RTT * bandwidth * 2 (bytes in flight)
* @brief Обновить размер окна на основе текущего RTT и пропускной способности
* @param epkt Экземпляр ETCP
* Размер окна = RTT * пропускная способность * 2 (байт в пути)
*/
void etcp_update_window(epkt_t* epkt);
/**
* @brief Get current RTT
* @param epkt ETCP instance
* @return RTT in timebase units
* @brief Получить текущее RTT
* @param epkt Экземпляр ETCP
* @return RTT в единицах времени
*/
uint16_t etcp_get_rtt(epkt_t* epkt);
/**
* @brief Get current jitter
* @param epkt ETCP instance
* @return Jitter in timebase units
* @brief Получить текущий джиттер
* @param epkt Экземпляр ETCP
* @return Джиттер в единицах времени
*/
uint16_t etcp_get_jitter(epkt_t* epkt);
/**
* @brief Сбросить состояние соединения (очистить очереди, метрики, таймеры)
* @param epkt Экземпляр ETCP
* Примечание: Сохраняет настройки пропускной способности и обратного вызова
*/
void etcp_reset(epkt_t* epkt);
/**
* @brief Получить статистику ETCP
* @param epkt Экземпляр ETCP
* @param retransmissions Указатель для возврата количества ретрансмиссий
* @param total_packets_sent Указатель для возврата общего количества отправленных пакетов
* @param unique_packets_sent Указатель для возврата количества уникальных отправленных пакетов
* @param bytes_sent_total Указатель для возврата общего количества отправленных байт
* @param bytes_received_total Указатель для возврата общего количества полученных байт
* @param ack_packets_count Указатель для возврата количества отправленных пакетов подтверждения
* @param control_packets_count Указатель для возврата количества отправленных управляющих пакетов
*/
void etcp_get_stats(epkt_t* epkt,
uint32_t* retransmissions,
uint32_t* total_packets_sent,
uint32_t* unique_packets_sent,
uint32_t* bytes_sent_total,
uint32_t* bytes_received_total,
uint32_t* ack_packets_count,
uint32_t* control_packets_count);
#ifdef __cplusplus
}
#endif

179
etcp.h1

@ -0,0 +1,179 @@
// etcp.h - Extended Transmission Control Protocol
#ifndef ETCP_H
#define ETCP_H
#include <stdint.h>
#include <stddef.h>
#include "ll_queue.h"
// Debug logging
#ifdef ETCP_DEBUG
#include <stdio.h>
#define ETCP_LOG(fmt, ...) printf("[ETCP] " fmt, ##__VA_ARGS__)
#else
#define ETCP_LOG(fmt, ...) ((void)0)
#endif
#ifdef __cplusplus
extern "C" {
#endif
// Forward declarations
typedef struct epkt epkt_t;
// Callback type for sending packets via UDP
typedef void (*etcp_tx_callback_t)(epkt_t* epkt, uint8_t* pkt, uint16_t len, void* arg);
// Main ETCP structure
struct epkt {
// Queues
ll_queue_t* tx_queue; // Queue of data to send
ll_queue_t* output_queue; // Output queue (reassembled data)
// Received packets sorted linked list
struct rx_packet* rx_list;
// Sent packets (for retransmission)
struct sent_packet* sent_list;
// Metrics
uint16_t rtt_last; // Last RTT (timebase 0.1us)
uint16_t rtt_avg_10; // Average RTT last 10 packets
uint16_t rtt_avg_100; // Average RTT last 100 packets
uint16_t jitter; // Jitter (averaged)
uint16_t bandwidth; // Current bandwidth (bytes per timebase)
uint32_t bytes_sent_total; // Total bytes sent
uint16_t last_sent_timestamp; // Timestamp of last sent packet
uint32_t bytes_allowed; // Calculated bytes allowed to send
// State
uint16_t next_tx_id; // Next ID for transmission
uint16_t last_rx_id; // Last received ID (for ACK)
uint16_t last_delivered_id; // Last delivered to output_queue ID
// Timers
void* next_tx_timer; // Timer for next transmission
void* retransmit_timer; // Timer for retransmissions
// Callback
etcp_tx_callback_t tx_callback;
void* tx_callback_arg;
// RTT history for averaging
uint16_t rtt_history[100];
uint8_t rtt_history_idx;
uint8_t rtt_history_count;
// Pending ACKs
uint16_t pending_ack_ids[32];
uint16_t pending_ack_timestamps[32];
uint8_t pending_ack_count;
// Pending retransmission requests
uint16_t pending_retransmit_ids[32];
uint8_t pending_retransmit_count;
// Window management
uint32_t unacked_bytes; // Number of bytes sent but not yet acknowledged
uint32_t window_size; // Current window size in bytes (calculated)
uint16_t last_acked_id; // Last acknowledged packet ID
uint16_t last_rx_ack_id; // Latest received ACK ID from receiver
uint16_t retrans_timer_period; // Current retransmission timer period (timebase)
uint16_t next_retrans_time; // Time of next retransmission check
uint8_t window_blocked; // Flag: transmission blocked by window limit
};
// API Functions
/**
* @brief Initialize new ETCP instance
* @return Pointer to new instance or NULL on error
*/
epkt_t* etcp_init(void);
/**
* @brief Free ETCP instance and all associated resources
* @param epkt Instance to free
*/
void etcp_free(epkt_t* epkt);
/**
* @brief Set callback for sending packets via UDP
* @param epkt ETCP instance
* @param cb Callback function
* @param arg User argument passed to callback
*/
void etcp_set_callback(epkt_t* epkt, etcp_tx_callback_t cb, void* arg);
/**
* @brief Process received UDP packet
* @param epkt ETCP instance
* @param pkt Packet data
* @param len Packet length
* @return 0 on success, -1 on error
*/
int etcp_rx_input(epkt_t* epkt, uint8_t* pkt, uint16_t len);
/**
* @brief Get total number of packets waiting in transmission queues
* @param epkt ETCP instance
* @return Number of packets
*/
int etcp_tx_queue_size(epkt_t* epkt);
/**
* @brief Put data into transmission queue
* @param epkt ETCP instance
* @param data Data to send
* @param len Data length
* @return 0 on success, -1 on error
*/
int etcp_tx_put(epkt_t* epkt, uint8_t* data, uint16_t len);
/**
* @brief Get output queue for reading received data
* @param epkt ETCP instance
* @return Pointer to output queue (ll_queue_t*)
*/
ll_queue_t* etcp_get_output_queue(epkt_t* epkt);
/**
* @brief Set bandwidth limit
* @param epkt ETCP instance
* @param bandwidth Bytes per timebase (0.1us)
*/
void etcp_set_bandwidth(epkt_t* epkt, uint16_t bandwidth);
/**
* @brief Update window size based on current RTT and bandwidth
* @param epkt ETCP instance
* Window size = RTT * bandwidth * 2 (bytes in flight)
*/
void etcp_update_window(epkt_t* epkt);
/**
* @brief Get current RTT
* @param epkt ETCP instance
* @return RTT in timebase units
*/
uint16_t etcp_get_rtt(epkt_t* epkt);
/**
* @brief Get current jitter
* @param epkt ETCP instance
* @return Jitter in timebase units
*/
uint16_t etcp_get_jitter(epkt_t* epkt);
/**
* @brief Reset connection state (clear queues, metrics, timers)
* @param epkt ETCP instance
* Note: Keeps bandwidth setting and callback
*/
void etcp_reset(epkt_t* epkt);
#ifdef __cplusplus
}
#endif
#endif // ETCP_H

337
etcp.txt

@ -1,55 +1,282 @@
etcp - extended transmission control protocol
Протокол для передачи-приёма, пободный TCP, реализованый отдельным модулем (etcp.c/h).
Задача протокола:
- передать пакеты через UDP (учитывая его особенности - потери, негарантированный порядок), восстанавливая порядок и потери.
- пакеты уже предварительно подогнаны под размер чтобы вмещались доп. заголовки и служебные фреймы.
-
На приёмной стороне создаются две очереди.
первая - сортированный linked-list - в нее добавляются принятые пакеты, но отсеиваются дубликаты и вставляются в нужное место.
и перемещаются в выходную очередь когда все нужные пакеты дошли.
при получении принятого пакета он сразу парсится, для всех hdr!=0 вызываем static upd_metric_for_transmitter(epkt*, buf*, size) - и он разгребает метрики, обновляя:
- rtt_last (roud-trip delay) для последнего пакета
- rtt average last 10
- rtt average last 100
- jitter как усредненное: jitter+=(abs(rtt_last_10-rtt_last)-jitter)*0.1f
- для retransmission request - если timestamp последней попытки передачи этого пакета больше rtt_last_10*1.2+jitter*2 то отправляем сейчас
при отправке также ограничиваем полосу пропускания: суммируем сколько байт отправлено всего (uint32_t), обновляем timestamp и расчетное число байт которое может быть отправлено на момент этого timestamp. и формируем таймер для отправки следующего пакета.
вторая - ll_queue - выходная осчередь с пакетами в строгом порядке (строгий инкремент по id).
struct epkt* = etcp_init() - инициализирует новый instance и выделяет под него память
etcp_free(struct epkt*)
etcp_rx_input(struct epkt*,uint8_t* pkt, uint16_t len) - принятый пакет на обработку
rx_output - через механизм ll_queue, функция не нужна.
новый пакет на передачу отправляется в очередь передачи ll_queue. используй callback для обработки очереди.
для передачи сформированных пакетов по udp:
etcp_set_callback(epkt*, &cbk) -> etcp_tx_output(struct epkt*,uint8_t* pkt, uint16_t len) - callback в управляющей структуре (отправка пакета в udp сокет)
внутренняя структура передачи:
при готовности отправить очередной пакет из очереди ll_queue пакеты перемещаются в linked_list (как отправленные но неподтвержденные), и освобождаются при получении подтверждения (ack).
int etcp_tx_queue_size(epkt*) - должна быть функция которая возвращает общее кол-во пакетов на передачу в очередях (входящей и рабочей)
Формат udp пакета:
<id> <timestamp> [<hdr> metrics] <hdr=0> <payload>
id - uint16_t циклический порядковый номер пакета (при передаче следующего пакета инкрементируется). при ретрансмиссии передается с этим же id и payload, но с обновленными остальными полями
timestamp - uint16_t текущий timestamp (циклическое, 16 бит, timebase = 0.1mS)
hdr - 1 байт:
payload - передаваемые полезные данные
metrics - опциональное поле для передачи служебных фреймов (ack, retransmission request, statistic reply)
hdr:
0x00 - hdr для payload (от следующего байта до конца пакета)
0x01 - hdr для отчета о timestamp (время приёма) принятого пакета, 4 байта: <id> <timestamp> - передаётся при очередной передаче пакета, для формирования статичтики на приёмной стороне. отчеты передаются для всех новых принятых пакетов с момента последней передачи. т.е. накапливаем timestamp-ы и передаём их. если пакет потерялся - не страшно.
0x10-0x2f - hdr для перезапроса пакетов (передачу каких пакетов надо повторить. значение определяет количество записей (номеров пакетов) от 1 до 32, если больше - 32 самых старых), далее по 2 байта идут ID пакетов. и в конце - 2 байта номер последнего пакета который ушел в выходную очередь (т.е. последний номер для успешно собранной цепочки)
если что-то еще надо можно добавить.
Если данных нет (очередь на передачу пустая) и нужно передать только метрику, то передаётся пакет с id=0 и без <hdr=0> <payload>. на приёмной стороне он определяется по отсутствию записи с hdr=0
ETCP - Extended Transmission Control Protocol
=============================================
Протокол для надежной передачи данных поверх UDP с восстановлением порядка,
повторной передачей потерянных пакетов и управлением потоком.
Основные задачи:
- Передача пакетов через UDP с учетом его особенностей (потери, негарантированный порядок)
- Восстановление правильного порядка пакетов на приемной стороне
- Повторная передача потерянных пакетов
- Управление потоком на основе измерения RTT и пропускной способности
- Обеспечение прогресса доставки даже при длительных потерях
Архитектура:
------------
На передающей стороне:
1. Очередь передачи (tx_queue) - данные, ожидающие отправки
2. Список отправленных пакетов (sent_list) - пакеты, ожидающие подтверждения
3. Таймеры для управления передачей и повторной отправкой
На приемной стороне:
1. Сортированный список принятых пакетов (rx_list) - пакеты в порядке ID
2. Выходная очередь (output_queue) - собранные в правильном порядке данные
3. Буферы для накопления подтверждений и запросов на повторную передачу
Структура данных:
-----------------
struct epkt {
// Очереди
ll_queue_t* tx_queue; // Очередь данных для отправки
ll_queue_t* output_queue; // Выходная очередь (собранные данные)
// Списки пакетов
struct rx_packet* rx_list; // Полученные пакеты (отсортированный список)
struct sent_packet* sent_list; // Отправленные пакеты (для повторной передачи)
// Метрики
uint16_t rtt_last; // Последнее RTT (в единицах времени 0.1 мс)
uint16_t rtt_avg_10; // Среднее RTT за последние 10 пакетов
uint16_t rtt_avg_100; // Среднее RTT за последние 100 пакетов
uint16_t jitter; // Джиттер (усредненный)
uint16_t bandwidth; // Текущая пропускная способность (байты за единицу времени)
uint32_t bytes_sent_total; // Общее количество отправленных байт
uint16_t last_sent_timestamp; // Временная метка последнего отправленного пакета
uint32_t bytes_allowed; // Рассчитанное количество разрешенных к отправке байт
// Состояние передачи
uint16_t next_tx_id; // Следующий ID для передачи
uint16_t last_sent_id; // Последний отправленный ID (для ретрансмиссии самого нового пакета)
uint16_t last_rx_id; // Последний полученный ID (для подтверждения)
uint16_t last_delivered_id; // Последний ID, переданный в output_queue
// Таймеры
void* next_tx_timer; // Таймер для следующей передачи
void* retransmit_timer; // Таймер для повторных передач
// Обратный вызов для отправки
etcp_tx_callback_t tx_callback;
void* tx_callback_arg;
// История RTT для усреднения
uint16_t rtt_history[100];
uint8_t rtt_history_idx;
uint8_t rtt_history_count;
// Ожидающие подтверждения
uint16_t pending_ack_ids[32];
uint16_t pending_ack_timestamps[32];
uint8_t pending_ack_count;
// Ожидающие запросы на повторную передачу
uint16_t pending_retransmit_ids[32];
uint8_t pending_retransmit_count;
// Управление окном
uint32_t unacked_bytes; // Количество байт, отправленных но еще не подтвержденных
uint32_t window_size; // Текущий размер окна в байтах (рассчитывается)
uint16_t last_acked_id; // Последний подтвержденный ID пакета
uint16_t last_rx_ack_id; // Последний полученный ID подтверждения от получателя
uint16_t retrans_timer_period; // Текущий период таймера повторной передачи (в единицах времени)
uint16_t next_retrans_time; // Время следующей проверки повторной передачи
uint8_t window_blocked; // Флаг: передача заблокирована из-за ограничения окна
// Отслеживание прогресса доставки
uint16_t oldest_missing_id; // Самый старый отсутствующий ID пакета
uint16_t missing_since_time; // Время, когда самый старый отсутствующий пакет был впервые обнаружен
};
Формат пакета:
---------------
Пакет состоит из обязательного заголовка и опциональных секций:
1. Обязательный заголовок (4 байта):
- ID пакета (2 байта): циклический порядковый номер (0 для пакетов только с метриками)
- Timestamp (2 байта): время отправки в единицах 0.1 мс (циклическое, 16 бит)
2. Опциональные секции (одна или несколько, каждая начинается с байта-заголовка):
а) Подтверждения (ACK) - заголовок 0x01:
[0x01] [count] [(id, timestamp) × count] [last_delivered_id] [last_rx_id]
- count: количество пар ID+timestamp (1 байт)
- Для каждого подтверждаемого пакета: ID (2 байта) + timestamp получения (2 байта)
- last_delivered_id (2 байта): последний ID, доставленный в выходную очередь
- last_rx_id (2 байта): последний полученный ID (новейший известный пакет)
б) Запросы на повторную передачу - заголовок 0x10-0x2F:
[0x10 + (count-1)] [IDs × count] [last_delivered_id] [last_rx_id]
- count: (заголовок & 0x0F) + 1 (от 1 до 32)
- Для каждого запрашиваемого пакета: ID (2 байта)
- last_delivered_id (2 байта): последний ID, доставленный в выходную очередь
- last_rx_id (2 байта): последний полученный ID
в) Полезные данные - заголовок 0x00:
[0x00] [данные...]
- Данные произвольной длины (до конца пакета)
Пакет может содержать несколько секций (например, ACK + данные). Секция с полезными данными
обычно идет последней, если присутствует.
Алгоритмы:
----------
1. Управление передачей:
- Передача происходит при наличии данных в tx_queue и доступной полосы пропускания
- Полоса пропускания контролируется через bytes_allowed, который накапливается со временем
- Размер окна рассчитывается как: window_size = RTT × bandwidth × 2
- Передача блокируется, если unacked_bytes превышает window_size
2. Подтверждение и RTT измерение:
- При получении пакета его ID и timestamp добавляются в pending_ack_ids
- При следующей отправке эти подтверждения включаются в пакет
- Получатель вычисляет RTT как разницу между текущим временем и полученным timestamp
- RTT усредняется за последние 10 и 100 пакетов
3. Повторная передача:
- Отправленные пакеты хранятся в sent_list до подтверждения
- Таймер retransmit_timer периодически проверяет пакеты старше 1.5×RTT
- Если пакет не подтвержден, его ID добавляется в pending_retransmit_ids
- Новейший неподтвержденный пакет (last_sent_id) повторно передается после 2×RTT
- Запросы на повторную передачу от получателя также обрабатываются
4. Сборка пакетов на приемной стороне:
- Полученные пакеты вставляются в отсортированный rx_list
- При обнаружении пропусков (gaps) отправляются запросы на повторную передачу
- Непрерывная последовательность пакетов перемещается в output_queue
- last_delivered_id отслеживает последний доставленный ID
5. Обеспечение прогресса доставки (forward progress):
- Отслеживается самый старый отсутствующий пакет (oldest_missing_id)
- Если пакет отсутствует дольше 3×RTT (минимум 6 мс), last_delivered_id продвигается вперед
- Это предотвращает бесконечное ожидание потерянных пакетов
6. Синхронизация состояния:
- Поля last_delivered_id и last_rx_id передаются в ACK и запросах на повторную передачу
- Получатель обновляет свой last_delivered_id, если полученное значение новее
- Это позволяет синхронизировать прогресс доставки между отправителем и получателем
Таймеры:
--------
1. Таймер передачи (next_tx_timer):
- Срабатывает, когда передача невозможна (нет полосы или окно заполнено)
- Перезапускает процесс передачи
2. Таймер повторной передачи (retransmit_timer):
- Период: max(RTT/2, 2 мс)
- Проверяет sent_list на наличие неподтвержденных пакетов старше 1.5×RTT
- Планирует повторную передачу
API функции:
------------
epkt_t* etcp_init(void);
Инициализирует новый экземпляр ETCP
void etcp_free(epkt_t* epkt);
Освобождает экземпляр ETCP и все связанные ресурсы
void etcp_set_callback(epkt_t* epkt, etcp_tx_callback_t cb, void* arg);
Устанавливает обратный вызов для отправки пакетов через UDP
int etcp_rx_input(epkt_t* epkt, uint8_t* pkt, uint16_t len);
Обрабатывает полученный UDP пакет
int etcp_tx_put(epkt_t* epkt, uint8_t* data, uint16_t len);
Помещает данные в очередь передачи
ll_queue_t* etcp_get_output_queue(epkt_t* epkt);
Возвращает выходную очередь для чтения полученных данных
void etcp_set_bandwidth(epkt_t* epkt, uint16_t bandwidth);
Устанавливает ограничение пропускной способности
int etcp_tx_queue_size(epkt_t* epkt);
Возвращает общее количество пакетов, ожидающих в очередях передачи
void etcp_reset(epkt_t* epkt);
Сбрасывает состояние соединения (очищает очереди, метрики, таймеры)
uint16_t etcp_get_rtt(epkt_t* epkt);
Возвращает текущее RTT
uint16_t etcp_get_jitter(epkt_t* epkt);
Возвращает текущий джиттер
Внутренние структуры:
---------------------
typedef struct rx_packet {
struct rx_packet* next;
uint16_t id;
uint16_t timestamp;
uint8_t* data;
uint16_t data_len;
uint8_t has_payload;
} rx_packet_t;
typedef struct sent_packet {
struct sent_packet* next;
uint16_t id;
uint16_t timestamp;
uint8_t* data;
uint16_t data_len; // Общая длина пакета
uint16_t payload_len; // Длина полезных данных (для учета окна)
uint16_t send_time; // Время отправки
uint8_t need_ack; // Требуется подтверждение
uint8_t need_retransmit; // Требуется повторная передача
} sent_packet_t;
Особенности реализации:
-----------------------
1. Циклические счетчики:
- ID пакетов: 16-битные, циклические (0-65535, затем 0)
- Timestamp: 16-битные, циклические (0-65535 единиц по 0.1 мс ≈ 6.55 секунд)
- Сравнение с учетом цикличности через функцию id_compare()
2. Единицы времени:
- Базовый интервал: 0.1 мс (100 микросекунд)
- Все таймеры и измерения RTT используют эту единицу
3. Ограничения:
- Максимум 32 ожидающих подтверждения или запроса на повторную передачу
- История RTT хранит до 100 измерений
- Максимальный размер окна: 2^32-1 байт
4. Обработка дубликатов:
- При получении пакета с уже существующим ID он игнорируется
- Повторная передача пакета имеет тот же ID, но новый timestamp
Пример использования:
---------------------
1. Инициализация:
epkt_t* epkt = etcp_init();
etcp_set_callback(epkt, udp_send_callback, udp_socket);
2. Отправка данных:
etcp_tx_put(epkt, data, len);
3. Обработка входящих пакетов:
etcp_rx_input(epkt, pkt, len);
4. Чтение полученных данных:
ll_queue_t* output = etcp_get_output_queue(epkt);
ll_entry_t* entry = queue_entry_get(output);
// Обработка entry...
5. Освобождение:
etcp_free(epkt);
Примечания:
-----------
- Протокол предназначен для работы в условиях умеренных потерь и реордеринга
- Механизм forward progress обеспечивает доставку даже при потере начальных пакетов
- Управление окном предотвращает перегрузку сети
- Поддержка двунаправленной связи (каждая сторона может быть одновременно отправителем и получателем)
- Минимальные накладные расходы: 4 байта на пакет + опциональные метаданные

109
ll_queue.c

@ -1,11 +1,39 @@
#include "ll_queue.h"
#include <stdlib.h>
#include <string.h>
#include <stdio.h>
#include <assert.h>
#include "u_async.h"
// Предварительное объявление для отложенного возобновления
static void queue_resume_timeout_cb(void* arg);
// Проверить и запустить ожидающие коллбэки
static void check_waiters(ll_queue_t* q) {
if (!q || !q->waiters) return;
queue_waiter_t** pprev = &q->waiters;
queue_waiter_t* waiter = q->waiters;
while (waiter) {
queue_waiter_t* next = waiter->next;
// Проверить условие: не больше max_packets и не больше max_bytes
if (q->count <= waiter->max_packets && q->total_bytes <= waiter->max_bytes) {
// Условие выполнено - вызвать коллбэк
waiter->callback(q, waiter->callback_arg);
// Удалить waiter из списка
*pprev = next;
free(waiter);
// pprev уже указывает на правильный следующий элемент
} else {
// Условие не выполнено - оставить в списке
pprev = &waiter->next;
}
waiter = next;
}
}
// ==================== Управление очередью ====================
ll_queue_t* queue_new(void) {
@ -15,11 +43,13 @@ ll_queue_t* queue_new(void) {
q->head = NULL;
q->tail = NULL;
q->count = 0;
q->total_bytes = 0;
q->size_limit = -1; // По умолчанию без ограничения
q->callback = NULL;
q->callback_arg = NULL;
q->callback_suspended = 0; // Коллбэки разрешены изначально
q->resume_timeout_id = NULL;
q->waiters = NULL;
return q;
}
@ -35,6 +65,14 @@ void queue_free(ll_queue_t* q) {
entry = next;
}
// Освободить все ожидающие коллбэки
queue_waiter_t* waiter = q->waiters;
while (waiter) {
queue_waiter_t* next = waiter->next;
free(waiter);
waiter = next;
}
// Отменить отложенное возобновление если запланировано
if (q->resume_timeout_id) {
uasync_cancel_timeout(q->resume_timeout_id);
@ -123,13 +161,18 @@ int queue_entry_put(ll_queue_t* q, ll_entry_t* entry) {
}
q->tail = entry;
q->count++;
q->total_bytes += entry->size;
// Если очередь была пустой до добавления (count был 0) и коллбэки разрешены - вызвать коллбэк
// Если коллбэки разрешены - вызвать коллбэк
// Это запускает автоматическую обработку очереди
if (q->count == 1 && !q->callback_suspended && q->callback) {
if (!q->callback_suspended && q->callback) {
printf("[LL_QUEUE DEBUG] queue_entry_put: calling callback, count=%d, suspended=%d\n", q->count, q->callback_suspended);
q->callback(q, entry, q->callback_arg);
}
// Проверить ожидающие коллбэки
check_waiters(q);
return 0;
}
@ -149,12 +192,17 @@ int queue_entry_put_first(ll_queue_t* q, ll_entry_t* entry) {
q->tail = entry;
}
q->count++;
q->total_bytes += entry->size;
// Если очередь была пустой до добавления (count был 0) и коллбэки разрешены - вызвать коллбэк
if (q->count == 1 && !q->callback_suspended && q->callback) {
// Если коллбэки разрешены - вызвать коллбэк
if (!q->callback_suspended && q->callback) {
printf("[LL_QUEUE DEBUG] queue_entry_put_first: calling callback, count=%d, suspended=%d\n", q->count, q->callback_suspended);
q->callback(q, entry, q->callback_arg);
}
// Проверить ожидающие коллбэки
check_waiters(q);
return 0;
}
@ -167,6 +215,7 @@ ll_entry_t* queue_entry_get(ll_queue_t* q) {
q->tail = NULL;
}
q->count--;
q->total_bytes -= entry->size;
entry->next = NULL; // Отсоединить от очереди
@ -174,6 +223,9 @@ ll_entry_t* queue_entry_get(ll_queue_t* q) {
// Это предотвращает рекурсию если во время обработки добавляются новые элементы
q->callback_suspended = 1;
// Проверить ожидающие коллбэки
check_waiters(q);
return entry;
}
@ -181,3 +233,52 @@ int queue_entry_count(ll_queue_t* q) {
if (!q) return 0;
return q->count;
}
// ==================== Асинхронное ожидание ====================
queue_waiter_t* queue_wait_threshold(ll_queue_t* q, int max_packets, size_t max_bytes,
queue_threshold_callback_t callback, void* arg) {
if (!q || !callback) return NULL;
// Создать новый waiter
queue_waiter_t* waiter = malloc(sizeof(queue_waiter_t));
if (!waiter) return NULL;
waiter->max_packets = max_packets;
waiter->max_bytes = max_bytes;
waiter->callback = callback;
waiter->callback_arg = arg;
waiter->next = NULL;
// Проверить условие немедленно
if (q->count <= max_packets && q->total_bytes <= max_bytes) {
// Условие уже выполнено - вызвать коллбэк и освободить waiter
callback(q, arg);
free(waiter);
return NULL;
}
// Добавить в список ожидающих
waiter->next = q->waiters;
q->waiters = waiter;
return waiter;
}
void queue_cancel_wait(ll_queue_t* q, queue_waiter_t* waiter) {
if (!q || !waiter) return;
// Найти и удалить waiter из списка
queue_waiter_t** pprev = &q->waiters;
queue_waiter_t* w = q->waiters;
while (w) {
if (w == waiter) {
*pprev = w->next;
free(w);
return;
}
pprev = &w->next;
w = w->next;
}
}

38
ll_queue.h

@ -17,11 +17,24 @@ struct ll_entry {
size_t size; // Размер данных элемента (байт)
};
// Структура условия ожидания (waiter)
struct queue_waiter {
int max_packets; // Максимальное количество пакетов
size_t max_bytes; // Максимальное количество байт
void (*callback)(ll_queue_t* q, void* arg); // Коллбэк для вызова
void* callback_arg; // Аргумент коллбэка
struct queue_waiter* next; // Следующий ожидающий в списке
};
typedef struct queue_waiter queue_waiter_t;
typedef void (*queue_threshold_callback_t)(ll_queue_t* q, void* arg);
// Структура очереди
struct ll_queue {
ll_entry_t* head; // Первый элемент (извлекается отсюда)
ll_entry_t* tail; // Последний элемент (добавляется сюда)
int count; // Текущее количество элементов
size_t total_bytes; // Общий размер данных всех элементов (байт)
int size_limit; // Максимальное количество (-1 = без ограничения)
queue_callback_t callback; // Функция коллбэка
@ -29,6 +42,8 @@ struct ll_queue {
int callback_suspended; // 1 если коллбэки приостановлены (во время обработки)
void* resume_timeout_id; // ID таймаута uasync для отложенного возобновления
queue_waiter_t* waiters; // Список ожидающих коллбэков
};
// ==================== Управление очередью ====================
@ -45,9 +60,11 @@ void queue_free(ll_queue_t* q);
// Установить функцию и аргумент коллбэка для очереди
// Коллбэк вызывается при добавлении элемента в пустую очередь (разрешенные коллбэки)
// обработчик должен обработать этот пакет и когда будет готов к приёму следующего - вызывает resume_callback. обработка строго по одному пакету.
void queue_set_callback(ll_queue_t* q, queue_callback_t cbk_fn, void* arg);
// Возобновить коллбэки после обработки элемента
// Возобновить коллбэки после обработки элемента переданного в коллбэке (тянуть дополнительные элементы из очереди не предусмотернные api нельзя).
// эта функция должна вызываться всегда после того как cbk_fn обработала пакет (можно с ожиданием через async), иначе очередь застрянет.
// Если в очереди остались элементы, запланирует вызов коллбэка через uasync_set_timeout(0)
// Это предотвращает накопление рекурсии в стеке вызовов
void queue_resume_callback(ll_queue_t* q);
@ -99,4 +116,23 @@ static inline size_t ll_entry_size(ll_entry_t* entry) {
return entry->size;
}
// ==================== Асинхронное ожидание ====================
// Зарегистрировать коллбэк, который будет вызван когда очередь будет иметь
// не более max_packets пакетов и не более max_bytes байт.
// Если условие уже выполнено, коллбэк вызывается немедленно.
// Можно зарегистрировать несколько ожиданий на одной очереди.
// Возвращает указатель на waiter для возможной отмены через queue_cancel_wait
queue_waiter_t* queue_wait_threshold(ll_queue_t* q, int max_packets, size_t max_bytes,
queue_threshold_callback_t callback, void* arg);
// Отменить ожидание (удалить waiter из списка)
void queue_cancel_wait(ll_queue_t* q, queue_waiter_t* waiter);
// Получить общий размер данных в очереди (байт)
static inline size_t queue_total_bytes(ll_queue_t* q) {
if (!q) return 0;
return q->total_bytes;
}
#endif // LL_QUEUE_H

195
pkt_normalizer.c

@ -1,9 +1,10 @@
#include "pkt_normalizer.h"
#include "u_async.h"
#include "settings.h"
#include <stdlib.h>
#include <string.h>
#include <stdint.h>
#include "pkt_normalizer.h"
#include "u_async.h"
#include "settings.h"
#include <stdlib.h>
#include <string.h>
#include <stdint.h>
#include <stdio.h>
static void packer_handler(ll_queue_t* q, ll_entry_t* unused, void* arg);
static void unpacker_handler(ll_queue_t* q, ll_entry_t* unused, void* arg);
@ -39,8 +40,9 @@ pn_struct* pkt_normalizer_init(int is_packer) {
free(pn);
return NULL;
}
pn->u.packer.len = 0;
queue_set_callback(pn->input, packer_handler, pn);
pn->u.packer.len = 0;
pn->u.packer.error_count = 0;
queue_set_callback(pn->input, packer_handler, pn);
} else {
pn->u.unpacker.buf = NULL;
pn->u.unpacker.len = 0;
@ -148,6 +150,7 @@ static void cancel_fragment_timeout(pn_struct* pn) {
static void send_buf(pn_struct* pn) {
if (pn->u.packer.len == 0) return;
size_t payload_len = pn->u.packer.len;
printf("[PN DEBUG] send_buf: packer len=%zu, output queue count=%d\n", payload_len, queue_entry_count(pn->output));
ll_entry_t* out = queue_entry_new(2 + payload_len);
if (!out) return;
uint8_t* d = ll_entry_data(out);
@ -157,50 +160,59 @@ static void send_buf(pn_struct* pn) {
pn->u.packer.len = 0;
}
static void packer_handler(ll_queue_t* q, ll_entry_t* unused, void* arg) {
(void)unused;
pn_struct* pn = arg;
size_t max = (size_t)settings.max_fragment_size;
while (queue_entry_count(q) > 0) {
ll_entry_t* entry = queue_entry_get(q);
size_t L = ll_entry_size(entry);
uint8_t* data = ll_entry_data(entry);
uint8_t header[2];
int hsize = get_header(header, L);
size_t needed = (size_t)hsize + L;
if (hsize < 0 || needed > max) {
// Fragment
if (pn->u.packer.len > 0) {
send_buf(pn);
}
size_t remaining = L;
size_t pos = 0;
int fragment_count = 0;
while (remaining > 0) {
size_t chunk;
size_t payload_len;
ll_entry_t* fout;
uint8_t* fd;
uint8_t frag_header[2];
int frag_hsize;
if (fragment_count == 0) {
// Первый фрагмент: FF + общая длина (2 байта)
chunk = remaining > (max - 5) ? (max - 5) : remaining; // 2+1+2+chunk <= max
payload_len = 1 + 2 + chunk; // FF + total_len + data
fout = queue_entry_new(2 + payload_len);
if (!fout) break;
fd = ll_entry_data(fout);
*(uint16_t*)fd = (uint16_t)payload_len;
fd += 2;
*fd++ = 0xFF;
*fd++ = (uint8_t)(L >> 8); // старший байт общей длины
*fd++ = (uint8_t)(L & 0xFF); // младший байт общей длины
} else {
// Проверим, можно ли отправить последний фрагмент как обычный блок
static void packer_handler(ll_queue_t* q, ll_entry_t* unused, void* arg) {
(void)unused;
pn_struct* pn = arg;
size_t max = (size_t)settings.max_fragment_size;
ll_entry_t* entry = queue_entry_get(q);
if (!entry) {
queue_resume_callback(q);
return;
}
size_t L = ll_entry_size(entry);
uint8_t* data = ll_entry_data(entry);
uint8_t header[2];
int hsize = get_header(header, L);
size_t needed = (size_t)hsize + L;
if (hsize < 0 || needed > max) {
// Fragment
if (pn->u.packer.len > 0) {
send_buf(pn);
}
size_t remaining = L;
size_t pos = 0;
int fragment_count = 0;
while (remaining > 0) {
size_t chunk;
size_t payload_len;
ll_entry_t* fout;
uint8_t* fd;
uint8_t frag_header[2];
int frag_hsize;
if (fragment_count == 0) {
// Первый фрагмент: FF + общая длина (2 байта)
chunk = remaining > (max - 5) ? (max - 5) : remaining; // 2+1+2+chunk <= max
payload_len = 1 + 2 + chunk; // FF + total_len + data
fout = queue_entry_new(2 + payload_len);
if (!fout) break;
fd = ll_entry_data(fout);
*(uint16_t*)fd = (uint16_t)payload_len;
fd += 2;
*fd++ = 0xFF;
*fd++ = (uint8_t)(L >> 8); // старший байт общей длины
*fd++ = (uint8_t)(L & 0xFF); // младший байт общей длины
} else {
// Не первый фрагмент
if (remaining <= max - 3) {
// Это последний возможный фрагмент (помещается в один пакет с префиксом FE)
// Пытаемся отправить как обычный блок
frag_hsize = get_header(frag_header, remaining);
if (frag_hsize > 0 && (size_t)frag_hsize + remaining + 2 <= max) {
// Отправляем как обычный блок
// Успешно: обычный блок
payload_len = frag_hsize + remaining;
chunk = remaining;
fout = queue_entry_new(2 + payload_len);
@ -211,8 +223,10 @@ static void packer_handler(ll_queue_t* q, ll_entry_t* unused, void* arg) {
memcpy(fd, frag_header, frag_hsize);
fd += frag_hsize;
} else {
// Промежуточный или последний фрагмент как FE
chunk = remaining > (max - 3) ? (max - 3) : remaining; // 2+1+chunk <= max
// Не удалось отправить как обычный блок - ошибка
pn->u.packer.error_count++;
// Отправляем как FE (нарушение спецификации)
chunk = remaining > (max - 3) ? (max - 3) : remaining;
payload_len = 1 + chunk; // FE + data
fout = queue_entry_new(2 + payload_len);
if (!fout) break;
@ -221,46 +235,69 @@ static void packer_handler(ll_queue_t* q, ll_entry_t* unused, void* arg) {
fd += 2;
*fd++ = 0xFE;
}
} else {
// Промежуточный фрагмент, отправляем как FE
chunk = remaining > (max - 3) ? (max - 3) : remaining;
payload_len = 1 + chunk; // FE + data
fout = queue_entry_new(2 + payload_len);
if (!fout) break;
fd = ll_entry_data(fout);
*(uint16_t*)fd = (uint16_t)payload_len;
fd += 2;
*fd++ = 0xFE;
}
memcpy(fd, data + pos, chunk);
queue_entry_put(pn->output, fout);
pos += chunk;
remaining -= chunk;
fragment_count++;
}
} else {
if (pn->u.packer.len + needed > max) {
send_buf(pn);
}
// Add to buffer
uint8_t* p = pn->u.packer.buf + pn->u.packer.len;
memcpy(p, header, (size_t)hsize);
memcpy(p + hsize, data, L);
pn->u.packer.len += needed;
}
queue_entry_free(entry);
}
if (pn->u.packer.len > 0) {
send_buf(pn);
}
memcpy(fd, data + pos, chunk);
queue_entry_put(pn->output, fout);
pos += chunk;
remaining -= chunk;
fragment_count++;
}
} else {
if (pn->u.packer.len + needed > max) {
send_buf(pn);
}
// Add to buffer
uint8_t* p = pn->u.packer.buf + pn->u.packer.len;
memcpy(p, header, (size_t)hsize);
memcpy(p + hsize, data, L);
pn->u.packer.len += needed;
}
queue_entry_free(entry);
if (pn->u.packer.len > 0) {
send_buf(pn);
}
queue_resume_callback(q);
}
int pkt_normalizer_get_error_count(const pn_struct* pn) {
if (!pn || pn->is_packer) {
return 0;
if (!pn) return 0;
if (pn->is_packer) {
return pn->u.packer.error_count;
}
return pn->u.unpacker.error_count;
}
void pkt_normalizer_reset_error_count(pn_struct* pn) {
if (!pn || pn->is_packer) {
return;
if (!pn) return;
if (pn->is_packer) {
pn->u.packer.error_count = 0;
} else {
pn->u.unpacker.error_count = 0;
}
pn->u.unpacker.error_count = 0;
}
static void unpacker_handler(ll_queue_t* q, ll_entry_t* unused, void* arg) {
void pkt_normalizer_flush(pn_struct* pn) {
if (!pn || !pn->is_packer) return;
if (pn->u.packer.len > 0) {
send_buf(pn);
}
}
static void unpacker_handler(ll_queue_t* q, ll_entry_t* unused, void* arg) {
(void)unused;
pn_struct* pn = arg;
while (queue_entry_count(q) > 0) {

4
pkt_normalizer.h

@ -22,6 +22,7 @@ struct pn_struct {
uint8_t* buf;
size_t len;
size_t cap;
int error_count;
} packer;
struct {
uint8_t* buf; /* буфер для сборки фрагментов */
@ -44,6 +45,9 @@ void pkt_normalizer_pair_deinit(pkt_normalizer_pair* pair);
/* Error handling */
int pkt_normalizer_get_error_count(const pn_struct* pn);
void pkt_normalizer_reset_error_count(pn_struct* pn);
/* Flush internal buffer (packer only) */
void pkt_normalizer_flush(pn_struct* pn);
struct pkt_normalizer_pair {
pn_struct* packer;

2842
stress.log

File diff suppressed because it is too large Load Diff

2834
stress.out

File diff suppressed because it is too large Load Diff

1355
stress_50.log

File diff suppressed because it is too large Load Diff

1098
stress_new.log

File diff suppressed because it is too large Load Diff

1
tests/simple_uasync.h

@ -2,6 +2,7 @@
#ifndef SIMPLE_UASYNC_H
#define SIMPLE_UASYNC_H
#define ETCP_DEBUG 1
#include <stdint.h>
// These functions are only available when linking with simple_uasync.o

1
tests/test_etcp_simple.c

@ -1,4 +1,5 @@
// test_etcp_simple.c - Simple test to verify ETCP sender-receiver communication
#define ETCP_DEBUG 1
#include "etcp.h"
#include "u_async.h"
#include "ll_queue.h"

12
tests/test_etcp_stress.c

@ -10,13 +10,13 @@
#include <assert.h>
#include <time.h>
#define NUM_PACKETS 1000
#define NUM_PACKETS 10000
#define MIN_PACKET_SIZE 1
#define MAX_PACKET_SIZE 1300
#define LOSS_PROBABILITY 0.0 // 0% packet loss for window testing
#define REORDER_PROBABILITY 0.0 // 0% reordering for testing
#define QUEUE_MAX_SIZE 1000 // Max packets in delay queue
#define MAX_DELAY_MS 0 // No delay for testing
#define LOSS_PROBABILITY 0.0 // 0% packet loss for testing reordering
#define REORDER_PROBABILITY 0.1 // 10% reordering for testing
#define QUEUE_MAX_SIZE 5000 // Max packets in delay queue
#define MAX_DELAY_MS 100 // Up to 100ms delay for reordering effect
#define TIME_BASE_MS 0.1 // uasync timebase is 0.1ms
// Packet in the network delay queue
@ -322,7 +322,7 @@ int main(void) {
// Let network deliver remaining packets and wait for retransmissions
printf("\nDelivering remaining packets and waiting for retransmissions...\n");
uint32_t start_time = net.current_time;
for (int i = 0; i < 50000 && (net.queue_size > 0 || i < 1000); i++) {
for (int i = 0; i < 200000 && (net.queue_size > 0 || i < 1000); i++) {
simple_uasync_advance_time(10); // Advance 1ms
net.current_time = simple_uasync_get_time(); // Keep in sync
deliver_packets(&net);

373
tests/test_pkt_normalizer.c

@ -25,7 +25,7 @@ typedef struct {
size_t len;
} test_packet_t;
#define MAX_TEST_PACKETS 20
#define MAX_TEST_PACKETS 20000
static test_packet_t sent_packets[MAX_TEST_PACKETS];
static test_packet_t received_packets[MAX_TEST_PACKETS];
static int sent_count = 0;
@ -136,10 +136,11 @@ static void process_queues(pkt_normalizer_pair* pair) {
do {
processed = 0;
iterations++;
if (iterations % 10 == 1) printf("process_queues iteration %d: counts: pi=%d po=%d ui=%d uo=%d\n", iterations, queue_entry_count(pair->packer->input), queue_entry_count(pair->packer->output), queue_entry_count(pair->unpacker->input), queue_entry_count(pair->unpacker->output));
/* Обработать входную очередь упаковщика */
if (pair->packer->input->callback &&
queue_entry_count(pair->packer->input) > 0) {
while (pair->packer->input->callback &&
queue_entry_count(pair->packer->input) > 0) {
pair->packer->input->callback(pair->packer->input,
pair->packer->input->head,
pair->packer->input->callback_arg);
@ -173,7 +174,7 @@ static void process_queues(pkt_normalizer_pair* pair) {
processed = 1;
}
if (iterations > 50) {
if (iterations > 100000) {
printf("ERROR: process_queues infinite loop detected\n");
break;
}
@ -282,18 +283,20 @@ static int test_mixed_packets(pkt_normalizer_pair* pair) {
/* Тест 6: стресс-тест со случайными размерами пакетов */
static int test_stress_random(pkt_normalizer_pair* pair) {
printf("\n--- Test 6: stress test with random packets ---\n");
const int NUM_PACKETS = 100;
const int MAX_PACKET_SIZE = 4000;
printf("\n--- Test 6: stress test with 10000 random packets (0-4096 bytes) ---\n");
const int NUM_PACKETS = 10000;
const int MAX_PACKET_SIZE = 4096;
int original_fragment_size = settings.max_fragment_size;
/* Часть 1: случайные пакеты без изменения размера фрагмента */
/* Часть 1: 10000 случайных пакетов без изменения размера фрагмента */
printf("Sending %d random packets...\n", NUM_PACKETS);
for (int i = 0; i < NUM_PACKETS; i++) {
reset_test_data();
pkt_normalizer_reset_error_count(pair->packer);
pkt_normalizer_reset_error_count(pair->unpacker);
/* Случайный размер пакета от 1 до MAX_PACKET_SIZE */
size_t size = (rand() % MAX_PACKET_SIZE) + 1;
/* Случайный размер пакета от 0 до MAX_PACKET_SIZE */
size_t size = rand() % (MAX_PACKET_SIZE + 1);
ll_entry_t* entry = create_test_packet(size);
TEST_ASSERT(entry != NULL, "create random packet");
@ -303,41 +306,117 @@ static int test_stress_random(pkt_normalizer_pair* pair) {
TEST_ASSERT(queue_entry_put(pair->packer->input, entry) == 0, "put random packet");
process_queues(pair);
/* Проверить, что пакет корректно обработан */
TEST_ASSERT(compare_packets() == 0, "random packet comparison");
TEST_ASSERT(pkt_normalizer_get_error_count(pair->unpacker) == 0, "no errors during random packet");
// Периодически обрабатывать очереди чтобы не переполнять
if (i % 500 == 0) {
process_queues(pair);
}
}
/* Часть 2: фрагментация со случайным размером фрагмента */
const int FRAGMENT_TESTS = 50;
for (int i = 0; i < FRAGMENT_TESTS; i++) {
reset_test_data();
pkt_normalizer_reset_error_count(pair->unpacker);
/* Случайный размер фрагмента от 50 до 500 */
int frag_size = (rand() % 451) + 50;
settings.max_fragment_size = frag_size;
/* Случайный размер пакета от frag_size+1 до 3000 (чтобы гарантировать фрагментацию) */
size_t packet_size = (rand() % (3000 - frag_size)) + frag_size + 1;
ll_entry_t* entry = create_test_packet(packet_size);
TEST_ASSERT(entry != NULL, "create fragmentable packet");
uint8_t* data = ll_entry_data(entry);
add_sent_packet(data, packet_size);
TEST_ASSERT(queue_entry_put(pair->packer->input, entry) == 0, "put fragmentable packet");
process_queues(pair);
/* Проверить, что пакет корректно обработан */
TEST_ASSERT(compare_packets() == 0, "fragmentable packet comparison");
TEST_ASSERT(pkt_normalizer_get_error_count(pair->unpacker) == 0, "no errors during fragmentation");
// Обработать оставшиеся пакеты
process_queues(pair);
printf("Sent %d packets, received %d packets\n", sent_count, received_count);
TEST_ASSERT(sent_count == NUM_PACKETS, "all packets sent");
TEST_ASSERT(received_count == NUM_PACKETS, "all packets received");
TEST_ASSERT(compare_packets() == 0, "all packets matched");
TEST_ASSERT(pkt_normalizer_get_error_count(pair->packer) == 0, "no packer errors");
TEST_ASSERT(pkt_normalizer_get_error_count(pair->unpacker) == 0, "no unpacker errors");
/* Часть 2: проверка дефрагментации (большой пакет фрагментируется, а не отправляется мелкими пакетами) */
printf("\n--- Verifying fragmentation behavior ---\n");
settings.max_fragment_size = 500;
reset_test_data();
pkt_normalizer_reset_error_count(pair->packer);
pkt_normalizer_reset_error_count(pair->unpacker);
// Большой пакет, который должен быть фрагментирован
size_t big_packet_size = 2000;
ll_entry_t* entry = create_test_packet(big_packet_size);
TEST_ASSERT(entry != NULL, "create big packet");
add_sent_packet(ll_entry_data(entry), big_packet_size);
// Подсчитать количество элементов в output очереди packer'а до обработки
int initial_output_count = queue_entry_count(pair->packer->output);
TEST_ASSERT(queue_entry_put(pair->packer->input, entry) == 0, "put big packet");
// Обработать только packer (чтобы фрагменты появились в его output)
if (pair->packer->input->callback &&
queue_entry_count(pair->packer->input) > 0) {
pair->packer->input->callback(pair->packer->input,
pair->packer->input->head,
pair->packer->input->callback_arg);
}
// Проверить количество фрагментов в output очереди packer'а
int fragment_count = queue_entry_count(pair->packer->output);
printf("Big packet %zu bytes, fragment size %d, produced %d fragments\n",
big_packet_size, settings.max_fragment_size, fragment_count);
// Должно быть больше 1 фрагмента
TEST_ASSERT(fragment_count > 1, "packet was fragmented");
// Ожидаемое количество фрагментов: ceil(2000 / (500 - overhead))
TEST_ASSERT(fragment_count >= 3 && fragment_count <= 6, "reasonable fragment count");
// Теперь обработать все фрагменты через unpacker и проверить сборку
process_queues(pair);
TEST_ASSERT(compare_packets() == 0, "fragmented packet correctly reassembled");
TEST_ASSERT(pkt_normalizer_get_error_count(pair->unpacker) == 0, "no unpacker errors");
/* Часть 3: проверка отправки буфера при пустой входной очереди */
printf("\n--- Verifying buffer flush on empty input queue ---\n");
settings.max_fragment_size = 500;
reset_test_data();
pkt_normalizer_reset_error_count(pair->packer);
pkt_normalizer_reset_error_count(pair->unpacker);
// Отправить маленький пакет, который поместится в буфер packer'а
size_t small_size = 100;
ll_entry_t* entry1 = create_test_packet(small_size);
TEST_ASSERT(entry1 != NULL, "create first small packet");
add_sent_packet(ll_entry_data(entry1), small_size);
TEST_ASSERT(queue_entry_put(pair->packer->input, entry1) == 0, "put first packet");
// Обработать packer input (пакет должен добавиться в буфер, но не отправиться)
if (pair->packer->input->callback &&
queue_entry_count(pair->packer->input) > 0) {
pair->packer->input->callback(pair->packer->input,
pair->packer->input->head,
pair->packer->input->callback_arg);
}
// Проверить, что в output очереди ничего нет (пакет ещё не отправлен)
int output_before = queue_entry_count(pair->packer->output);
printf("After first packet: packer output queue count = %d\n", output_before);
// Отправим второй маленький пакет, который заставит packer отправить буфер
ll_entry_t* entry2 = create_test_packet(small_size);
TEST_ASSERT(entry2 != NULL, "create second small packet");
add_sent_packet(ll_entry_data(entry2), small_size);
TEST_ASSERT(queue_entry_put(pair->packer->input, entry2) == 0, "put second packet");
// Обработать packer input
if (pair->packer->input->callback &&
queue_entry_count(pair->packer->input) > 0) {
pair->packer->input->callback(pair->packer->input,
pair->packer->input->head,
pair->packer->input->callback_arg);
}
// Теперь в output должна быть хотя бы одна запись (буфер отправлен)
int output_after = queue_entry_count(pair->packer->output);
printf("After second packet: packer output queue count = %d\n", output_after);
TEST_ASSERT(output_after > output_before, "buffer flushed when new packet arrives");
// Обработать все оставшиеся очереди
process_queues(pair);
TEST_ASSERT(compare_packets() == 0, "both packets correctly processed");
TEST_ASSERT(pkt_normalizer_get_error_count(pair->unpacker) == 0, "no unpacker errors");
settings.max_fragment_size = original_fragment_size;
return 0;
}
@ -350,6 +429,7 @@ static int test_edge_cases(pkt_normalizer_pair* pair) {
/* Сбросить данные теста */
reset_test_data();
pkt_normalizer_reset_error_count(pair->unpacker);
pkt_normalizer_flush(pair->packer);
/* Тест 1: минимальный размер пакета (1 байт) */
{
@ -367,6 +447,7 @@ static int test_edge_cases(pkt_normalizer_pair* pair) {
{
reset_test_data();
pkt_normalizer_reset_error_count(pair->unpacker);
pkt_normalizer_flush(pair->packer);
size_t size = 239;
ll_entry_t* entry = create_test_packet(size);
TEST_ASSERT(entry != NULL, "create 239-byte packet");
@ -393,6 +474,7 @@ static int test_edge_cases(pkt_normalizer_pair* pair) {
{
reset_test_data();
pkt_normalizer_reset_error_count(pair->unpacker);
pkt_normalizer_flush(pair->packer);
size_t size = 3839;
ll_entry_t* entry = create_test_packet(size);
TEST_ASSERT(entry != NULL, "create 3839-byte packet");
@ -407,6 +489,7 @@ static int test_edge_cases(pkt_normalizer_pair* pair) {
{
reset_test_data();
pkt_normalizer_reset_error_count(pair->unpacker);
pkt_normalizer_flush(pair->packer);
size_t size = 3840;
ll_entry_t* entry = create_test_packet(size);
TEST_ASSERT(entry != NULL, "create 3840-byte packet");
@ -421,6 +504,7 @@ static int test_edge_cases(pkt_normalizer_pair* pair) {
{
reset_test_data();
pkt_normalizer_reset_error_count(pair->unpacker);
pkt_normalizer_flush(pair->packer);
settings.max_fragment_size = 500;
size_t size = 480; /* достаточно мало, чтобы поместиться с заголовком */
ll_entry_t* entry = create_test_packet(size);
@ -436,6 +520,7 @@ static int test_edge_cases(pkt_normalizer_pair* pair) {
{
reset_test_data();
pkt_normalizer_reset_error_count(pair->unpacker);
pkt_normalizer_flush(pair->packer);
settings.max_fragment_size = 500;
size_t size = 520; /* потребует фрагментации */
ll_entry_t* entry = create_test_packet(size);
@ -447,10 +532,11 @@ static int test_edge_cases(pkt_normalizer_pair* pair) {
TEST_ASSERT(pkt_normalizer_get_error_count(pair->unpacker) == 0, "no errors");
}
/* Тест 7: максимальный размер пакета с фрагментацией (65535) - ограничим 10000 для скорости */
/* Тест 7: максимальный размер пакета с фрагментации (65535) - ограничим 10000 для скорости */
{
reset_test_data();
pkt_normalizer_reset_error_count(pair->unpacker);
pkt_normalizer_flush(pair->packer);
settings.max_fragment_size = 1000;
size_t size = 10000;
ll_entry_t* entry = create_test_packet(size);
@ -466,6 +552,204 @@ static int test_edge_cases(pkt_normalizer_pair* pair) {
return 0;
}
/* Тест 8: 10000 пакетов случайного размера (1-1024 байт) */
static int test_10000_packets(pkt_normalizer_pair* pair) {
printf("\n--- Test 8: 10000 random packets (1-1024 bytes) ---\n");
const int NUM_PACKETS = 10000;
const int MAX_PACKET_SIZE = 1024;
reset_test_data();
pkt_normalizer_reset_error_count(pair->packer);
pkt_normalizer_reset_error_count(pair->unpacker);
printf("Sending %d packets...\n", NUM_PACKETS);
for (int i = 0; i < NUM_PACKETS; i++) {
// Случайный размер от 1 до MAX_PACKET_SIZE
size_t size = (rand() % MAX_PACKET_SIZE) + 1;
ll_entry_t* entry = create_test_packet(size);
TEST_ASSERT(entry != NULL, "create random packet");
uint8_t* data = ll_entry_data(entry);
add_sent_packet(data, size);
TEST_ASSERT(queue_entry_put(pair->packer->input, entry) == 0, "put random packet");
// Периодически обрабатывать очереди чтобы не переполнять
if (i % 500 == 0) {
process_queues(pair);
}
}
// Обработать оставшиеся пакеты
process_queues(pair);
printf("Sent %d packets, received %d packets\n", sent_count, received_count);
TEST_ASSERT(sent_count == NUM_PACKETS, "all packets sent");
TEST_ASSERT(received_count == NUM_PACKETS, "all packets received");
TEST_ASSERT(compare_packets() == 0, "all packets matched");
TEST_ASSERT(pkt_normalizer_get_error_count(pair->packer) == 0, "no packer errors");
TEST_ASSERT(pkt_normalizer_get_error_count(pair->unpacker) == 0, "no unpacker errors");
return 0;
}
/* Тест 9: проверка дефрагментации (большой пакет фрагментируется, не отправляется мелкими пакетами) */
static int test_fragmentation_verify(pkt_normalizer_pair* pair) {
printf("\n--- Test 9: fragmentation verification ---\n");
int original_fragment_size = settings.max_fragment_size;
// Установить маленький размер фрагмента для теста
settings.max_fragment_size = 500;
reset_test_data();
pkt_normalizer_reset_error_count(pair->packer);
pkt_normalizer_reset_error_count(pair->unpacker);
// Большой пакет, который должен быть фрагментирован
size_t big_packet_size = 2000; // Должен быть разбит на ~4 фрагмента (500 байт каждый с накладными расходами)
ll_entry_t* entry = create_test_packet(big_packet_size);
TEST_ASSERT(entry != NULL, "create big packet");
add_sent_packet(ll_entry_data(entry), big_packet_size);
// Подсчитать количество элементов в output очереди packer'а до обработки
int initial_output_count = queue_entry_count(pair->packer->output);
TEST_ASSERT(queue_entry_put(pair->packer->input, entry) == 0, "put big packet");
// Обработать только packer (чтобы фрагменты появились в его output)
// Вместо process_queues вызовем обработку только packer input и output
if (pair->packer->input->callback &&
queue_entry_count(pair->packer->input) > 0) {
pair->packer->input->callback(pair->packer->input,
pair->packer->input->head,
pair->packer->input->callback_arg);
}
// Проверить количество фрагментов в output очереди packer'а
int fragment_count = queue_entry_count(pair->packer->output);
printf("Big packet %zu bytes, fragment size %d, produced %d fragments\n",
big_packet_size, settings.max_fragment_size, fragment_count);
// Должно быть больше 1 фрагмента
TEST_ASSERT(fragment_count > 1, "packet was fragmented");
// Ожидаемое количество фрагментов: ceil(2000 / (500 - overhead))
// overhead зависит от заголовков. Минимум 1 байт заголовка на фрагмент + 2 байта длины
// Приблизительно 4-5 фрагментов
TEST_ASSERT(fragment_count >= 3 && fragment_count <= 6, "reasonable fragment count");
// Теперь обработать все фрагменты через unpacker и проверить сборку
process_queues(pair);
TEST_ASSERT(compare_packets() == 0, "fragmented packet correctly reassembled");
TEST_ASSERT(pkt_normalizer_get_error_count(pair->unpacker) == 0, "no unpacker errors");
settings.max_fragment_size = original_fragment_size;
return 0;
}
/* Тест 10: проверка отправки буфера при пустой входной очереди */
static int test_empty_queue_flush(pkt_normalizer_pair* pair) {
printf("\n--- Test 10: empty queue buffer flush ---\n");
int original_fragment_size = settings.max_fragment_size;
settings.max_fragment_size = 500;
reset_test_data();
pkt_normalizer_reset_error_count(pair->packer);
pkt_normalizer_reset_error_count(pair->unpacker);
// Отправить маленький пакет, который поместится в буфер packer'а
size_t small_size = 100;
ll_entry_t* entry1 = create_test_packet(small_size);
TEST_ASSERT(entry1 != NULL, "create first small packet");
add_sent_packet(ll_entry_data(entry1), small_size);
TEST_ASSERT(queue_entry_put(pair->packer->input, entry1) == 0, "put first packet");
// Обработать packer input (пакет должен добавиться в буфер, но не отправиться)
if (pair->packer->input->callback &&
queue_entry_count(pair->packer->input) > 0) {
pair->packer->input->callback(pair->packer->input,
pair->packer->input->head,
pair->packer->input->callback_arg);
}
// Проверить, что буфер не пуст (packer накопил данные)
// packer должен иметь len > 0, но мы не можем получить доступ к внутреннему состоянию.
// Вместо этого проверим, что в output очереди ничего нет (пакет ещё не отправлен)
int output_before = queue_entry_count(pair->packer->output);
printf("After first packet: packer output queue count = %d\n", output_before);
// Теперь симулируем ситуацию, когда входящая очередь пуста и вызывается callback
// В реальности packer_handler вызывается при добавлении нового пакета.
// Но мы можем проверить, что при следующем пакете буфер будет отправлен.
// Отправим второй маленький пакет, который заставит packer отправить буфер
ll_entry_t* entry2 = create_test_packet(small_size);
TEST_ASSERT(entry2 != NULL, "create second small packet");
add_sent_packet(ll_entry_data(entry2), small_size);
TEST_ASSERT(queue_entry_put(pair->packer->input, entry2) == 0, "put second packet");
// Обработать packer input
if (pair->packer->input->callback &&
queue_entry_count(pair->packer->input) > 0) {
pair->packer->input->callback(pair->packer->input,
pair->packer->input->head,
pair->packer->input->callback_arg);
}
// Теперь в output должна быть хотя бы одна запись (буфер отправлен)
int output_after = queue_entry_count(pair->packer->output);
printf("After second packet: packer output queue count = %d\n", output_after);
TEST_ASSERT(output_after > output_before, "buffer flushed when new packet arrives");
// Обработать все оставшиеся очереди
process_queues(pair);
TEST_ASSERT(compare_packets() == 0, "both packets correctly processed");
TEST_ASSERT(pkt_normalizer_get_error_count(pair->unpacker) == 0, "no unpacker errors");
settings.max_fragment_size = original_fragment_size;
return 0;
}
/* Тест async wait: проверка асинхронного ожидания порога в очереди */
static void test_callback(ll_queue_t* q_arg, void* arg) {
(void)q_arg;
int* flag = (int*)arg;
*flag = 1;
}
static int test_async_wait(void) {
printf("\n--- Test async wait: threshold waiter ---\n");
ll_queue_t* q = queue_new();
TEST_ASSERT(q != NULL, "create queue for async wait test");
int callback_called = 0;
// Добавим элемент, чтобы очередь была непустой
ll_entry_t* entry = queue_entry_new(10);
TEST_ASSERT(entry != NULL, "create entry");
TEST_ASSERT(queue_entry_put(q, entry) == 0, "put entry");
// Регистрируем ожидание: когда очередь будет <= 0 пакетов и <= 0 байт
// Условие не выполнено, поэтому коллбэк не должен вызваться сразу
queue_waiter_t* waiter = queue_wait_threshold(q, 0, 0, test_callback, &callback_called);
TEST_ASSERT(waiter != NULL, "waiter registered");
TEST_ASSERT(callback_called == 0, "callback not called immediately");
// Извлекаем элемент - условие должно выполниться и коллбэк вызваться
ll_entry_t* retrieved = queue_entry_get(q);
TEST_ASSERT(retrieved != NULL, "retrieved entry");
queue_entry_free(retrieved);
// Коллбэк должен был быть вызван внутри queue_entry_get через check_waiters
TEST_ASSERT(callback_called == 1, "callback called after threshold met");
// Отменяем waiter (хотя он уже должен быть удален)
queue_cancel_wait(q, waiter);
queue_free(q);
return 0;
}
/* Тест 5: проверка жизненного цикла */
static int test_lifecycle(void) {
printf("\n--- Test 5: lifecycle ---\n");
@ -519,7 +803,10 @@ int main(void) {
result |= test_mixed_packets(pair);
result |= test_stress_random(pair);
result |= test_edge_cases(pair);
result |= test_fragmentation_verify(pair);
result |= test_empty_queue_flush(pair);
result |= test_async_wait();
pkt_normalizer_pair_deinit(pair);
reset_test_data();

87
u_async.c

@ -234,43 +234,62 @@ err_t uasync_remove_socket(void* s_id) {
void uasync_mainloop(void) {
while (1) {
// Process timeouts first
process_timeouts();
// Prepare select with copies of masters (point 1: no loop/FD_ZERO here)
fd_set readfds = master_readfds;
fd_set writefds = master_writefds;
fd_set exceptfds = master_exceptfds;
struct timeval tv;
get_next_timeout(&tv);
struct timeval* ptv = (tv.tv_sec == 0 && tv.tv_usec == 0 && !timeout_head) ? NULL : &tv;
int nfds = select(max_fd + 1, &readfds, &writefds, &exceptfds, ptv);
if (nfds < 0) {
if (errno == EINTR) continue;
perror("select");
break;
uasync_poll(-1); /* infinite timeout */
}
}
void uasync_poll(int timeout_tb) {
/* Process expired timeouts */
process_timeouts();
/* Prepare select with copies of masters */
fd_set readfds = master_readfds;
fd_set writefds = master_writefds;
fd_set exceptfds = master_exceptfds;
struct timeval tv;
get_next_timeout(&tv);
/* If timeout_tb >= 0, compute timeout as min(timeout_tb, existing timer) */
if (timeout_tb >= 0) {
struct timeval user_tv;
user_tv.tv_sec = timeout_tb / 10000;
user_tv.tv_usec = (timeout_tb % 10000) * 100;
/* If no internal timer or user timeout is smaller */
if (tv.tv_sec == 0 && tv.tv_usec == 0 && !timeout_head) {
tv = user_tv;
} else if (user_tv.tv_sec < tv.tv_sec ||
(user_tv.tv_sec == tv.tv_sec && user_tv.tv_usec < tv.tv_usec)) {
tv = user_tv;
}
}
struct timeval* ptv = (tv.tv_sec == 0 && tv.tv_usec == 0 && !timeout_head) ? NULL : &tv;
int nfds = select(max_fd + 1, &readfds, &writefds, &exceptfds, ptv);
if (nfds < 0) {
if (errno == EINTR) return;
perror("select");
return;
}
// Process sockets with faster dispatch (point 2: loop only up to max_fd, but use map for O(1) node lookup)
// This is O(max_fd) worst-case, but in practice fast; only checks if FD_ISSET.
for (int fd = 0; nfds > 0 && fd <= max_fd; fd++) {
struct socket_node* node = fd_to_node[fd];
if (!node) continue; // Skip unmapped fds
/* Process sockets with faster dispatch */
for (int fd = 0; nfds > 0 && fd <= max_fd; fd++) {
struct socket_node* node = fd_to_node[fd];
if (!node) continue;
if (node->except_cbk && FD_ISSET(fd, &exceptfds)) {
node->except_cbk(fd, node->user_data);
nfds--;
}
if (node->read_cbk && FD_ISSET(fd, &readfds)) {
node->read_cbk(fd, node->user_data);
nfds--;
}
if (node->write_cbk && FD_ISSET(fd, &writefds)) {
node->write_cbk(fd, node->user_data);
nfds--;
}
if (node->except_cbk && FD_ISSET(fd, &exceptfds)) {
node->except_cbk(fd, node->user_data);
nfds--;
}
if (node->read_cbk && FD_ISSET(fd, &readfds)) {
node->read_cbk(fd, node->user_data);
nfds--;
}
if (node->write_cbk && FD_ISSET(fd, &writefds)) {
node->write_cbk(fd, node->user_data);
nfds--;
}
}
}

3
u_async.h

@ -30,4 +30,7 @@ err_t uasync_cancel_timeout(void* t_id);
void* uasync_add_socket(int fd, socket_callback_t read_cbk, socket_callback_t write_cbk, socket_callback_t except_cbk, void* user_arg);
err_t uasync_remove_socket(void* s_id);
// Single iteration of event loop with timeout (timebase units)
void uasync_poll(int timeout_tb);
#endif // UASYNC_H

Loading…
Cancel
Save