You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 
 
 

645 lines
21 KiB

/**
* @file dummynet.c
* @brief Реализация UDP dummynet эмулятора
*/
#include "dummynet.h"
#include "../lib/u_async.h"
#include "../lib/ll_queue.h"
#include "../lib/memory_pool.h"
#include "../lib/socket_compat.h"
#include "../lib/debug_config.h"
#include "../lib/platform_compat.h"
#include "etcp_connections.h"
#include <stdlib.h>
#include <string.h>
#include <stdio.h>
#include <errno.h>
#include "../lib/mem.h"
/* Timebase: 0.1 мс */
#define TIMEBASE_US 100
#define MS_TO_TB(ms) ((ms) * 10)
/* Burst allowance для shaper (в timebase units) */
#define SHAPER_BURST_TB 10
/* Debug category */
#ifndef DEBUG_CATEGORY_DUMMYNET
#define DEBUG_CATEGORY_DUMMYNET ((debug_category_t)1 << 14)
#endif
/* Forward declarations */
static void dummynet_read_callback(socket_t sock, void* user_arg);
static void dummynet_delay_callback(void* user_arg);
static void dummynet_shaper_callback(void* user_arg);
/**
* @brief Пакет в обработке (внутренняя структура)
*/
struct dummynet_pkt {
uint8_t data[DUMMYNET_MAX_PKT_SIZE];
size_t len;
struct sockaddr_storage src_addr;
socklen_t src_len;
uint64_t recv_time_tb;
int direction;
};
/**
* @brief Направление передачи
*/
struct dummynet_dir {
uint32_t delay_fixed_ms;
uint32_t delay_random_ms;
uint32_t bandwidth_kbps;
uint32_t max_queue_pkts;
uint32_t loss_permille;
struct sockaddr_storage dest_addr;
socklen_t dest_len;
int configured;
struct ll_queue* queue;
void* shaper_timer;
int shaper_pending;
/* Отладочные счетчики */
uint64_t delay_timer_set; /* Число установок delay таймеров */
uint64_t delay_timer_fire; /* Число срабатываний delay таймеров */
uint64_t shaper_timer_set; /* Число установок shaper таймеров */
uint64_t shaper_timer_fire; /* Число срабатываний shaper таймеров */
struct dummynet_stats stats;
};
/**
* @brief Контекст dummynet
*/
struct dummynet {
struct UASYNC* ua;
socket_t sock;
struct sockaddr_storage bind_addr;
socklen_t bind_len;
uint16_t listen_port;
struct dummynet_dir dirs[DUMMYNET_DIR_COUNT];
struct memory_pool* pkt_pool;
};
/**
* @brief Аргумент для delay callback
*/
struct delay_cb_arg {
struct dummynet* dn;
struct ll_entry* entry;
};
/**
* @brief Аргумент для shaper callback
*/
struct shaper_cb_arg {
struct dummynet* dn;
int dir_idx;
};
/* Генерация случайного числа 0..max */
static inline uint32_t random_range(uint32_t max) {
if (max == 0) return 0;
return (uint32_t)rand() % (max + 1);
}
/* Определяет направление по порту источника */
static int dummynet_get_direction(struct dummynet* dn, uint16_t src_port) {
if (src_port == dn->listen_port - 1) {
return DUMMYNET_FORWARD;
} else if (src_port == dn->listen_port + 1) {
return DUMMYNET_BACKWARD;
}
return -1;
}
/* Отправляет пакет и планирует следующую отправку */
static void dummynet_process_queue(struct dummynet* dn, int dir_idx) {
struct dummynet_dir* dir = &dn->dirs[dir_idx];
struct ll_entry* entry = queue_data_get(dir->queue);
if (!entry) {
dir->shaper_pending = 0;
return;
}
struct dummynet_pkt* pkt = (struct dummynet_pkt*)entry->data;
ssize_t n = socket_sendto(dn->sock, pkt->data, pkt->len,
(struct sockaddr*)&dir->dest_addr, dir->dest_len);
if (n < 0) {
DEBUG_ERROR(DEBUG_CATEGORY_DUMMYNET, "sendto failed: %s", socket_strerror(socket_get_error()));
} else {
dir->stats.sent++;
DEBUG_DEBUG(DEBUG_CATEGORY_DUMMYNET, "Dir %d: sent packet (%zd bytes)", dir_idx, n);
}
size_t bytes_sent = (n > 0) ? (size_t)n : 0;
dir->stats.queue_size = queue_entry_count(dir->queue);
queue_entry_free(entry);
/* Планируем следующую отправку */
if (dir->stats.queue_size > 0) {
int delay_tb = 0;
if (dir->bandwidth_kbps > 0) {
uint64_t delay_calc = (uint64_t)bytes_sent * 80ULL / (uint64_t)dir->bandwidth_kbps;
delay_tb = (int)delay_calc;
if (delay_tb == 0) delay_tb = 1;
}
struct shaper_cb_arg* sarg = (struct shaper_cb_arg*)u_malloc(sizeof(struct shaper_cb_arg));
if (sarg) {
sarg->dn = dn;
sarg->dir_idx = dir_idx;
dir->shaper_timer = uasync_set_timeout(dn->ua, delay_tb, sarg, dummynet_shaper_callback, "dummynet_shaper");
if (dir->shaper_timer) {
dir->shaper_timer_set++;
} else {
u_free(sarg);
dir->shaper_pending = 0;
}
} else {
dir->shaper_pending = 0;
}
} else {
dir->shaper_pending = 0;
}
}
/* Callback шейпера */
static void dummynet_shaper_callback(void* user_arg) {
struct shaper_cb_arg* arg = (struct shaper_cb_arg*)user_arg;
struct dummynet* dn = arg->dn;
int dir_idx = arg->dir_idx;
u_free(arg);
dn->dirs[dir_idx].shaper_timer_fire++;
dummynet_process_queue(dn, dir_idx);
}
/* Callback чтения из UDP сокета */
static void dummynet_read_callback(socket_t sock, void* user_arg) {
struct dummynet* dn = (struct dummynet*)user_arg;
uint8_t buf[DUMMYNET_MAX_PKT_SIZE];
struct sockaddr_storage src_addr;
socklen_t src_len = sizeof(src_addr);
ssize_t n = socket_recvfrom(sock, buf, sizeof(buf), (struct sockaddr*)&src_addr, &src_len);
if (n < 0) {
if (socket_get_error() != ERR_WOULDBLOCK && socket_get_error() != ERR_AGAIN) {
DEBUG_ERROR(DEBUG_CATEGORY_DUMMYNET, "recvfrom failed: %s", socket_strerror(socket_get_error()));
}
return;
}
if (n == 0) return;
uint16_t src_port = 0;
if (src_addr.ss_family == AF_INET) {
src_port = ntohs(((struct sockaddr_in*)&src_addr)->sin_port);
} else if (src_addr.ss_family == AF_INET6) {
src_port = ntohs(((struct sockaddr_in6*)&src_addr)->sin6_port);
}
int dir_idx = dummynet_get_direction(dn, src_port);
if (dir_idx < 0) {
DEBUG_WARN(DEBUG_CATEGORY_DUMMYNET, "Packet from unknown port %u, expected %u or %u",
src_port, dn->listen_port - 1, dn->listen_port + 1);
return;
}
struct dummynet_dir* dir = &dn->dirs[dir_idx];
if (!dir->configured) {
DEBUG_WARN(DEBUG_CATEGORY_DUMMYNET, "Direction %d not configured, dropping packet", dir_idx);
return;
}
dir->stats.recv++;
/* Random loss check */
if (dir->loss_permille > 0) {
uint32_t roll = (uint32_t)rand() % 1000;
if (roll < dir->loss_permille) {
dir->stats.lost++;
DEBUG_INFO(DEBUG_CATEGORY_DUMMYNET, "Packet lost randomly (roll=%u, threshold=%u)",
roll, dir->loss_permille);
return;
}
}
/* Выделяем пакет */
struct ll_entry* entry = queue_entry_new_from_pool(dn->pkt_pool);
if (!entry) {
DEBUG_ERROR(DEBUG_CATEGORY_DUMMYNET, "Failed to allocate packet entry");
return;
}
struct dummynet_pkt* pkt = (struct dummynet_pkt*)entry->data;
pkt->len = n;
memcpy(pkt->data, buf, n);
pkt->src_addr = src_addr;
pkt->src_len = src_len;
pkt->recv_time_tb = get_time_tb();
pkt->direction = dir_idx;
/* Вычисляем задержку */
uint32_t delay_ms = dir->delay_fixed_ms + random_range(dir->delay_random_ms);
if (delay_ms > DUMMYNET_MAX_DELAY_MS) {
delay_ms = DUMMYNET_MAX_DELAY_MS;
}
int delay_tb = MS_TO_TB(delay_ms);
DEBUG_DEBUG(DEBUG_CATEGORY_DUMMYNET, "Dir %d: packet recv (%zd bytes), delay %u ms",
dir_idx, n, delay_ms);
/* Создаём таймер задержки */
struct delay_cb_arg* targ = (struct delay_cb_arg*)u_malloc(sizeof(struct delay_cb_arg));
if (!targ) {
DEBUG_ERROR(DEBUG_CATEGORY_DUMMYNET, "Failed to allocate timer arg");
dir->stats.dropped++;
queue_entry_free(entry);
return;
}
targ->dn = dn;
targ->entry = entry;
void* timer = uasync_set_timeout(dn->ua, delay_tb, targ, dummynet_delay_callback, "dummynet_delay");
if (!timer) {
DEBUG_ERROR(DEBUG_CATEGORY_DUMMYNET, "Failed to set delay timer");
dir->stats.dropped++;
u_free(targ);
queue_entry_free(entry);
return;
}
dir->delay_timer_set++;
}
/* Callback таймера задержки */
static void dummynet_delay_callback(void* user_arg) {
struct delay_cb_arg* arg = (struct delay_cb_arg*)user_arg;
struct dummynet* dn = arg->dn;
struct ll_entry* entry = arg->entry;
u_free(arg);
struct dummynet_pkt* pkt = (struct dummynet_pkt*)entry->data;
int dir_idx = pkt->direction;
struct dummynet_dir* dir = &dn->dirs[dir_idx];
dir->delay_timer_fire++;
/* Проверяем размер очереди */
int queue_count = queue_entry_count(dir->queue);
if ((uint32_t)queue_count >= dir->max_queue_pkts) {
dir->stats.dropped++;
// DEBUG_WARN(DEBUG_CATEGORY_DUMMYNET, "Dir %d: queue full (%d packets), dropping packet", dir_idx, queue_count);
queue_entry_free(entry);
return;
}
/* Добавляем в очередь */
uint32_t id = (uint32_t)(uintptr_t)entry;
if (queue_data_put(dir->queue, entry) != 0) {
dir->stats.dropped++;
DEBUG_WARN(DEBUG_CATEGORY_DUMMYNET, "Dir %d: queue_data_put failed", dir_idx);
queue_entry_free(entry);
return;
}
dir->stats.queue_size = queue_count + 1;
if ((uint32_t)(queue_count + 1) > dir->stats.queue_max) {
dir->stats.queue_max = queue_count + 1;
}
DEBUG_DEBUG(DEBUG_CATEGORY_DUMMYNET, "Dir %d: packet added to queue (size=%d)",
dir_idx, queue_count + 1);
/* Запускаем шейпер если не активен */
if (!dir->shaper_pending) {
dir->shaper_pending = 1;
struct shaper_cb_arg* sarg = (struct shaper_cb_arg*)u_malloc(sizeof(struct shaper_cb_arg));
if (sarg) {
sarg->dn = dn;
sarg->dir_idx = dir_idx;
int shaper_delay_tb = 0;
if (dir->bandwidth_kbps > 0) {
uint64_t delay_calc = (uint64_t)pkt->len * 80ULL / (uint64_t)dir->bandwidth_kbps;
shaper_delay_tb = (int)delay_calc;
if (shaper_delay_tb == 0) shaper_delay_tb = 1;
}
dir->shaper_timer = uasync_set_timeout(dn->ua, shaper_delay_tb, sarg, dummynet_shaper_callback, "dummynet_shaper0");
if (dir->shaper_timer) dir->shaper_timer_set++;
else { u_free(sarg); dir->shaper_pending = 0; }
} else {
dir->shaper_pending = 0;
}
}
}
/* Создаёт dummynet контекст */
struct dummynet* dummynet_create(struct UASYNC* ua, const char* bind_ip, uint16_t listen_port) {
if (!ua || listen_port == 0 || listen_port < 2 || listen_port > 65533) {
DEBUG_ERROR(DEBUG_CATEGORY_DUMMYNET, "Invalid parameters");
return NULL;
}
struct dummynet* dn = (struct dummynet*)u_calloc(1, sizeof(struct dummynet));
if (!dn) {
DEBUG_ERROR(DEBUG_CATEGORY_DUMMYNET, "Failed to allocate dummynet context");
return NULL;
}
dn->ua = ua;
dn->listen_port = listen_port;
/* Создаём пул для пакетов */
dn->pkt_pool = memory_pool_init(sizeof(struct ll_entry) + sizeof(struct dummynet_pkt));
if (!dn->pkt_pool) {
DEBUG_ERROR(DEBUG_CATEGORY_DUMMYNET, "Failed to create packet pool");
u_free(dn);
return NULL;
}
/* Создаём очереди для направлений */
for (int i = 0; i < DUMMYNET_DIR_COUNT; i++) {
dn->dirs[i].queue = queue_new(ua, 0, 0, 0, "dummynet1");
if (!dn->dirs[i].queue) {
DEBUG_ERROR(DEBUG_CATEGORY_DUMMYNET, "Failed to create queue for direction %d", i);
dummynet_destroy(dn);
return NULL;
}
}
/* Создаём UDP сокет */
dn->sock = socket_create_udp(AF_INET);
if (dn->sock == SOCKET_INVALID) {
DEBUG_ERROR(DEBUG_CATEGORY_DUMMYNET, "Failed to create UDP socket");
dummynet_destroy(dn);
return NULL;
}
if (socket_set_nonblocking(dn->sock) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_DUMMYNET, "Failed to set nonblocking");
dummynet_destroy(dn);
return NULL;
}
/* Привязываем сокет */
struct sockaddr_in bind_addr;
memset(&bind_addr, 0, sizeof(bind_addr));
bind_addr.sin_family = AF_INET;
bind_addr.sin_port = htons(listen_port);
if (bind_ip && strcmp(bind_ip, "0.0.0.0") != 0) {
if (inet_pton(AF_INET, bind_ip, &bind_addr.sin_addr) != 1) {
DEBUG_ERROR(DEBUG_CATEGORY_DUMMYNET, "Invalid bind IP: %s", bind_ip);
dummynet_destroy(dn);
return NULL;
}
} else {
bind_addr.sin_addr.s_addr = INADDR_ANY;
}
if (bind(dn->sock, (struct sockaddr*)&bind_addr, sizeof(bind_addr)) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_DUMMYNET, "Failed to bind to %s:%u: %s",
bind_ip ? bind_ip : "0.0.0.0", listen_port, socket_strerror(socket_get_error()));
dummynet_destroy(dn);
return NULL;
}
/* Добавляем сокет в uasync */
void* sock_id = uasync_add_socket_t(ua, dn->sock, dummynet_read_callback, NULL, NULL, dn);
if (!sock_id) {
DEBUG_ERROR(DEBUG_CATEGORY_DUMMYNET, "Failed to add socket to uasync");
dummynet_destroy(dn);
return NULL;
}
DEBUG_INFO(DEBUG_CATEGORY_DUMMYNET, "Dummynet created on %s:%u (forward: %u->%u, backward: %u->%u)",
bind_ip ? bind_ip : "0.0.0.0", listen_port,
listen_port - 1, listen_port + 1,
listen_port + 1, listen_port - 1);
return dn;
}
/* Уничтожает dummynet контекст */
void dummynet_destroy(struct dummynet* dn) {
if (!dn) return;
DEBUG_INFO(DEBUG_CATEGORY_DUMMYNET, "Destroying dummynet");
if (dn->sock != SOCKET_INVALID) {
uasync_remove_socket_t(dn->ua, dn->sock);
socket_close_wrapper(dn->sock);
}
for (int i = 0; i < DUMMYNET_DIR_COUNT; i++) {
if (dn->dirs[i].shaper_timer) {
uasync_cancel_timeout(dn->ua, dn->dirs[i].shaper_timer);
}
if (dn->dirs[i].queue) {
struct ll_entry* entry;
while ((entry = queue_data_get(dn->dirs[i].queue)) != NULL) {
queue_entry_free(entry);
}
queue_free(dn->dirs[i].queue);
}
}
if (dn->pkt_pool) {
memory_pool_destroy(dn->pkt_pool);
}
u_free(dn);
}
/* Настраивает параметры направления */
int dummynet_set_direction(struct dummynet* dn, int direction,
uint32_t delay_fixed_ms, uint32_t delay_random_ms,
uint32_t bandwidth_kbps, uint32_t max_queue_pkts,
uint32_t loss_permille,
const char* dest_ip, uint16_t dest_port) {
if (!dn || direction < 0 || direction >= DUMMYNET_DIR_COUNT) {
return -1;
}
if (!dest_ip || dest_port == 0) {
DEBUG_ERROR(DEBUG_CATEGORY_DUMMYNET, "Invalid destination");
return -1;
}
if (loss_permille > DUMMYNET_MAX_LOSS_PERMILLE) {
DEBUG_WARN(DEBUG_CATEGORY_DUMMYNET, "Loss permille %u clamped to %d",
loss_permille, DUMMYNET_MAX_LOSS_PERMILLE);
loss_permille = DUMMYNET_MAX_LOSS_PERMILLE;
}
struct dummynet_dir* dir = &dn->dirs[direction];
dir->delay_fixed_ms = delay_fixed_ms;
dir->delay_random_ms = delay_random_ms;
dir->bandwidth_kbps = bandwidth_kbps;
dir->max_queue_pkts = max_queue_pkts;
dir->loss_permille = loss_permille;
memset(&dir->dest_addr, 0, sizeof(dir->dest_addr));
struct sockaddr_in* dest_sin = (struct sockaddr_in*)&dir->dest_addr;
dest_sin->sin_family = AF_INET;
dest_sin->sin_port = htons(dest_port);
if (inet_pton(AF_INET, dest_ip, &dest_sin->sin_addr) != 1) {
DEBUG_ERROR(DEBUG_CATEGORY_DUMMYNET, "Invalid destination IP: %s", dest_ip);
return -1;
}
dir->dest_len = sizeof(struct sockaddr_in);
dir->configured = 1;
DEBUG_INFO(DEBUG_CATEGORY_DUMMYNET, "Direction %d configured: delay=%u+%u ms, bw=%u kbps, "
"queue=%u pkts, loss=%u.%u%%, dest=%s:%u",
direction, delay_fixed_ms, delay_random_ms, bandwidth_kbps,
max_queue_pkts, loss_permille / 10, loss_permille % 10,
dest_ip, dest_port);
return 0;
}
/* Получает статистику направления */
const struct dummynet_stats* dummynet_get_stats(struct dummynet* dn, int direction) {
if (!dn || direction < 0 || direction >= DUMMYNET_DIR_COUNT) {
return NULL;
}
return &dn->dirs[direction].stats;
}
/* Сбрасывает статистику */
void dummynet_reset_stats(struct dummynet* dn, int direction) {
if (!dn) return;
if (direction < 0) {
for (int i = 0; i < DUMMYNET_DIR_COUNT; i++) {
memset(&dn->dirs[i].stats, 0, sizeof(struct dummynet_stats));
}
} else if (direction >= 0 && direction < DUMMYNET_DIR_COUNT) {
memset(&dn->dirs[direction].stats, 0, sizeof(struct dummynet_stats));
}
}
/* Возвращает файловый дескриптор сокета */
int dummynet_get_socket(struct dummynet* dn) {
if (!dn) return SOCKET_INVALID;
return dn->sock;
}
/* Возвращает listen_port */
uint16_t dummynet_get_listen_port(struct dummynet *dn) {
if (!dn) return 0;
return dn->listen_port;
}
/* Возвращает указатель на UASYNC */
struct UASYNC *dummynet_get_uasync(struct dummynet *dn) {
if (!dn) return NULL;
return dn->ua;
}
/* Возвращает размер очереди направления */
int dummynet_get_queue_size(struct dummynet *dn, int direction) {
if (!dn || direction < 0 || direction >= DUMMYNET_DIR_COUNT) {
return -1;
}
return queue_entry_count(dn->dirs[direction].queue);
}
/* Получает отладочные счетчики направления */
void dummynet_get_debug_counters(struct dummynet *dn, int direction,
uint64_t *delay_set, uint64_t *delay_fire,
uint64_t *shaper_set, uint64_t *shaper_fire) {
if (!dn || direction < 0 || direction >= DUMMYNET_DIR_COUNT) {
*delay_set = *delay_fire = *shaper_set = *shaper_fire = 0;
return;
}
struct dummynet_dir *dir = &dn->dirs[direction];
*delay_set = dir->delay_timer_set;
*delay_fire = dir->delay_timer_fire;
*shaper_set = dir->shaper_timer_set;
*shaper_fire = dir->shaper_timer_fire;
}
// === Inline filter: send_hook для etcp_udp_send ===
struct dummynet_filter {
struct UASYNC* ua;
uint32_t loss_permille;
uint32_t delay_ms;
uint32_t jitter_ms;
struct ETCP_LINK* link;
struct dummynet_stats stats;
};
static ssize_t dummynet_filter_send_hook(socket_t fd, const void* buf, size_t len,
const struct sockaddr* addr, socklen_t addr_len,
struct ETCP_LINK* link, void* ctx) {
struct dummynet_filter* df = (struct dummynet_filter*)ctx;
(void)link;
df->stats.recv++;
if (df->loss_permille > 0) {
uint32_t roll = (uint32_t)rand() % 1000;
if (roll < df->loss_permille) {
df->stats.lost++;
return (ssize_t)len;
}
}
ssize_t sent = socket_sendto(fd, buf, len, addr, addr_len);
if (sent > 0) df->stats.sent++;
return sent;
}
struct dummynet_filter* dummynet_filter_create(struct UASYNC* ua) {
struct dummynet_filter* df = u_calloc(1, sizeof(*df));
if (!df) return NULL;
df->ua = ua;
return df;
}
void dummynet_filter_set_loss(struct dummynet_filter* df, uint32_t loss_permille) {
if (!df) return;
if (loss_permille > 1000) loss_permille = 1000;
df->loss_permille = loss_permille;
}
void dummynet_filter_set_delay(struct dummynet_filter* df, uint32_t delay_ms, uint32_t jitter_ms) {
if (!df) return;
df->delay_ms = delay_ms;
df->jitter_ms = jitter_ms;
}
void dummynet_filter_attach(struct dummynet_filter* df, struct ETCP_LINK* link) {
if (!df || !link) return;
dummynet_filter_detach(df);
df->link = link;
link->send_hook = dummynet_filter_send_hook;
link->send_hook_ctx = df;
}
void dummynet_filter_detach(struct dummynet_filter* df) {
if (!df || !df->link) return;
df->link->send_hook = NULL;
df->link->send_hook_ctx = NULL;
df->link = NULL;
}
void dummynet_filter_destroy(struct dummynet_filter* df) {
if (!df) return;
dummynet_filter_detach(df);
u_free(df);
}
const struct dummynet_stats* dummynet_filter_get_stats(struct dummynet_filter* df) {
return df ? &df->stats : NULL;
}