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.
 
 
 
 
 
 

724 lines
24 KiB

// 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 */
}