diff --git a/AUDIT_FINDINGS.txt b/AUDIT_FINDINGS.txt deleted file mode 100644 index 08d3b428..00000000 --- a/AUDIT_FINDINGS.txt +++ /dev/null @@ -1,253 +0,0 @@ -======================================================================== - uTun3 AUDIT FINDINGS — Сводный отчёт -======================================================================== - - Аудит от 14.05.2026 — обход основных модулей - (use-after-free, out-of-bounds, утечки, reentrancy) - -======================================================================== -HIGH (19 проблем) -======================================================================== - -1. u_async.c:137-165 use-after-free - Частичный успех u_realloc в socket_array_add_internal: если часть - реаллоков прошла, а часть нет — старые массивы уже освобождены, - новые не выделены, sa->sockets/sa->fd_to_index/sa->index_to_fd - указывают на освобождённую память. - -2. u_async.c:842,1167 логическая ошибка - process_posted_tasks никогда не вызывается на Linux. - В epoll-пути: drain_wakeup_pipe не вызывает handle_wakeup/process_posted_tasks. - В poll-пути: wakeup-слот (i==0) обрабатывается с drain_wakeup_pipe и continue, - минуя зарегистрированный коллбэк. - Следствие: uasync_post из другого потока не доставляет задачи — posted_tasks - накапливаются бесконечно. - -3. u_async.c:1383,1398,1418 утечка памяти - При сбое создания wakeup (pipe/socket) в uasync_create не освобождаются: - ua->sockets, ua->timeout_heap, ua->timeout_pool, ua->epoll_fd. - -4. memory_pool.c:51-69 double-free - Нет детекта double-free. Повторный free создаёт цикл во freelist, - memory_pool_alloc выдаёт один и тот же блок двум потребителям. - -5. etcp.c:679-683 утечка памяти - input_queue_cb: при исчерпании inflight_pool вызывается - queue_entry_free(in_pkt) без queue_dgram_free(&in_pkt->ll). - dgram из data_pool утекает. - -6. etcp_connections.c:87,876-926 use-after-free - burst_resp_timer не отменяется в etcp_link_close. - Коллбэк burst_resp_timeout_cb срабатывает на освобождённом link. - -7. etcp_connections.c:846-854 use-after-free - stats_timer не отменяется при провале insert_link в etcp_link_new. - link_stats_timer_cb срабатывает на освобождённом link. - -8. etcp_connections.c:879-886 use-after-free - Ветка conn==NULL в etcp_link_close пропускает отмену ВСЕХ таймеров - (stats_timer, init_timer, shaper_timer, keepalive_timer, burst_resp_timer). - Любой из них сработает на освобождённом link. - -9. etcp_connections.c:1388,1531,1561 out-of-bounds read - INIT request парсинг: проверка pkt_len < 12 должна быть pkt_len < 26 - (23 байта заголовка). Пакет с длиной 12-25 пройдёт проверку, - но читает за границей расшифрованных данных. - -10. pkt_normalizer.c:33 null deref - if (!pn) pn->alloc_errors++ — разыменование NULL после u_calloc. - -11. pkt_normalizer.c:300 null deref - pn_buf_renew в цикле without проверки возврата → memcpy(NULL, ...) - при сбое memory_pool_alloc. - -12. eim_nat.c:80-93 логическая ошибка - Переполнение uint16_t порта: next_port с 65535→0, затем порты - ниже port_start выделяются (включая порт 0). Соединения нерабочие. - -13. eim_nat.c (весь файл) логическая ошибка - Нет механизма истечения/удаления NAT-записей. Таблица навсегда - исчерпывается при долгой работе, eim_nat_egress начинает отказывать. - -14. route_bgp.c:749,762-766 утечка памяти - route_bgp_process_nodeinfo: при пересоздании nodeinfo1 старая - очередь paths сохраняется в локальную переменную, но затем - перезаписывается новой queue_new() — утечка. - -15. route_bgp.c:904-907 утечка памяти - route_bgp_send_nodeinfo: при провале etcp_send не освобождаются - p (u_malloc) и e (queue_entry_new). - -16. route_ping.c:126-134 утечка памяти - route_ping_cancel_for_conn: nat_check_arg (p->arg, u_calloc) не - освобождается. В нормальном и таймаутном путях коллбэк освобождает, - но cancel-путь коллбэк не вызывает. - -17. ll_queue.c:321,439 reentrancy - Синхронный вызов коллбэка при первом queue_data_put (count==1) - без отложенного возобновления через uasync → рекурсивный - ack_timeout_check с устаревшим указателем current. - -18. secure_channel.c:379,711 потенциальный OOB write - sc_encrypt пишет в caller buffer без параметра capacity. - Контракт API опасен: вызывающий должен сам гарантировать размер буфера. - -19. secure_channel.c:432,756 потенциальный OOB write - sc_decrypt — аналогично, нет проверки capacity выходного буфера. - - -======================================================================== -MEDIUM (18 проблем) -======================================================================== - -20. u_async.c:591-639 use-after-free - Stale handle в uasync_cancel_timeout/uasync_call_soon_cancel. - Нет инвалидации (generation counter), пул может переиспользовать память. - -21. memory_pool.c:8,38,61,95 race condition - Глобальная g_total_free изменяется без блокировок, общая для всех пулов. - -22. memory_pool.c:51-69 потенциальный OOB write - Нет проверки принадлежности блока пулу. Free в чужой пул вызывает - memset с чужим object_size — частичная очистка или запись за границу. - -23. timeout_heap.c:43-45 утечка / несоответствие документации - Документация обещает free() при free_callback==NULL, код этого не делает. - -24. etcp.c:718-763 reentrancy - Вложенный ack_timeout_check (через input_send_q_cb→etcp_request_pkt→ - wait_ack_cb) обрабатывает не тот entry, т.к. outer loop использует - current, который мог быть перемещён/удалён внутренним вызовом. - -25. etcp.c:747,944 утечка / дрейф счётчиков - Непроверенный queue_data_put_with_index: при провале счётчики - (retransmissions_count, inflight_bytes) не откатываются. - -26. etcp_connections.c:1731,1763 out-of-bounds read - INIT_RESPONSE парсинг: нужно >=19 байт, проверяется только >=3. - -27. etcp_connections.c:748-753 use-after-free + бесконечный цикл - etcp_socket_remove: при link->conn==NULL remove_link пропускается, - num_channels не уменьшается, цикл перечитывает освобождённый conn->links[0]. - -28. etcp_connections.c:1736-1738 логическая ошибка - link_status перезаписывается значением remote_keepalive, игнорируя - локальный recv_keepalive. Временно помечает link UP даже когда - локальный приём мёртв. - -29. pkt_normalizer.c:43 integer overflow/underflow - frag_size = mtu - ACK_REZERV - UDP_HDR_SIZE - UDP_SC_HDR_SIZE. - При mtu < 158 беззнаковое вычитание даёт ~65379 → огромные буферы. - -30. pkt_normalizer.c:240 утечка памяти - При провале queue_data_put фрагмент и dgram не освобождаются. - -31. eim_nat.c:97,173 null deref - Нет проверки ctx->initialized/ctx->table в egress/ingress. - -32. eim_nat.c:268,293 логическая ошибка - Статические пробросы (port forward) вне port range недоступны. - -33. secure_channel.c:66 key hygiene - Нет sc_free_ctx — сессионный ключ не зануляется при уничтожении контекста. - -34. secure_channel.c:322,620 key hygiene - ECDH shared secret остаётся на стеке незанулённым после - sc_derive_session_key. - -35. secure_channel.c:352,680 nonce entropy - При постоянном провале random_bytes seed остаётся 0 (static init). - Nonce становится полностью предсказуемым (counter+timestamp). - -36. route_bgp.c:1122 утечка памяти - route_bgp_send_nat_info: etcp_send без проверки → pkt/e утекают при провале. - -37. route_bgp.c:1157 утечка памяти (аналогично #36) - route_bgp_send_nat_check_req: etcp_send без проверки. - -38. config_parser.c:1055-1066 логическая ошибка - update_config_keys дописывает [global] секцию в конец файла. - При повторном вызове получается дубликат секции, парсер может сбойнуть. - - -======================================================================== -LOW (13 проблем) -======================================================================== - -39. u_async.c:1594-1617 утечка (abort) - abort() в destroy до освобождения posted_tasks и мьютекса. - -40. u_async.c:1608 race condition - posted_tasks_head читается без posted_lock в destroy. - -41. memory_pool.c:47 truncation - size_t → uint32_t в u_calloc_impl (object_size > UINT32_MAX). - -42. memory_pool.h:16 signedness - int free_count может обернуться при двойном освобождении. - -43. timeout_heap.c:31 dead code - freed_count никогда не обновляется, timeout_heap_get_freed_count не реализован. - -44. timeout_heap.c:66-67 integer overflow - capacity*2 и sizeof(TimeoutEntry)*new_cap могут переполниться. - -45. ll_queue.c:106-127 dangling pointer - queue_free освобождает hash_table, но не чистит entry->hash_next. - -46. pkt_normalizer.c:374 десинхронизация потока - При провале ll_alloc_lldgram 2 байта заголовка уже consumed, - но фрагмент сброшен → потеря данных. - -47. eim_nat.c:105,181 unaligned access (ARM) - *(uint16_t*)(ip_data + 6) — невыровненный доступ на ARM. - -48. secure_channel.c:108-131 избыточность - CRC32 поверх CCM (аутентифицированное шифрование) избыточен. - -49. secure_channel.h:51-52 dead code - Неиспользуемые поля send_nonce[13]/recv_nonce[13]. - -50. config_parser.c:312 undefined behavior - ~0U << (32 - cidr) при cidr==0 — сдвиг на ширину типа. - -51. route_ping.c:227 integer overflow - uint16_t avg_rtt = timeout_ms * 10 — переполнение при timeout_ms > 6553. - - -======================================================================== -Файлы без ошибок -======================================================================== - -- route_node.c — чисто -- tun_if.c — чисто -- etcp_api.c — только порядок queue_entry_free/dgram_free (LOW, см. etcp.c Issue #7 старого аудита) -- etcp_loadbalancer.c — не аудирован (250 строк) - - -======================================================================== -Рекомендуемый порядок исправления -======================================================================== - -1. КРИТИЧЕСКИ ФУНКЦИОНАЛЬНЫЕ: - u_async.c #2 (posted_tasks не работают на Linux) - u_async.c #3 (утечки при сбое старта) - -2. USE-AFTER-FREE: - etcp_connections.c #6,#7,#8,#5 (висячие таймеры) - u_async.c #1 (частичный realloc) - ll_queue.c #17 (reentrancy) - -3. УТЕЧКИ ПАМЯТИ: - etcp.c #5 (dgram из data_pool) - route_bgp.c #14,#15,#36,#37 - route_ping.c #16 - -4. OOB / NULL DEREF: - etcp_connections.c #9,#26 (INIT bounds) - pkt_normalizer.c #10,#11 (null deref) - -5. ЛОГИКА: - eim_nat.c #12,#13 (порты + expiry) - memory_pool.c #4 (double-free) - config_parser.c #38 (дубликат секции) diff --git a/filelist.txt b/filelist.txt deleted file mode 100644 index 6305f11a..00000000 --- a/filelist.txt +++ /dev/null @@ -1,164 +0,0 @@ -config.h -lib\debug_config.c -lib\debug_config.h -lib\ll_queue.c -lib\ll_queue.h -lib\mem.c -lib\mem.h -lib\memory_pool.c -lib\memory_pool.h -lib\platform_compat.c -lib\platform_compat.h -lib\sha256.c -lib\sha256.h -lib\socket_compat.c -lib\socket_compat.h -lib\timeout_heap.c -lib\timeout_heap.h -lib\u_async.c -lib\u_async.h -lib\wintun.h -net_emulator\net_emulator.c -net_emulator\net_emulator.h -src\config_parser.c -src\config_parser.h -src\config_updater.c -src\config_updater.h -src\control_server.c -src\control_server.h -src\crc32.c -src\crc32.h -src\dummynet.c -src\dummynet.h -src\etcp.c -src\etcp.h -src\etcp_api.c -src\etcp_api.h -src\etcp_connections.c -src\etcp_connections.h -src\etcp_debug.c -src\etcp_debug.h -src\etcp_loadbalancer.c -src\etcp_loadbalancer.h -src\packet_dump.c -src\packet_dump.h -src\pkt_normalizer.c -src\pkt_normalizer.h -src\route_bgp.c -src\route_bgp.h -src\route_lib.c -src\route_lib.h -src\routing.c -src\routing.h -src\secure_channel.c -src\secure_channel.h -src\tun_if.c -src\tun_if.h -src\tun_linux.c -src\tun_route.c -src\tun_route.h -src\tun_windows.c -src\utun.c -src\utun_instance.c -src\utun_instance.h -tests\bench_timeout_heap.c -tests\bench_uasync_timeouts.c -tests\debug_full_test.c -tests\debug_performance.c -tests\debug_simple.c -tests\detailed_test.c -tests\simple_test.c -tests\test_bgp_route_exchange.c -tests\test_config_debug.c -tests\test_control_server.c -tests\test_control_simple.c -tests\test_crash_debug.c -tests\test_crypto.c -tests\test_debug_categories.c -tests\test_dummynet.c -tests\test_ecc_encrypt.c -tests\test_etcp_100_packets.c -tests\test_etcp_api.c -tests\test_etcp_crypto.c -tests\test_etcp_dummynet.c -tests\test_etcp_exit.c -tests\test_etcp_link_id.c -tests\test_etcp_minimal.c -tests\test_etcp_simple_traffic.c -tests\test_etcp_two_instances.c -tests\test_intensive_memory_pool.c -tests\test_intensive_memory_pool_new.c -tests\test_ll_queue.c -tests\test_memory_pool_and_config.c -tests\test_minimal.c -tests\test_minimal_exit.c -tests\test_offset.c -tests\test_packet_dump.c -tests\test_pkt_normalizer_etcp.c -tests\test_pkt_normalizer_standalone.c -tests\test_poll_exact.c -tests\test_poll_multi.c -tests\test_route_lib.c -tests\test_routing_mesh.c -tests\test_simple.c -tests\test_simple2.c -tests\test_socket.c -tests\test_utils.h -tests\test_u_async_comprehensive.c -tests\test_u_async_performance.c -tests\track_test.c -tests\working_crypto_test.c -tinycrypt\lib\include\tinycrypt\aes.h -tinycrypt\lib\include\tinycrypt\cbc_mode.h -tinycrypt\lib\include\tinycrypt\ccm_mode.h -tinycrypt\lib\include\tinycrypt\cmac_mode.h -tinycrypt\lib\include\tinycrypt\constants.h -tinycrypt\lib\include\tinycrypt\ctr_mode.h -tinycrypt\lib\include\tinycrypt\ctr_prng.h -tinycrypt\lib\include\tinycrypt\ecc.h -tinycrypt\lib\include\tinycrypt\ecc_dh.h -tinycrypt\lib\include\tinycrypt\ecc_dsa.h -tinycrypt\lib\include\tinycrypt\ecc_platform_specific.h -tinycrypt\lib\include\tinycrypt\hmac.h -tinycrypt\lib\include\tinycrypt\hmac_prng.h -tinycrypt\lib\include\tinycrypt\sha256.h -tinycrypt\lib\include\tinycrypt\utils.h -tinycrypt\lib\source\aes_decrypt.c -tinycrypt\lib\source\aes_encrypt.c -tinycrypt\lib\source\cbc_mode.c -tinycrypt\lib\source\ccm_mode.c -tinycrypt\lib\source\cmac_mode.c -tinycrypt\lib\source\ctr_mode.c -tinycrypt\lib\source\ctr_prng.c -tinycrypt\lib\source\ecc.c -tinycrypt\lib\source\ecc_dh.c -tinycrypt\lib\source\ecc_dsa.c -tinycrypt\lib\source\ecc_platform_specific.c -tinycrypt\lib\source\hmac.c -tinycrypt\lib\source\hmac_prng.c -tinycrypt\lib\source\sha256.c -tinycrypt\lib\source\utils.c -tinycrypt\tests\test_aes.c -tinycrypt\tests\test_cbc_mode.c -tinycrypt\tests\test_ccm_mode.c -tinycrypt\tests\test_client_server.c -tinycrypt\tests\test_cmac_mode.c -tinycrypt\tests\test_ctr_mode.c -tinycrypt\tests\test_ctr_prng.c -tinycrypt\tests\test_ecc_dh.c -tinycrypt\tests\test_ecc_dsa.c -tinycrypt\tests\test_ecc_utils.c -tinycrypt\tests\test_hmac.c -tinycrypt\tests\test_hmac_prng.c -tinycrypt\tests\test_sha256.c -tinycrypt\tests\include\test_ecc_utils.h -tinycrypt\tests\include\test_utils.h -tools\bping\bping.c -tools\etcpmon\etcpmon_client.c -tools\etcpmon\etcpmon_client.h -tools\etcpmon\etcpmon_graph.c -tools\etcpmon\etcpmon_graph.h -tools\etcpmon\etcpmon_gui.c -tools\etcpmon\etcpmon_gui.h -tools\etcpmon\etcpmon_main.c -tools\etcpmon\etcpmon_protocol.h diff --git a/nat.txt b/nat.txt deleted file mode 100644 index 167c8fd2..00000000 --- a/nat.txt +++ /dev/null @@ -1,219 +0,0 @@ -Диаграмма вызовов функций при определении типа NAT (NAT detection) -===================================================================== - -=== ИНИЦИАЛИЗАЦИЯ BGP === - -route_bgp_init(instance) - └── etcp_set_new_conn_cbk(instance, route_bgp_etcp_conn_cbk, NULL) - └── При создании нового ETCP_CONN: - route_bgp_etcp_conn_cbk(conn, arg) - ├── etcp_conn_set_up_cbk(conn, route_bgp_on_conn_up, bgp) - └── etcp_conn_set_down_cbk(conn, route_bgp_on_conn_down, bgp) - - -=== УСТАНОВКА СОЕДИНЕНИЯ (сервер получает INIT_REQUEST) === - -etcp_connections_read_callback_socket(sock, arg) - └── Получает ETCP_INIT_REQUEST - ├── etcp_connection_create(instance, name) [если новый peer] - │ └── Вызывает instance->etcp_new_conn_cbk [route_bgp_etcp_conn_cbk] - ├── etcp_link_new(etcp, e_sock, remote_addr, is_server=1) - ├── Отправляет ETCP_INIT_RESPONSE (с NAT info) - ├── link->initialized = 1 - └── etcp_conn_ready(link->etcp) - └── conn->ready_cbk - └── etcp_on_up(etcp) - └── etcp->up_cbk = route_bgp_on_conn_up(etcp, bgp) - - -=== УСТАНОВКА СОЕДИНЕНИЯ (клиент получает INIT_RESPONSE) === - -etcp_connections_read_callback_socket(sock, arg) - └── Получает ETCP_INIT_RESPONSE - ├── Сохраняет NAT address: link->nat_ip, link->nat_port - ├── link->initialized = 1 - └── etcp_conn_ready(link->etcp) - └── etcp_on_up(etcp) - └── etcp->up_cbk = route_bgp_on_conn_up(etcp, bgp) - - -=== ЗАПУСК ПРОВЕРКИ NAT (при поднятии соединения) === - -route_bgp_on_conn_up(etcp, arg) - └── route_bgp_new_conn(conn) - ├── route_bgp_add_to_senders(bgp, conn) - └── Сканирует ВСЕ линки ВСЕХ соединений: - if (link->initialized && link->conn && link->nat_check_status < NAT_CHECK_IN_PROGRESS) - └── route_bgp_start_link_nat_check(bgp, link) - ├── if (nat_check_status == IN_PROGRESS) return - ├── route_bgp_find_third_node(bgp, exclude=link->etcp) - │ └── Ищет другой ETCP_CONN в senders_list - ├── Определяет target_ip/target_port: - │ ├── Если link->nat_ip != 0: использует nat_ip/nat_port - │ └── Иначе: использует link->remote_addr - ├── Проверяет is_local_subnet(target_ip) [если запрещено - skip] - ├── route_bgp_extract_nat_addr(conn, &nat_ip, &nat_port) - │ └── Ищет первый линк с nat_ip != 0 - ├── Создает nat_check_arg { link, nat_ip, nat_port } - └── route_ping_send_req_addr(bgp, third_conn, target_ip, target_port, - count=3, interval=500ms, timeout=1000ms, - wait_timeout=5000ms, - cb=nat_link_check_cb, arg=nat_check_arg, - pubkey=peer_pubkey) - - -=== ОТПРАВКА ЗАПРОСА НА ПИНГ (через third node) === - -route_ping_send_req_addr(bgp, to_conn, target_ip, target_port, ...) - ├── Создает BGP_PING_REQUEST { cmd, subcmd=PING_REQ, request_id, count, interval, timeout, target_ipv4, target_port, pubkey } - ├── etcp_send(to_conn, entry) [отправляет запрос third node] - ├── Создает route_ping_pending { request_id, callback=nat_link_check_cb, arg, timeout_timer } - └── Добавляет в bgp->ping_pending - - -=== ОБРАБОТКА PING_REQ (на third node) === - -route_bgp_receive_cbk(data, len, from_conn, bgp) - └── ROUTE_SUBCMD_PING_REQ - └── route_ping_handle_req(bgp, from_conn, data, len) - ├── Создает route_ping_series_ctx { reply_conn, request_id, target_addr, pubkey, local_sock, count_total } - ├── Находит первый IPv4 ETCP_SOCKET -> local_sock - └── etcp_send_ping_to_socket(instance, local_sock, pubkey, &target_addr, - timeout_ms, route_ping_single_cb, ctx, NULL, 0) - └── Отправляет ETCP пинг на target_addr - - -=== ОТВЕТ НА ПИНГ (целевой узел) === - -etcp_connections_read_callback_socket() - └── Получает ETCP пинг -> отправляет pong - - -=== ПОЛУЧЕНИЕ ОТВЕТА НА ПИНГ (third node) === - -route_ping_single_cb(success, rtt, arg, nonce, resp_data, resp_data_len) - ├── Обновляет статистику: count_sent++, count_ok++, sum_rtt += rtt - ├── Если count_sent < count_total: - │ └── etcp_send_ping_to_socket(...) [следующий пинг] - └── Иначе: - └── route_ping_series_finish(ctx) - ├── Вычисляет avg_rtt - └── Создает BGP_PING_RESPONSE { cmd, subcmd=PING_RESP, request_id, count_sent, count_ok, avg_rtt } - └── etcp_send(reply_conn, entry) [отправляет результат инициатору] - - -=== ПОЛУЧЕНИЕ PING_RESP (на инициаторе) === - -route_bgp_receive_cbk(data, len, from_conn, bgp) - └── ROUTE_SUBCMD_PING_RESP - └── route_ping_handle_resp(bgp, from_conn, data, len) - ├── Ищет route_ping_pending по request_id - ├── Отменяет timeout_timer - └── Вызывает callback: nat_link_check_cb(success, avg_rtt, count_sent, count_ok, arg) - - -=== CALLBACK ПРОВЕРКИ NAT === - -nat_link_check_cb(success, avg_rtt, count_sent, count_ok, arg) - ├── Определяет nat_type = success ? NAT_TYPE_OPEN : NAT_TYPE_RESTRICTED - ├── link->nat_type = nat_type - ├── link->nat_check_status = success ? NAT_CHECK_OPEN : NAT_CHECK_RESTRICTED - └── route_bgp_send_nat_info(link->etcp, socket_id, nat_ip, nat_port, nat_type) - ├── Создает BGP_NAT_INFO { cmd, subcmd=NAT_INFO, socket_id, nat_ip[4], nat_port, nat_type } - └── etcp_send(conn, entry) [отправляет результат клиенту] - - -=== ПРИЕМ NAT_INFO (на клиенте) === - -route_bgp_receive_cbk(data, len, from_conn, bgp) - └── ROUTE_SUBCMD_NAT_INFO - └── route_bgp_handle_nat_info(bgp, from_conn, data, len) - ├── Ищет линки по remote_socket_id - ├── Устанавливает link->nat_type = info->nat_type - └── Логирует результат NAT detection - - -=== ЗАПРОС НА ПРОВЕРКУ NAT (NAT_CHECK_REQ - новый механизм) === - -route_bgp_send_nat_check_req(conn, socket_id) - ├── Создает BGP_NAT_CHECK_REQ { cmd, subcmd=NAT_CHECK_REQ, socket_id, interface_ip, interface_port } - │ ├── Ищет линк с remote_socket_id == socket_id - │ └── Берет interface_addr из линка -> interface_ip/interface_port - └── etcp_send(conn, entry) - -route_bgp_handle_nat_check_req(bgp, from_conn, data, len) - ├── Принимает BGP_NAT_CHECK_REQ - ├── Сохраняет interface_ip/interface_port в struct sockaddr_storage dummy - ├── Ищет линк с remote_socket_id == req->socket_id - └── Если линк найден и не в процессе проверки: - └── route_bgp_start_link_nat_check(bgp, target_link) - [далее по стандартному flow] - - -=== СВЯЗЬ СТРУКТУР === - -struct ETCP_SOCKET - ├── local_addr [адрес бинда из конфига] - ├── interface_addr [IP интерфейса + порт: интерфейс/default route/конфиг] - └── ... - -struct ETCP_LINK - ├── remote_addr [адрес удаленного узла] - ├── nat_ip [NAT IP, полученный из INIT_RESPONSE] - ├── nat_port [NAT port, полученный из INIT_RESPONSE] - ├── nat_type [OPEN/RESTRICTED, результат проверки] - ├── nat_check_status [NONE/WAITING/IN_PROGRESS/OPEN/RESTRICTED] - ├── initialized [0/1, флаг завершения handshake] - └── ... - -struct ETCP_CONN - ├── links [список ETCP_LINK] - ├── up_cbk [route_bgp_on_conn_up] - ├── down_cbk [route_bgp_on_conn_down] - ├── ready_cbk [вызывается при etcp_conn_ready] - ├── initialized [0/1, флаг готовности] - ├── links_up [0/1, есть ли активные линки] - └── ... - - -=== СТАТУСЫ NAT CHECK === - -NAT_CHECK_NONE = 0 [не проверялся] -NAT_CHECK_WAITING = 1 [ожидает] -NAT_CHECK_IN_PROGRESS = 2 [проверка запущена] -NAT_CHECK_OPEN = 3 [результат: OPEN] -NAT_CHECK_RESTRICTED = 4 [результат: RESTRICTED] - -NAT_TYPE_UNKNOWN = 0 -NAT_TYPE_OPEN = 1 -NAT_TYPE_RESTRICTED = 2 - -NAT_VERIFIED_UNKNOWN = 4 -NAT_VERIFIED_OPEN = 5 -NAT_VERIFIED_RESTRICTED = 6 -NAT_VERIFIED_DIRECT = 7 - - -=== СООБЩЕНИЯ ПРОТОКОЛА === - -ETCP_INIT_REQUEST (0x02) [клиент -> сервер] -ETCP_INIT_RESPONSE (0x03) [сервер -> клиент, включает NAT info] -ETCP_INIT_REQUEST_NOINIT (0x04) -ETCP_INIT_RESPONSE_NOINIT (0x05) -ETCP_PING (0x06) -ETCP_PONG (0x07) - -ROUTE_SUBCMD_NODEINFO (0x04) -ROUTE_SUBCMD_REQUEST_TABLE (0x05) -ROUTE_SUBCMD_WITHDRAW (0x06) -ROUTE_SUBCMD_PING_REQ (0x07) -ROUTE_SUBCMD_PING_RESP (0x08) -ROUTE_SUBCMD_NAT_INFO (0x09) -ROUTE_SUBCMD_NAT_CHECK_REQ (0x0A) - - -=== ПОЛНЫЙ FLOW (одна строка) === - -Handshake: etcp_connections_read_callback_socket -> INIT_REQUEST/RESPONSE -> link->initialized=1 -> etcp_conn_ready -> etcp_on_up -> route_bgp_on_conn_up -> route_bgp_new_conn -> route_bgp_start_link_nat_check -> route_ping_send_req_addr -> (через ETCP) -> route_ping_handle_req -> etcp_send_ping_to_socket -> (UDP) -> route_ping_single_cb -> route_ping_series_finish -> route_ping_handle_resp -> nat_link_check_cb -> route_bgp_send_nat_info -> route_bgp_handle_nat_info - -NAT_CHECK_REQ flow: route_bgp_send_nat_check_req -> BGP_NAT_CHECK_REQ -> route_bgp_handle_nat_check_req -> route_bgp_start_link_nat_check -> [стандартный ping flow] diff --git a/src/Makefile.am b/src/Makefile.am index 9d20ce6a..38f7f096 100644 --- a/src/Makefile.am +++ b/src/Makefile.am @@ -1,4 +1,5 @@ bin_PROGRAMS = utun +noinst_LIBRARIES = libutun.a # Sources that are always compiled utun_CORE_SOURCES = \ @@ -14,6 +15,7 @@ utun_CORE_SOURCES = \ route_node_lmdb.c \ route_connectivity.c \ conn_mgr.c \ + db_sync.c \ routing.c \ tun_if.c \ tun_route.c \ @@ -53,10 +55,64 @@ utun_CORE_SOURCES = \ lwip_tcp/lwip_tcp_in.c \ lwip_tcp/lwip_tcp_out.c +# libutun: all core sources except main() +libutun_a_SOURCES = \ + utun_instance.c \ + config_parser.c \ + config_updater.c \ + route_lib.c \ + route6_lib.c \ + route_bgp.c \ + route_ping.c \ + route_node.c \ + route_node_lmdb.c \ + route_connectivity.c \ + conn_mgr.c \ + db_sync.c \ + routing.c \ + tun_if.c \ + tun_route.c \ + tun_linux.c \ + tun_freebsd.c \ + tun_windows.c \ + etcp.c \ + etcp_connections.c \ + etcp_bbr.c \ + etcp_loadbalancer.c \ + etcp_debug.c \ + etcp_dump.c \ + secure_channel.c \ + crc32.c \ + stcp_link.c \ + stcp.c \ + stcp_server.c \ + stcp_client.c \ + pkt_normalizer.c \ + packet_dump.c \ + etcp_api.c \ + etcp_connect.c \ + control_server.c \ + msg_transport.c \ + firewall.c \ + eim_nat.c \ + nat_transport.c \ + dummynet.c \ + proxy/tcp_proxy_client.c \ + etcp_router.c \ + proxy/tcp_proxy_server.c \ + proxy/udp_proxy.c \ + proxy/socks_proxy.c \ + proxy/icmp_proxy.c \ + lwip_tcp/lwip_pbuf.c \ + lwip_tcp/lwip_tcp.c \ + lwip_tcp/lwip_tcp_in.c \ + lwip_tcp/lwip_tcp_out.c +libutun_a_CFLAGS = $(utun_CFLAGS) + # Platform-specific TUN libs (Windows only) utun_TUN_LIBS = @TUN_LIBS@ -utun_SOURCES = $(utun_CORE_SOURCES) $(utun_TUN_SOURCES) +utun_SOURCES = utun.c # Include paths utun_CFLAGS = \ @@ -67,6 +123,7 @@ utun_CFLAGS = \ # Libraries utun_LDADD = \ + libutun.a \ $(top_builddir)/lib/libuasync.a \ -lpthread \ -lm \ diff --git a/src/_db_arch.txt b/src/_db_arch.txt new file mode 100644 index 00000000..ccc6b83f --- /dev/null +++ b/src/_db_arch.txt @@ -0,0 +1,22 @@ +Архитектура таблицы c быстрой репликацией между несколькими узлами: + +1. В конец таблицы каждый узел может самостоятельно добавлять данные (со своим timestamp, желательно время правильное) +2. узлы между собой синхронизируются, распространяя обновления по соседям + +3. TTL: изменение имеет время жизни. узел удаляет собственные записи которые не были отправлены никому в течении суток + +4. Формат записи: + +ID (64bit monotonic autoincrement) +chain_hash (256bit) - цепочка хешей. считается так: chain_hash[index+1]=hash(chain_hash[index],ID[index],timestamp[index],datahash[index]). если хеш совпадает - это признак того что все строки выше синхронизированы. +timestamp (64bit) - время создания записи. строки должны сортироваться по возрастанию concat(datahash(младшая часть числа,64 bit)+timestamp(старшая часть)) +datahash (64bit) - хеш данных этой строки +data (varchar) - данные в формате json + +алгоритм репликации (2 узла): +узел A запрашивает синхронизацию и передает свой last id +test_id= min(mast id, peer last id) +узел B передает: test_id, chain_hash(test id), hash(test id-1), hash(test id-2), hash(test id-4), hash(test id-8) итд 16,32,..., до первого элемента включитально (лимитируем вылетевший индекс первым элементом) +узел A - сравнивает, находит проверяемый диапазон id, отправляет другому узлу 16 (можно больше) хешей (линейно разбив проверяемый диапазон на более короткие поддиапазоны). таким образом узлы уточняют первый ид который не совпал. +Когда первая различающияся запись найдена, узел отправляет хеш этой и n (например 32) последующих записей (если записей много). если записей мало (<4) то узел отправляет сразу содержимое записей. +(надо додумать алгоритм, чтобы оптимизировать количество итераций - лучше передать больше данных за раз чем много итераций с ожиданием ответной стороны) diff --git a/src/config_parser.c b/src/config_parser.c index 2875f26f..55b27b52 100644 --- a/src/config_parser.c +++ b/src/config_parser.c @@ -399,6 +399,14 @@ static int parse_global(const char *key, const char *value, struct global_config if (strcmp(key, "db_path") == 0) { return assign_string(global->db_path, sizeof(global->db_path), value); } + if (strcmp(key, "db_sync_enabled") == 0) { + global->db_sync_enabled = atoi(value); + return 0; + } + if (strcmp(key, "db_sync_ttl") == 0) { + global->db_sync_ttl = (uint32_t)atoi(value); + return 0; + } /* debug_categories key is deprecated - use [debug] section instead */ if (strcmp(key, "enable_timestamp") == 0) { global->enable_timestamp = atoi(value); @@ -792,6 +800,8 @@ static struct utun_config* parse_config_internal(FILE *fp, const char *filename) cfg->global.firewall_bypass_all = 0; cfg->global.control_allows = NULL; cfg->global.control_allow_count = 0; + cfg->global.db_sync_enabled = 0; + cfg->global.db_sync_ttl = 86400; section_type_t cur_section = SECTION_UNKNOWN; struct CFG_SERVER *cur_server = NULL; diff --git a/src/config_parser.h b/src/config_parser.h index 6ef955ee..653cb299 100644 --- a/src/config_parser.h +++ b/src/config_parser.h @@ -115,6 +115,8 @@ struct global_config { // Debug and logging configuration char log_file[256]; // Path to log file (empty = stdout) char db_path[256]; // Path to LMDB nodeinfo database (empty = disabled) + int db_sync_enabled; // 1 = enable distributed DB sync (default: 0) + uint32_t db_sync_ttl; // TTL for unsent records in seconds (default: 86400) char debug_level[16]; // debug level: error, warn, info, debug, trace int enable_timestamp; // enable timestamps in logs int enable_function_names; // enable function names in logs diff --git a/src/db_sync.c b/src/db_sync.c new file mode 100644 index 00000000..c55cf22e --- /dev/null +++ b/src/db_sync.c @@ -0,0 +1,940 @@ +// db_sync.c — Distributed content-addressed table with LMDB + peer sync via etcp_router + +#include "db_sync.h" +#include "etcp_api.h" +#include "etcp.h" +#include "etcp_router.h" +#include "utun_instance.h" +#include "route_bgp.h" +#include "../lib/debug_config.h" +#include "../lib/mem.h" +#include "../lib/u_async.h" +#include "../lib/sha256.h" +#include "../lib/liblmdb/lmdb.h" +#include "../lib/platform_compat.h" +#include + +#define DEBUG_CATEGORY_DB_SYNC DEBUG_CATEGORY_DEBUG + +// ---- LMDB key/value layout ---- +// Key: 16 bytes = timestamp:8 BE || datahash:8 BE (sorted by timestamp first) +// Value: id:8 || chain_hash:32 || creator:8 || flags:1 || data_len:4 || data:variable +#define DB_VAL_OFF_ID 0 +#define DB_VAL_OFF_CHAIN_HASH 8 +#define DB_VAL_OFF_CREATOR 40 +#define DB_VAL_OFF_FLAGS 48 +#define DB_VAL_OFF_DATA_LEN 49 +#define DB_VAL_OFF_DATA 53 +#define DB_VAL_HDR_SIZE 53 + +// ---- Module state ---- +struct DB_SYNC { + struct UTUN_INSTANCE* inst; + MDB_env* env; + MDB_dbi dbi_records; + MDB_dbi dbi_meta; + uint64_t next_id; + uint64_t last_connected_tb; // timebase when any peer was last connected (for TTL) + uint64_t last_timestamp_us; // гарантия монотонности timestamp + + struct DB_SYNC_PEER { + uint64_t node_id; + uint8_t synced; // 0=not, 1=in_progress, 2=synced + }* peers; + int peer_count; + int peer_capacity; + + void* ttl_timer; + void* peer_check_timer; + uint8_t enabled; +}; + + +// ---- Forward declarations ---- +static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry); +static void db_sync_on_new_conn(struct ETCP_CONN* conn, void* arg); +static void db_sync_on_conn_up(struct ETCP_CONN* conn, void* arg); +static void db_sync_on_conn_down(struct ETCP_CONN* conn, void* arg); +static void db_sync_peer_check_cb(void* arg); +static void db_sync_ttl_cleanup_cb(void* arg); +static int db_sync_send(struct DB_SYNC* db, uint64_t dst_node_id, const uint8_t* payload, size_t len); +static void db_sync_initiate_sync(struct DB_SYNC* db, uint64_t peer_node_id); + +// ---- LMDB helpers ---- +static int db_lmdb_open(struct DB_SYNC* db, const char* path) { + int rc = mdb_env_create(&db->env); + if (rc) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync: mdb_env_create failed: %s", mdb_strerror(rc)); return -1; } + rc = mdb_env_set_mapsize(db->env, DB_SYNC_DEFAULT_MAPSIZE); + if (rc) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync: mdb_env_set_mapsize failed: %s", mdb_strerror(rc)); mdb_env_close(db->env); db->env=NULL; return -1; } + rc = mdb_env_set_maxdbs(db->env, 3); + if (rc) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync: mdb_env_set_maxdbs failed: %s", mdb_strerror(rc)); mdb_env_close(db->env); db->env=NULL; return -1; } + rc = mdb_env_open(db->env, path, 0, 0644); + if (rc) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync: mdb_env_open(%s) failed: %s", path, mdb_strerror(rc)); mdb_env_close(db->env); db->env=NULL; return -1; } + + MDB_txn* txn; + rc = mdb_txn_begin(db->env, NULL, 0, &txn); + if (rc) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync: txn_begin failed: %s", mdb_strerror(rc)); mdb_env_close(db->env); db->env=NULL; return -1; } + rc = mdb_dbi_open(txn, "records", MDB_CREATE, &db->dbi_records); + if (rc) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync: dbi_open(records) failed: %s", mdb_strerror(rc)); mdb_txn_abort(txn); mdb_env_close(db->env); db->env=NULL; return -1; } + rc = mdb_dbi_open(txn, "meta", MDB_CREATE, &db->dbi_meta); + if (rc) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync: dbi_open(meta) failed: %s", mdb_strerror(rc)); mdb_txn_abort(txn); mdb_env_close(db->env); db->env=NULL; return -1; } + rc = mdb_txn_commit(txn); + if (rc) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync: txn_commit failed: %s", mdb_strerror(rc)); mdb_env_close(db->env); db->env=NULL; return -1; } + + // Read next_id from meta + MDB_txn* rt; + rc = mdb_txn_begin(db->env, NULL, MDB_RDONLY, &rt); + if (rc == MDB_SUCCESS) { + MDB_val key, data; key.mv_size = 7; key.mv_data = (char*)"next_id"; + rc = mdb_get(rt, db->dbi_meta, &key, &data); + if (rc == MDB_SUCCESS && data.mv_size >= 8) db->next_id = *(uint64_t*)data.mv_data; + else db->next_id = 1; + mdb_txn_abort(rt); + } + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: LMDB opened at %s, next_id=%llu", path, (unsigned long long)db->next_id); + return 0; +} + +static void db_lmdb_close(struct DB_SYNC* db) { + if (!db->env) return; + mdb_env_close(db->env); db->env = NULL; + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: LMDB closed"); +} + +// Build 16-byte key: timestamp:8 BE || datahash:8 BE +static void db_key_build(uint8_t key[16], uint64_t timestamp, uint64_t datahash) { + uint64_t ts_be = htobe64(timestamp), dh_be = htobe64(datahash); + memcpy(key, &ts_be, 8); memcpy(key + 8, &dh_be, 8); +} + +// Get chain_hash at a given sorted position (0-based) via cursor. Returns 0 on success. +static int db_chain_hash_at(struct DB_SYNC* db, uint32_t pos, uint8_t chain_hash_out[32]) { + MDB_txn* txn; int rc = mdb_txn_begin(db->env, NULL, MDB_RDONLY, &txn); + if (rc) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync: chain_hash_at txn_begin failed: %s", mdb_strerror(rc)); return -1; } + MDB_cursor* cursor; rc = mdb_cursor_open(txn, db->dbi_records, &cursor); + if (rc) { mdb_txn_abort(txn); return -1; } + MDB_val key, val; + rc = mdb_cursor_get(cursor, &key, &val, MDB_FIRST); + for (uint32_t i = 0; rc == 0 && i < pos; i++) rc = mdb_cursor_get(cursor, &key, &val, MDB_NEXT); + if (rc == 0 && val.mv_size >= DB_VAL_HDR_SIZE) memcpy(chain_hash_out, (uint8_t*)val.mv_data + DB_VAL_OFF_CHAIN_HASH, 32); + else { mdb_cursor_close(cursor); mdb_txn_abort(txn); return -1; } + mdb_cursor_close(cursor); mdb_txn_abort(txn); + return 0; +} + +// Get count of records +static uint32_t db_count(struct DB_SYNC* db) { + MDB_txn* txn; int rc = mdb_txn_begin(db->env, NULL, MDB_RDONLY, &txn); + if (rc) return 0; + MDB_stat stat; rc = mdb_stat(txn, db->dbi_records, &stat); + mdb_txn_abort(txn); + return rc == 0 ? (uint32_t)stat.ms_entries : 0; +} + +// Get datahash at given sorted position +static int db_datahash_at(struct DB_SYNC* db, uint32_t pos, uint64_t* datahash) { + MDB_txn* txn; int rc = mdb_txn_begin(db->env, NULL, MDB_RDONLY, &txn); + if (rc) return -1; + MDB_cursor* cursor; rc = mdb_cursor_open(txn, db->dbi_records, &cursor); + if (rc) { mdb_txn_abort(txn); return -1; } + MDB_val key, val; rc = mdb_cursor_get(cursor, &key, &val, MDB_FIRST); + for (uint32_t i=0; rc==0 && ienv, NULL, MDB_RDONLY, &txn); + if (rc) return -1; + MDB_cursor* cursor; rc = mdb_cursor_open(txn, db->dbi_records, &cursor); + if (rc) { mdb_txn_abort(txn); return -1; } + MDB_val search_key, val; search_key.mv_size = 16; search_key.mv_data = (void*)key; + rc = mdb_cursor_get(cursor, &search_key, &val, MDB_SET_RANGE); + if (rc == 0) { + rc = mdb_cursor_get(cursor, &search_key, &val, MDB_PREV); + if (rc == 0 && val.mv_size >= DB_VAL_HDR_SIZE) { + memcpy(chain_hash, (uint8_t*)val.mv_data + DB_VAL_OFF_CHAIN_HASH, 32); + mdb_cursor_close(cursor); mdb_txn_abort(txn); return 0; + } + } + // No previous record → use zero hash + memset(chain_hash, 0, 32); + mdb_cursor_close(cursor); mdb_txn_abort(txn); + return 0; +} + +// ---- Insert record ---- +static int db_record_insert(struct DB_SYNC* db, uint64_t id, uint64_t timestamp, uint64_t datahash, const char* json, size_t json_len) { + uint8_t key[16]; db_key_build(key, timestamp, datahash); + + // Check duplicate + MDB_txn* txn; int rc = mdb_txn_begin(db->env, NULL, 0, &txn); + if (rc) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync: insert txn_begin failed: %s", mdb_strerror(rc)); return -1; } + MDB_val dup_key, dup_val; dup_key.mv_size = 16; dup_key.mv_data = key; + rc = mdb_get(txn, db->dbi_records, &dup_key, &dup_val); + if (rc == MDB_SUCCESS) { mdb_txn_abort(txn); return 1; } // already exists + + // Compute chain_hash + uint8_t prev_ch[32]; db_prev_chain_hash(db, key, prev_ch); + uint8_t chain_hash[32]; db_chain_hash_compute(prev_ch, id, timestamp, datahash, chain_hash); + + // Build value + size_t val_size = DB_VAL_HDR_SIZE + json_len; + uint8_t* val_buf = u_malloc(val_size); + if (!val_buf) { mdb_txn_abort(txn); return -1; } + uint64_t creator = db->inst->node_id; + memcpy(val_buf + DB_VAL_OFF_ID, &id, 8); + memcpy(val_buf + DB_VAL_OFF_CHAIN_HASH, chain_hash, 32); + memcpy(val_buf + DB_VAL_OFF_CREATOR, &creator, 8); + val_buf[DB_VAL_OFF_FLAGS] = 0; + *(uint32_t*)(val_buf + DB_VAL_OFF_DATA_LEN) = (uint32_t)json_len; + if (json_len > 0) memcpy(val_buf + DB_VAL_OFF_DATA, json, json_len); + + MDB_val val; val.mv_size = val_size; val.mv_data = val_buf; + MDB_val k; k.mv_size = 16; k.mv_data = key; + rc = mdb_put(txn, db->dbi_records, &k, &val, 0); + u_free(val_buf); + + if (rc) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync: mdb_put failed: %s", mdb_strerror(rc)); mdb_txn_abort(txn); return -1; } + + // Update chain_hash for subsequent records (out-of-order insert case) + // Check if there are records AFTER this key + MDB_cursor* cursor; rc = mdb_cursor_open(txn, db->dbi_records, &cursor); + if (rc == 0) { + MDB_val ckey, cval; ckey.mv_size = 16; ckey.mv_data = key; + rc = mdb_cursor_get(cursor, &ckey, &cval, MDB_SET); + if (rc == 0) { + uint8_t running_ch[32]; memcpy(running_ch, chain_hash, 32); + uint64_t running_id = id; + uint64_t running_ts = timestamp; + uint64_t running_dh = datahash; + // Move to next + while (mdb_cursor_get(cursor, &ckey, &cval, MDB_NEXT) == 0) { + if (cval.mv_size < DB_VAL_HDR_SIZE) continue; + uint64_t n_id = *(uint64_t*)((uint8_t*)cval.mv_data + DB_VAL_OFF_ID); + uint64_t n_ts = be64toh(*(uint64_t*)ckey.mv_data); + uint64_t n_dh = be64toh(*(uint64_t*)((uint8_t*)ckey.mv_data + 8)); + uint8_t new_ch[32]; db_chain_hash_compute(running_ch, n_id, n_ts, n_dh, new_ch); + memcpy((uint8_t*)cval.mv_data + DB_VAL_OFF_CHAIN_HASH, new_ch, 32); + mdb_cursor_put(cursor, &ckey, &cval, MDB_CURRENT); + memcpy(running_ch, new_ch, 32); + } + } + mdb_cursor_close(cursor); + } + + // Update next_id in meta + MDB_val mk, md; mk.mv_size = 7; mk.mv_data = (char*)"next_id"; md.mv_size = 8; md.mv_data = &id; id++; + mdb_put(txn, db->dbi_meta, &mk, &md, 0); + db->next_id = id; + + rc = mdb_txn_commit(txn); + if (rc) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync: insert commit failed: %s", mdb_strerror(rc)); return -1; } + return 0; +} + +// ---- Send ---- +static int db_sync_send(struct DB_SYNC* db, uint64_t dst_node_id, const uint8_t* payload, size_t len) { + struct ll_entry* entry = queue_entry_new(0); + if (!entry) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync_send: queue_entry_new failed"); return -1; } + uint8_t* buf = u_malloc(len + 1); + if (!buf) { queue_entry_free(entry); return -1; } + buf[0] = ETCP_RT_ID_DB_SYNC; memcpy(buf + 1, payload, len); + entry->dgram = buf; entry->len = len + 1; + int ret = etcp_route_send(db->inst, dst_node_id, entry, 0); + // etcp_route_send always takes ownership and frees entry in all code paths + return ret; +} + +// ---- Peer management ---- +static struct DB_SYNC_PEER* db_peer_find(struct DB_SYNC* db, uint64_t node_id) { + for (int i=0; ipeer_count; i++) if (db->peers[i].node_id == node_id) return &db->peers[i]; + return NULL; +} +static struct DB_SYNC_PEER* db_peer_add(struct DB_SYNC* db, uint64_t node_id) { + struct DB_SYNC_PEER* p = db_peer_find(db, node_id); + if (p) return p; + if (db->peer_count >= db->peer_capacity) { + int nc = db->peer_capacity ? db->peer_capacity * 2 : 8; + struct DB_SYNC_PEER* np = u_realloc(db->peers, nc * sizeof(struct DB_SYNC_PEER)); + if (!np) return NULL; + db->peers = np; db->peer_capacity = nc; + } + p = &db->peers[db->peer_count++]; + p->node_id = node_id; p->synced = 0; + return p; +} + +// ---- Sync protocol: INIT_SYNC (A→B) ---- +static void db_handle_init_sync(struct DB_SYNC* db, uint64_t src_node_id, const uint8_t* payload, size_t len) { + if (len < 36) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync: INIT_SYNC too short %zu from %016llx", len, (unsigned long long)src_node_id); return; } + uint32_t peer_count = *(uint32_t*)payload; + const uint8_t* peer_last_ch = payload + 4; + + uint32_t my_count = db_count(db); + uint32_t test_pos = (peer_count < my_count ? peer_count : my_count); + if (test_pos > 0) test_pos--; + + // Send INIT_RESP + uint8_t resp[4096]; + uint32_t resp_off = 0; + resp[resp_off++] = DB_MSG_INIT_RESP; + memcpy(resp+resp_off, &test_pos, 4); resp_off += 4; + + uint8_t my_ch[32]; + if (my_count == 0) { + memset(my_ch, 0, 32); + } else if (db_chain_hash_at(db, test_pos, my_ch) != 0) { + memset(my_ch, 0, 32); + } + memcpy(resp+resp_off, my_ch, 32); resp_off += 32; + + // Always include sparse_count at position 37 (after test_pos + chain_hash) + // For complete case: sparse_count=0, no sparse entries follow + // For sparse case: sparse entries start at position 38 + + // Check if chain_hash matches + if (memcmp(my_ch, peer_last_ch, 32) == 0) { + resp[resp_off++] = 0; // sparse_count = 0 + db_sync_send(db, src_node_id, resp, resp_off); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: INIT_RESP complete to %016llx test_pos=%u", (unsigned long long)src_node_id, test_pos); + return; + } + + // Skip sparse_count slot for now, write it at the end + uint32_t sparse_count_pos = resp_off; + resp_off++; // placeholder for sparse_count + + // Build sparse hashes: chain_hash at positions test_pos-1, -2, -4, -8, -16, ... + int sparse_count = 0; + for (uint32_t k = 0; k < 16; k++) { + uint32_t step = (uint32_t)(1u << k); // 1, 2, 4, 8, ... + if (test_pos < step) break; + uint32_t pos = test_pos - step; + if (db_chain_hash_at(db, pos, my_ch) != 0) break; + if (resp_off + 36 > sizeof(resp)) break; + memcpy(resp+resp_off, &pos, 4); resp_off += 4; + memcpy(resp+resp_off, my_ch, 32); resp_off += 32; + sparse_count++; + } + resp[sparse_count_pos] = (uint8_t)sparse_count; + + db_sync_send(db, src_node_id, resp, resp_off); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: INIT_RESP to %016llx test_pos=%u sparse=%d", (unsigned long long)src_node_id, test_pos, sparse_count); +} + +// ---- Sync protocol: INIT_RESP (B→A) ---- +static void db_handle_init_resp(struct DB_SYNC* db, uint64_t src_node_id, const uint8_t* payload, size_t len) { + if (len < 37) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync: INIT_RESP too short %zu from %016llx", len, (unsigned long long)src_node_id); return; } + uint32_t test_pos = *(uint32_t*)payload; + const uint8_t* peer_ch = payload + 4; + uint8_t sparse_count = payload[36]; + + uint8_t my_ch[32]; + if (db_chain_hash_at(db, test_pos, my_ch) != 0) { memset(my_ch, 0, 32); } + + // Check if complete + if (sparse_count == 0 && memcmp(my_ch, peer_ch, 32) == 0) { + struct DB_SYNC_PEER* p = db_peer_find(db, src_node_id); + if (p) p->synced = 2; + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: sync complete with %016llx", (unsigned long long)src_node_id); + return; + } + + // If peer has empty table (chain_hash is all zeros), send everything + { + uint8_t zero[32]; memset(zero, 0, 32); + if (memcmp(peer_ch, zero, 32) == 0 && sparse_count == 0) { + uint32_t my_count = db_count(db); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: peer %016llx has empty table, sending all %u records", (unsigned long long)src_node_id, my_count); + // Send all records in batches + uint32_t sent = 0; + while (sent < my_count) { + uint32_t batch = my_count - sent; + if (batch > DB_SEND_DATA_MAX) batch = DB_SEND_DATA_MAX; + uint8_t sdbuf[8192]; uint32_t off=0; + sdbuf[off++] = DB_MSG_SEND_DATA; memcpy(sdbuf+off, &sent, 4); off+=4; + uint16_t* rcp = (uint16_t*)(sdbuf+off); off+=2; uint16_t rc=0; + MDB_txn* txn; int rrc = mdb_txn_begin(db->env, NULL, MDB_RDONLY, &txn); + if (rrc==0) { MDB_cursor* cursor; rrc=mdb_cursor_open(txn, db->dbi_records, &cursor); + if (rrc==0) { MDB_val key,val; rrc=mdb_cursor_get(cursor, &key, &val, MDB_FIRST); + uint32_t cur=0; while (rrc==0 && cur < sent) { rrc=mdb_cursor_get(cursor,&key,&val,MDB_NEXT); cur++; } + while (rrc==0 && rc < batch && off+val.mv_size+16 < 8000) { + uint64_t ts = be64toh(*(uint64_t*)key.mv_data); uint64_t dh = be64toh(*(uint64_t*)((uint8_t*)key.mv_data+8)); + uint64_t rid = *(uint64_t*)((uint8_t*)val.mv_data+DB_VAL_OFF_ID); uint32_t dlen = *(uint32_t*)((uint8_t*)val.mv_data+DB_VAL_OFF_DATA_LEN); + memcpy(sdbuf+off, &rid,8); off+=8; memcpy(sdbuf+off, &ts,8); off+=8; memcpy(sdbuf+off, &dh,8); off+=8; + memcpy(sdbuf+off, &dlen,4); off+=4; if (dlen>0) { memcpy(sdbuf+off, (uint8_t*)val.mv_data+DB_VAL_OFF_DATA, dlen); off+=dlen; } + rc++; cur++; rrc=mdb_cursor_get(cursor, &key, &val, MDB_NEXT); } + mdb_cursor_close(cursor); } mdb_txn_abort(txn); } + *rcp = rc; sent += rc; + db_sync_send(db, src_node_id, sdbuf, off); + if (rc == 0) break; // no more records + } + return; + } + } + + if (memcmp(my_ch, peer_ch, 32) == 0) { + struct DB_SYNC_PEER* p = db_peer_find(db, src_node_id); + if (p) p->synced = 2; + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: chain_hash match with %016llx at pos %u — synced", (unsigned long long)src_node_id, test_pos); + // Send our extra records if any (beyond test_pos) + uint32_t my_count = db_count(db); + if (my_count > test_pos + 1) { + // We have more records — send them + uint32_t from = test_pos + 1; + uint32_t remaining = my_count - from; + uint32_t batch = remaining < DB_SEND_DATA_MAX ? remaining : DB_SEND_DATA_MAX; + // Continue in SEND_DATA handler... + } + return; + } + + // Find divergence range using sparse hashes + uint32_t div_start = 0, div_end = test_pos; + const uint8_t* sparse = payload + 37; + for (int i = 0; i < sparse_count && sparse + 36 <= payload + len; i++) { + uint32_t pos = *(uint32_t*)sparse; + const uint8_t* peer_sparse_ch = sparse + 4; + uint8_t my_sparse_ch[32]; + if (db_chain_hash_at(db, pos, my_sparse_ch) == 0) { + if (memcmp(my_sparse_ch, peer_sparse_ch, 32) == 0) { if (pos + 1 > div_start) div_start = pos + 1; } + else { if (pos < div_end) div_end = pos; } + } + sparse += 36; + } + + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: divergence with %016llx range [%u, %u]", (unsigned long long)src_node_id, div_start, div_end); + + if (div_end - div_start <= 1) { + // Divergence is small — request data directly via REFINE round + uint8_t ref[512]; uint32_t ref_off = 0; + ref[ref_off++] = DB_MSG_REFINE; + memcpy(ref+ref_off, &div_start, 4); ref_off += 4; + memcpy(ref+ref_off, &div_end, 4); ref_off += 4; + uint8_t hcount = 0; ref[ref_off++] = hcount; // no hashes, just requesting data + db_sync_send(db, src_node_id, ref, ref_off); + } else { + // Send REFINE with our hashes + uint8_t ref[512]; uint32_t ref_off = 0; + ref[ref_off++] = DB_MSG_REFINE; + memcpy(ref+ref_off, &div_start, 4); ref_off += 4; + memcpy(ref+ref_off, &div_end, 4); ref_off += 4; + uint32_t range = div_end - div_start; + uint8_t count = range < DB_REFINE_HASHES ? (uint8_t)range : DB_REFINE_HASHES; + ref[ref_off++] = count; + for (uint8_t i = 0; i < count; i++) { + uint32_t pos = div_start + (range * i / count); + uint64_t dh; + if (db_datahash_at(db, pos, &dh) == 0) { + memcpy(ref+ref_off, &pos, 4); ref_off += 4; + memcpy(ref+ref_off, &dh, 8); ref_off += 8; + } + } + db_sync_send(db, src_node_id, ref, ref_off); + } +} + +// ---- Sync protocol: REFINE ---- +static void db_handle_refine(struct DB_SYNC* db, uint64_t src_node_id, const uint8_t* payload, size_t len) { + if (len < 9) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync: REFINE too short %zu", len); return; } + uint32_t from = *(uint32_t*)payload; + uint32_t to = *(uint32_t*)(payload+4); + uint8_t hcount = payload[8]; + + if (hcount == 0) { + // Peer requests data in range [from, to] + uint32_t my_count = db_count(db); + uint32_t send_from = from; + uint32_t send_count = (to - from + 1) < DB_SEND_DATA_MAX ? (to - from + 1) : DB_SEND_DATA_MAX; + if (send_from >= my_count) return; // nothing to send + + // Build SEND_DATA + uint8_t* buf = u_malloc(8192); + if (!buf) return; + uint32_t off = 0; buf[off++] = DB_MSG_SEND_DATA; + memcpy(buf+off, &send_from, 4); off += 4; + uint16_t rec_count = 0; + uint16_t* cnt_ptr = (uint16_t*)(buf+off); off += 2; + + MDB_txn* txn; int rc = mdb_txn_begin(db->env, NULL, MDB_RDONLY, &txn); + if (rc == 0) { + MDB_cursor* cursor; rc = mdb_cursor_open(txn, db->dbi_records, &cursor); + if (rc == 0) { + MDB_val key, val; rc = mdb_cursor_get(cursor, &key, &val, MDB_FIRST); + uint32_t cur = 0; + while (rc == 0 && cur < send_from) { rc = mdb_cursor_get(cursor, &key, &val, MDB_NEXT); cur++; } + while (rc == 0 && rec_count < send_count && off + val.mv_size + 16 < 8000) { + uint64_t ts = be64toh(*(uint64_t*)key.mv_data); + uint64_t dh = be64toh(*(uint64_t*)((uint8_t*)key.mv_data + 8)); + uint64_t rid = *(uint64_t*)((uint8_t*)val.mv_data + DB_VAL_OFF_ID); + uint32_t dlen = *(uint32_t*)((uint8_t*)val.mv_data + DB_VAL_OFF_DATA_LEN); + const uint8_t* dptr = (uint8_t*)val.mv_data + DB_VAL_OFF_DATA; + memcpy(buf+off, &rid, 8); off += 8; + memcpy(buf+off, &ts, 8); off += 8; + memcpy(buf+off, &dh, 8); off += 8; + memcpy(buf+off, &dlen, 4); off += 4; + if (dlen > 0) { memcpy(buf+off, dptr, dlen); off += dlen; } + rec_count++; + rc = mdb_cursor_get(cursor, &key, &val, MDB_NEXT); + } + mdb_cursor_close(cursor); + } + mdb_txn_abort(txn); + } + *cnt_ptr = rec_count; + db_sync_send(db, src_node_id, buf, off); + u_free(buf); + return; + } + + // Peer sent hashes at positions — find mismatch + uint32_t first_mismatch = to + 1; + const uint8_t* hp = payload + 9; + for (uint8_t i = 0; i < hcount && hp + 12 <= payload + len; i++) { + uint32_t pos = *(uint32_t*)hp; + uint64_t peer_dh = *(uint64_t*)(hp + 4); + uint64_t my_dh; + if (db_datahash_at(db, pos, &my_dh) == 0 && my_dh == peer_dh) { + if (pos >= from) { if (pos + 1 < first_mismatch) first_mismatch = pos + 1; } + } else { + if (pos < first_mismatch) first_mismatch = pos; + } + hp += 12; + } + + // Send SEND_DATA from first_mismatch + uint32_t my_count = db_count(db); + uint32_t send_count = 4; + if (first_mismatch > to || first_mismatch >= my_count) { send_count = 0; } + + uint8_t sdbuf[8192]; uint32_t off = 0; + sdbuf[off++] = DB_MSG_SEND_DATA; + memcpy(sdbuf+off, &first_mismatch, 4); off += 4; + uint16_t rc_count = 0; + uint16_t* rcp = (uint16_t*)(sdbuf+off); off += 2; + + if (send_count > 0) { + MDB_txn* txn; int rrc = mdb_txn_begin(db->env, NULL, MDB_RDONLY, &txn); + if (rrc == 0) { + MDB_cursor* cursor; rrc = mdb_cursor_open(txn, db->dbi_records, &cursor); + if (rrc == 0) { + MDB_val key, val; rrc = mdb_cursor_get(cursor, &key, &val, MDB_FIRST); + uint32_t cur = 0; + while (rrc == 0 && cur < first_mismatch) { rrc = mdb_cursor_get(cursor, &key, &val, MDB_NEXT); cur++; } + while (rrc == 0 && rc_count < send_count && off + val.mv_size + 16 < 8000) { + uint64_t ts = be64toh(*(uint64_t*)key.mv_data); + uint64_t dh = be64toh(*(uint64_t*)((uint8_t*)key.mv_data + 8)); + uint64_t rid = *(uint64_t*)((uint8_t*)val.mv_data + DB_VAL_OFF_ID); + uint32_t dlen = *(uint32_t*)((uint8_t*)val.mv_data + DB_VAL_OFF_DATA_LEN); + const uint8_t* dptr = (uint8_t*)val.mv_data + DB_VAL_OFF_DATA; + memcpy(sdbuf+off, &rid, 8); off += 8; + memcpy(sdbuf+off, &ts, 8); off += 8; + memcpy(sdbuf+off, &dh, 8); off += 8; + memcpy(sdbuf+off, &dlen, 4); off += 4; + if (dlen > 0) { memcpy(sdbuf+off, dptr, dlen); off += dlen; } + rc_count++; rrc = mdb_cursor_get(cursor, &key, &val, MDB_NEXT); + } + mdb_cursor_close(cursor); + } + mdb_txn_abort(txn); + } + } + *rcp = rc_count; + db_sync_send(db, src_node_id, sdbuf, off); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: SEND_DATA to %016llx from=%u count=%u", (unsigned long long)src_node_id, first_mismatch, rc_count); +} + +// ---- Sync protocol: SEND_DATA ---- +static void db_handle_send_data(struct DB_SYNC* db, uint64_t src_node_id, const uint8_t* payload, size_t len) { + if (len < 6) return; + uint32_t from = *(uint32_t*)payload; + uint16_t count = *(uint16_t*)(payload+4); + const uint8_t* ptr = payload + 6; + + for (uint16_t i = 0; i < count; i++) { + if (ptr + 28 > payload + len) break; + uint64_t rid = *(uint64_t*)ptr; ptr += 8; + uint64_t ts = *(uint64_t*)ptr; ptr += 8; + uint64_t dh = *(uint64_t*)ptr; ptr += 8; + uint32_t dlen = *(uint32_t*)ptr; ptr += 4; + if (ptr + dlen > payload + len) break; + db_record_insert(db, rid, ts, dh, (const char*)ptr, dlen); + ptr += dlen; + } + + // Send our extra records too (bidirectional) + uint32_t my_count = db_count(db); + struct DB_SYNC_PEER* p = db_peer_find(db, src_node_id); + uint32_t peer_known = from + count; + if (my_count > peer_known && p && p->synced == 1) { + // Send our records beyond what peer knows + uint32_t send_count = my_count - peer_known; + if (send_count > DB_SEND_DATA_MAX) send_count = DB_SEND_DATA_MAX; + uint8_t sdbuf[8192]; uint32_t off=0; + sdbuf[off++] = DB_MSG_SEND_DATA; memcpy(sdbuf+off, &peer_known, 4); off+=4; + uint16_t* rcp = (uint16_t*)(sdbuf+off); off+=2; + uint16_t rc=0; + MDB_txn* txn; int rrc = mdb_txn_begin(db->env, NULL, MDB_RDONLY, &txn); + if (rrc==0) { + MDB_cursor* cursor; rrc = mdb_cursor_open(txn, db->dbi_records, &cursor); + if (rrc==0) { + MDB_val key,val; rrc = mdb_cursor_get(cursor, &key, &val, MDB_FIRST); + uint32_t cur=0; + while (rrc==0 && cur < peer_known) { rrc=mdb_cursor_get(cursor,&key,&val,MDB_NEXT); cur++; } + while (rrc==0 && rc0) { memcpy(sdbuf+off, rdptr, rdlen); off+=rdlen; } + rc++; rrc=mdb_cursor_get(cursor, &key, &val, MDB_NEXT); + } + mdb_cursor_close(cursor); + } + mdb_txn_abort(txn); + } + *rcp = rc; + db_sync_send(db, src_node_id, sdbuf, off); + } + + // Verify sync complete + uint32_t new_count = db_count(db); + uint8_t my_last_ch[32]; + if (new_count > 0 && db_chain_hash_at(db, new_count - 1, my_last_ch) == 0) { + uint8_t sd[37]; sd[0] = DB_MSG_SYNC_DONE; + memcpy(sd+1, &new_count, 4); memcpy(sd+5, my_last_ch, 32); + db_sync_send(db, src_node_id, sd, 37); + } + if (p) p->synced = 2; + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: received %u records from %016llx", count, (unsigned long long)src_node_id); +} + +// ---- Sync protocol: SYNC_DONE ---- +static void db_handle_sync_done(struct DB_SYNC* db, uint64_t src_node_id, const uint8_t* payload, size_t len) { + if (len < 36) return; + uint32_t peer_count = *(uint32_t*)payload; + const uint8_t* peer_ch = payload + 4; + uint32_t my_count = db_count(db); + uint8_t my_ch[32]; + if (my_count > 0) db_chain_hash_at(db, my_count - 1, my_ch); + else memset(my_ch, 0, 32); + if (my_count != peer_count || memcmp(my_ch, peer_ch, 32) != 0) { + DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "db_sync: SYNC_DONE mismatch with %016llx my=%u peer=%u — re-syncing", + (unsigned long long)src_node_id, my_count, peer_count); + db_sync_initiate_sync(db, src_node_id); + return; + } + struct DB_SYNC_PEER* p = db_peer_find(db, src_node_id); + if (p) p->synced = 2; + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: sync confirmed with %016llx count=%u", (unsigned long long)src_node_id, my_count); +} + +// ---- Sync protocol: ACK_PUSH ---- +static void db_handle_ack_push(struct DB_SYNC* db, uint64_t src_node_id, const uint8_t* payload, size_t len) { + (void)src_node_id; + if (len < 16) return; + uint64_t dh = *(uint64_t*)payload; + uint64_t ts = *(uint64_t*)(payload + 8); + uint8_t key[16]; db_key_build(key, ts, dh); + MDB_txn* txn; int rc = mdb_txn_begin(db->env, NULL, 0, &txn); + if (rc) return; + MDB_val k, val; k.mv_size = 16; k.mv_data = key; + rc = mdb_get(txn, db->dbi_records, &k, &val); + if (rc == 0 && val.mv_size >= DB_VAL_HDR_SIZE) { + uint64_t creator = *(uint64_t*)((uint8_t*)val.mv_data + DB_VAL_OFF_CREATOR); + uint8_t flags = *(uint8_t*)((uint8_t*)val.mv_data + DB_VAL_OFF_FLAGS); + if (creator == db->inst->node_id && !(flags & DB_REC_FLAG_WAS_SENT)) { + flags |= DB_REC_FLAG_WAS_SENT; + *(uint8_t*)((uint8_t*)val.mv_data + DB_VAL_OFF_FLAGS) = flags; + mdb_put(txn, db->dbi_records, &k, &val, 0); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: ACK_PUSH mark sent dh=%016llx ts=%llu from %016llx", (unsigned long long)dh, (unsigned long long)ts, (unsigned long long)src_node_id); + } + } + mdb_txn_commit(txn); +} + +// ---- Sync protocol: PUSH ---- +static void db_handle_push(struct DB_SYNC* db, uint64_t src_node_id, const uint8_t* payload, size_t len) { + if (len < 28) return; + const uint8_t* ptr = payload; + uint64_t rid = *(uint64_t*)ptr; ptr += 8; + uint64_t ts = *(uint64_t*)ptr; ptr += 8; + uint64_t dh = *(uint64_t*)ptr; ptr += 8; + uint32_t dlen = *(uint32_t*)ptr; ptr += 4; + if (ptr + dlen > payload + len) return; + int ret = db_record_insert(db, rid, ts, dh, (const char*)ptr, dlen); + if (ret == 0) { + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: PUSH inserted from %016llx id=%llu dh=%016llx", (unsigned long long)src_node_id, (unsigned long long)rid, (unsigned long long)dh); + uint8_t ack[17]; ack[0] = DB_MSG_ACK_PUSH; memcpy(ack+1, &dh, 8); memcpy(ack+9, &ts, 8); + db_sync_send(db, src_node_id, ack, 17); + } +} + +// ---- Main receive callback ---- +static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { + if (!entry || entry->len < 2) { if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } return; } + struct DB_SYNC* db = conn ? conn->instance->db_sync : NULL; + if (!db || !db->enabled) { queue_dgram_free(entry); queue_entry_free(entry); return; } + + uint64_t src_node_id = conn ? conn->peer_node_id : 0; + + uint8_t type = entry->dgram[1]; + const uint8_t* payload = entry->dgram + 2; + size_t plen = entry->len - 2; + + switch (type) { + case DB_MSG_INIT_SYNC: db_handle_init_sync(db, src_node_id, payload, plen); break; + case DB_MSG_INIT_RESP: db_handle_init_resp(db, src_node_id, payload, plen); break; + case DB_MSG_REFINE: db_handle_refine(db, src_node_id, payload, plen); break; + case DB_MSG_SEND_DATA: db_handle_send_data(db, src_node_id, payload, plen); break; + case DB_MSG_PUSH: db_handle_push(db, src_node_id, payload, plen); break; + case DB_MSG_ACK_PUSH: db_handle_ack_push(db, src_node_id, payload, plen); break; + case DB_MSG_SYNC_DONE: db_handle_sync_done(db, src_node_id, payload, plen); break; + default: DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "db_sync: unknown msg type 0x%02x from %016llx", type, (unsigned long long)src_node_id); break; + } + queue_dgram_free(entry); queue_entry_free(entry); +} + +// ---- Connection callbacks ---- +static void db_sync_on_new_conn(struct ETCP_CONN* conn, void* arg) { + (void)arg; + if (!conn || !conn->instance || !conn->instance->db_sync) return; + etcp_conn_add_up_cbk(conn, db_sync_on_conn_up, NULL); + etcp_conn_add_down_cbk(conn, db_sync_on_conn_down, NULL); +} + +static void db_sync_on_conn_up(struct ETCP_CONN* conn, void* arg) { + (void)arg; + if (!conn || !conn->instance || !conn->instance->db_sync) return; + struct DB_SYNC* db = conn->instance->db_sync; + uint64_t peer_id = conn->peer_node_id; + if (peer_id == 0 || peer_id == db->inst->node_id) return; + db->last_connected_tb = get_time_tb(); + struct DB_SYNC_PEER* p = db_peer_add(db, peer_id); + if (p && p->synced == 0) { + p->synced = 1; + db_sync_initiate_sync(db, peer_id); + } + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: peer up %016llx", (unsigned long long)peer_id); +} + +static void db_sync_on_conn_down(struct ETCP_CONN* conn, void* arg) { + (void)arg; + if (!conn || !conn->instance || !conn->instance->db_sync) return; + struct DB_SYNC* db = conn->instance->db_sync; + uint64_t peer_id = conn->peer_node_id; + if (peer_id == 0) return; + struct DB_SYNC_PEER* p = db_peer_find(db, peer_id); + if (p) p->synced = 0; + int any_up = 0; + for (int i = 0; i < db->peer_count; i++) { if (db->peers[i].synced >= 1) { any_up = 1; break; } } + if (!any_up) db->last_connected_tb = 0; + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: peer down %016llx", (unsigned long long)peer_id); +} + +// ---- Initiate sync ---- +static void db_sync_initiate_sync(struct DB_SYNC* db, uint64_t peer_node_id) { + uint32_t my_count = db_count(db); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: initiate_sync to %016llx my_count=%u", (unsigned long long)peer_node_id, my_count); + uint8_t msg[37]; msg[0] = DB_MSG_INIT_SYNC; + memcpy(msg+1, &my_count, 4); + if (my_count > 0) { + if (db_chain_hash_at(db, my_count-1, msg+5) != 0) memset(msg+5, 0, 32); + } else { + memset(msg+5, 0, 32); + } + db_sync_send(db, peer_node_id, msg, 37); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: INIT_SYNC → %016llx my_count=%u", (unsigned long long)peer_node_id, my_count); +} + +// ---- Periodic timers ---- +static void db_sync_peer_check_cb(void* arg) { + struct DB_SYNC* db = (struct DB_SYNC*)arg; + if (!db || !db->enabled) return; + struct ROUTE_BGP* bgp = db->inst->bgp; + if (!bgp) { db->peer_check_timer = uasync_set_timeout(db->inst->ua, DB_SYNC_PEER_CHECK_INTERVAL * 10000u, db, db_sync_peer_check_cb, "db_sync_peer"); return; } + + // Check all known peers with active connections + struct ll_entry* e = bgp->senders_list->head; + while (e) { + struct ROUTE_BGP_CONN_ITEM* item = (struct ROUTE_BGP_CONN_ITEM*)e->data; + if (item->conn && item->conn->peer_node_id != 0 && item->conn->links_up) { + uint64_t pid = item->conn->peer_node_id; + if (pid != db->inst->node_id) { + struct DB_SYNC_PEER* p = db_peer_add(db, pid); + if (p && p->synced == 0) { p->synced = 1; db_sync_initiate_sync(db, pid); } + } + } + e = e->next; + } + db->peer_check_timer = uasync_set_timeout(db->inst->ua, DB_SYNC_PEER_CHECK_INTERVAL * 10000u, db, db_sync_peer_check_cb, "db_sync_peer"); +} + +static void db_sync_ttl_cleanup_cb(void* arg) { + struct DB_SYNC* db = (struct DB_SYNC*)arg; + if (!db || !db->enabled || !db->env) { if (db) db->ttl_timer = uasync_set_timeout(db->inst->ua, DB_SYNC_TTL_INTERVAL * 10000u, db, db_sync_ttl_cleanup_cb, "db_sync_ttl"); return; } + uint64_t my_id = db->inst->node_id; + uint64_t now_us = get_time_us(); + uint64_t ttl_us = (uint64_t)(db->inst->config->global.db_sync_ttl) * 1000000uLL; + uint64_t cutoff = now_us - ttl_us; + + MDB_txn* txn; int rc = mdb_txn_begin(db->env, NULL, 0, &txn); + if (rc != 0) { db->ttl_timer = uasync_set_timeout(db->inst->ua, DB_SYNC_TTL_INTERVAL * 10000u, db, db_sync_ttl_cleanup_cb, "db_sync_ttl"); return; } + + MDB_cursor* cursor; rc = mdb_cursor_open(txn, db->dbi_records, &cursor); + if (rc == 0) { + MDB_val key, val; rc = mdb_cursor_get(cursor, &key, &val, MDB_FIRST); + while (rc == 0) { + if (val.mv_size >= DB_VAL_HDR_SIZE) { + uint64_t creator = *(uint64_t*)((uint8_t*)val.mv_data + DB_VAL_OFF_CREATOR); + uint8_t flags = *(uint8_t*)((uint8_t*)val.mv_data + DB_VAL_OFF_FLAGS); + uint64_t ts = be64toh(*(uint64_t*)key.mv_data); + if (creator == my_id && !(flags & DB_REC_FLAG_WAS_SENT) && ts < cutoff) { + mdb_cursor_del(cursor, 0); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: TTL deleted unsent record ts=%llu dh=%016llx", (unsigned long long)ts, (unsigned long long)be64toh(*(uint64_t*)((uint8_t*)key.mv_data+8))); + } + } + rc = mdb_cursor_get(cursor, &key, &val, MDB_NEXT); + } + mdb_cursor_close(cursor); + } + mdb_txn_commit(txn); + db->ttl_timer = uasync_set_timeout(db->inst->ua, DB_SYNC_TTL_INTERVAL * 10000u, db, db_sync_ttl_cleanup_cb, "db_sync_ttl"); +} + +// ---- Public API ---- +int db_sync_init(struct UTUN_INSTANCE* inst) { + if (!inst) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync_init: NULL instance"); return -1; } + struct DB_SYNC* db = u_calloc(1, sizeof(struct DB_SYNC)); + if (!db) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync_init: u_calloc failed"); return -1; } + db->inst = inst; db->last_connected_tb = 0; db->last_timestamp_us = 0; + if (!inst->config->global.db_sync_enabled) { + db->enabled = 0; inst->db_sync = db; + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: disabled by config"); + return 0; + } + db->enabled = 1; + inst->db_sync = db; + + // Open LMDB + const char* db_path = inst->config->global.db_path; + char sync_path[512]; + if (db_path[0]) { + snprintf(sync_path, sizeof(sync_path), "%s/sync", db_path); + } else { + snprintf(sync_path, sizeof(sync_path), "/tmp/utun_db_sync"); + } + if (db_lmdb_open(db, sync_path) != 0) { + DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "db_sync_init: LMDB open failed, sync disabled"); + db->enabled = 0; + } + + // Register in etcp_router + etcp_router_bind(inst, ETCP_RT_ID_DB_SYNC, db_sync_recv_cb); + + // Subscribe to ETCP connection events + etcp_add_new_conn_cbk(inst, db_sync_on_new_conn, NULL); + + // For existing connections: add up/down callbacks + { + struct ETCP_CONN* conn = inst->connections; + while (conn) { + etcp_conn_add_up_cbk(conn, db_sync_on_conn_up, NULL); + etcp_conn_add_down_cbk(conn, db_sync_on_conn_down, NULL); + conn = conn->next; + } + } + + // Start periodic timers + db->peer_check_timer = uasync_set_timeout(inst->ua, DB_SYNC_PEER_CHECK_INTERVAL * 10000u, db, db_sync_peer_check_cb, "db_sync_peer"); + db->ttl_timer = uasync_set_timeout(inst->ua, DB_SYNC_TTL_INTERVAL * 10000u, db, db_sync_ttl_cleanup_cb, "db_sync_ttl"); + + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: initialized (enabled=%d)", db->enabled); + return 0; +} + +void db_sync_destroy(struct UTUN_INSTANCE* inst) { + if (!inst || !inst->db_sync) return; + struct DB_SYNC* db = inst->db_sync; + inst->db_sync = NULL; + + etcp_router_unbind(inst, ETCP_RT_ID_DB_SYNC); + + if (db->peer_check_timer) { uasync_cancel_timeout(inst->ua, db->peer_check_timer); db->peer_check_timer = NULL; } + if (db->ttl_timer) { uasync_cancel_timeout(inst->ua, db->ttl_timer); db->ttl_timer = NULL; } + + // Remove up/down callbacks from existing connections + { struct ETCP_CONN* conn = inst->connections; while (conn) { etcp_conn_remove_up_cbk(conn, db_sync_on_conn_up, NULL); etcp_conn_remove_down_cbk(conn, db_sync_on_conn_down, NULL); conn = conn->next; } } + + db_lmdb_close(db); + if (db->peers) u_free(db->peers); + u_free(db); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: destroyed"); +} + +int db_sync_insert_len(struct UTUN_INSTANCE* inst, const char* json_data, size_t len) { + if (!inst || !inst->db_sync || !json_data || len == 0) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync_insert: invalid args"); return -1; } + struct DB_SYNC* db = inst->db_sync; + if (!db->enabled) return -1; + + // Compute datahash and timestamp + uint64_t datahash = db_datahash((const uint8_t*)json_data, len); + uint64_t now_us = get_time_us(); + // Guarantee monotonicity + if (now_us <= db->last_timestamp_us) now_us = db->last_timestamp_us + 1; + db->last_timestamp_us = now_us; + uint64_t timestamp = now_us; + uint64_t id = db->next_id; + + int ret = db_record_insert(db, id, timestamp, datahash, json_data, len); + if (ret != 0) return ret; // 1 = duplicate, -1 = error + + // Broadcast PUSH to all synced peers + uint8_t pbuf[2048]; + uint32_t off = 0; pbuf[off++] = DB_MSG_PUSH; + memcpy(pbuf+off, &id, 8); off += 8; + memcpy(pbuf+off, ×tamp, 8); off += 8; + memcpy(pbuf+off, &datahash, 8); off += 8; + memcpy(pbuf+off, &len, 4); off += 4; + if (len > 0 && off + len <= sizeof(pbuf)) { memcpy(pbuf+off, json_data, len); off += len; } + for (int i = 0; i < db->peer_count; i++) if (db->peers[i].synced >= 1 && db->peers[i].node_id != inst->node_id) db_sync_send(db, db->peers[i].node_id, pbuf, off); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: insert id=%llu dh=%016llx ts=%llu len=%zu", (unsigned long long)id, (unsigned long long)datahash, (unsigned long long)timestamp, len); + return 0; +} + +int db_sync_insert(struct UTUN_INSTANCE* inst, const char* json_data) { + if (!json_data) return -1; + return db_sync_insert_len(inst, json_data, strlen(json_data)); +} + +uint32_t db_sync_count(struct UTUN_INSTANCE* inst) { + if (!inst || !inst->db_sync) return 0; + return db_count(inst->db_sync); +} diff --git a/src/db_sync.h b/src/db_sync.h new file mode 100644 index 00000000..7b23aaa7 --- /dev/null +++ b/src/db_sync.h @@ -0,0 +1,41 @@ +// db_sync.h — Distributed content-addressed table with LMDB storage and peer sync via etcp_router +#ifndef DB_SYNC_H +#define DB_SYNC_H + +#include +#include + +struct UTUN_INSTANCE; + +// etcp_router service ID +#define ETCP_RT_ID_DB_SYNC 0x20 + +// Message types +#define DB_MSG_INIT_SYNC 0x01 +#define DB_MSG_INIT_RESP 0x02 +#define DB_MSG_REFINE 0x03 +#define DB_MSG_SEND_DATA 0x04 +#define DB_MSG_PUSH 0x05 +#define DB_MSG_ACK_PUSH 0x06 +#define DB_MSG_SYNC_DONE 0x07 + +// Defaults +#define DB_SYNC_DEFAULT_TTL 86400 +#define DB_SYNC_DEFAULT_MAPSIZE (100UL * 1024 * 1024) +#define DB_SYNC_PEER_CHECK_INTERVAL 5 +#define DB_SYNC_TTL_INTERVAL 3600 + +// Record flags +#define DB_REC_FLAG_WAS_SENT 0x01 + +// Sync protocol constants +#define DB_REFINE_HASHES 16 +#define DB_SEND_DATA_MAX 32 + +int db_sync_init(struct UTUN_INSTANCE* inst); +void db_sync_destroy(struct UTUN_INSTANCE* inst); +int db_sync_insert(struct UTUN_INSTANCE* inst, const char* json_data); +int db_sync_insert_len(struct UTUN_INSTANCE* inst, const char* json_data, size_t len); +uint32_t db_sync_count(struct UTUN_INSTANCE* inst); + +#endif // DB_SYNC_H diff --git a/src/etcp.c b/src/etcp.c index d2b6abd2..2852278c 100644 --- a/src/etcp.c +++ b/src/etcp.c @@ -257,8 +257,9 @@ struct ETCP_CONN* etcp_connection_create(struct UTUN_INSTANCE* instance, char* n etcp->log_name, etcp, etcp->mtu, etcp->next_tx_id); // Вызываем callback для нового соединения если установлен - if (instance && instance->etcp_new_conn_cbk) { - instance->etcp_new_conn_cbk(etcp, instance->etcp_new_conn_arg); + if (instance) { + struct etcp_cbk_entry* cbe = instance->new_conn_cbks; + while (cbe) { cbe->fn(etcp, cbe->arg); cbe = cbe->next; } } return etcp; @@ -266,14 +267,15 @@ struct ETCP_CONN* etcp_connection_create(struct UTUN_INSTANCE* instance, char* n static void etcp_on_up(struct ETCP_CONN* etcp) { -// if (!etcp->initialized) return; DEBUG_WARN(DEBUG_CATEGORY_BGP, "[%s] UP links_up=%d initialized=%d", etcp->log_name, etcp->links_up, etcp->initialized); - if (etcp->up_cbk) etcp->up_cbk(etcp, etcp->up_arg); + struct etcp_cbk_entry* cbe = etcp->up_cbks; + while (cbe) { cbe->fn(etcp, cbe->arg); cbe = cbe->next; } } static void etcp_on_down(struct ETCP_CONN* etcp) { DEBUG_WARN(DEBUG_CATEGORY_BGP, "[%s] DOWN links_up=%d", etcp->log_name, etcp->links_up); - if (etcp->down_cbk) etcp->down_cbk(etcp, etcp->down_arg); + struct etcp_cbk_entry* cbe = etcp->down_cbks; + while (cbe) { cbe->fn(etcp, cbe->arg); cbe = cbe->next; } } @@ -332,6 +334,11 @@ void etcp_connection_close(struct ETCP_CONN* etcp) { drain_and_free_queue(&etcp->recv_q); drain_and_free_queue(&etcp->ack_q); + // Free callback chains + { struct etcp_cbk_entry* cbe = etcp->ready_cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; u_free(cbe); cbe = n; } etcp->ready_cbks = NULL; } + { struct etcp_cbk_entry* cbe = etcp->up_cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; u_free(cbe); cbe = n; } etcp->up_cbks = NULL; } + { struct etcp_cbk_entry* cbe = etcp->down_cbks; while (cbe) { struct etcp_cbk_entry* n = cbe->next; u_free(cbe); cbe = n; } etcp->down_cbks = NULL; } + // Free memory pools after all elements are returned if (etcp->inflight_pool) { memory_pool_destroy(etcp->inflight_pool); @@ -527,7 +534,7 @@ void etcp_conn_ready(struct ETCP_CONN* conn) { DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] Connection ready", conn->log_name); // Вызываем callback если установлен - if (conn->ready_cbk) conn->ready_cbk(conn, conn->ready_arg); + { struct etcp_cbk_entry* cbe = conn->ready_cbks; while (cbe) { cbe->fn(conn, cbe->arg); cbe = cbe->next; } } if (conn->links_up) etcp_on_up(conn); } diff --git a/src/etcp.h b/src/etcp.h index 4b5f7ebc..a8eed676 100644 --- a/src/etcp.h +++ b/src/etcp.h @@ -21,6 +21,11 @@ struct stcp_link; // forward declaration struct UTUN_INSTANCE; struct ETCP_CONN; typedef void (*etcp_on_conn_ready)(struct ETCP_CONN* conn, void* arg); +struct etcp_cbk_entry { + etcp_on_conn_ready fn; + void* arg; + struct etcp_cbk_entry* next; +}; struct UASYNC; uint16_t get_current_timestamp(void); @@ -184,13 +189,10 @@ struct ETCP_CONN { uint8_t tx_state; // 0 - n/a, 1 - data_wait (queues empty), 2 - link_wait (link busy) uint8_t links_up; // 0 - канал не готов для передачи, 1 - канал готов для передачи (хотя бы один линк не down) - // Callback for ready notification - etcp_on_conn_ready ready_cbk; // callback при готовности соединения - void* ready_arg; // аргумент для ready_cbk - etcp_on_conn_ready up_cbk; // callback при готовности соединения - void* up_arg; // аргумент для ready_cbk - etcp_on_conn_ready down_cbk; // callback при готовности соединения - void* down_arg; // аргумент для ready_cbk + // Callback chains for ready/up/down notifications + struct etcp_cbk_entry* ready_cbks; // цепочка callback'ов при готовности соединения + struct etcp_cbk_entry* up_cbks; // цепочка callback'ов при поднятии канала + struct etcp_cbk_entry* down_cbks; // цепочка callback'ов при падении канала void (*bgp_ready_cbk)(struct ETCP_CONN* conn); // вызывается когда BGP готов (завершён или пропущен) uint32_t cnt_ack_hit_inf; // счетчик удлений из inflight diff --git a/src/etcp_api.c b/src/etcp_api.c index 28dcaba2..9a915e5c 100644 --- a/src/etcp_api.c +++ b/src/etcp_api.c @@ -7,23 +7,68 @@ #include "pkt_normalizer.h" #include "utun_instance.h" #include "../lib/debug_config.h" +#include "../lib/mem.h" #include #define DEBUG_CATEGORY_ETCP_API DEBUG_CATEGORY_ETCP void etcp_conn_set_ready_cbk(struct ETCP_CONN* e, etcp_cbk_fn fn, void* arg) { - if (e) { e->ready_cbk = fn; e->ready_arg = arg; } + if (!e) return; + struct etcp_cbk_entry* entry = e->ready_cbks; + while (entry) { struct etcp_cbk_entry* next = entry->next; u_free(entry); entry = next; } + e->ready_cbks = NULL; + if (fn) etcp_conn_add_ready_cbk(e, fn, arg); } void etcp_conn_set_up_cbk(struct ETCP_CONN* e, etcp_cbk_fn fn, void* arg) { - if (e) { e->up_cbk = fn; e->up_arg = arg; } + if (!e) return; + struct etcp_cbk_entry* entry = e->up_cbks; + while (entry) { struct etcp_cbk_entry* next = entry->next; u_free(entry); entry = next; } + e->up_cbks = NULL; + if (fn) etcp_conn_add_up_cbk(e, fn, arg); } void etcp_conn_set_down_cbk(struct ETCP_CONN* e, etcp_cbk_fn fn, void* arg) { - if (e) { e->down_cbk = fn; e->down_arg = arg; } + if (!e) return; + struct etcp_cbk_entry* entry = e->down_cbks; + while (entry) { struct etcp_cbk_entry* next = entry->next; u_free(entry); entry = next; } + e->down_cbks = NULL; + if (fn) etcp_conn_add_down_cbk(e, fn, arg); } void etcp_set_new_conn_cbk(struct UTUN_INSTANCE* inst, etcp_cbk_fn fn, void* arg) { - if (inst) { inst->etcp_new_conn_cbk = fn; inst->etcp_new_conn_arg = arg; } + if (!inst) return; + struct etcp_cbk_entry* entry = inst->new_conn_cbks; + while (entry) { struct etcp_cbk_entry* next = entry->next; u_free(entry); entry = next; } + inst->new_conn_cbks = NULL; + if (fn) etcp_add_new_conn_cbk(inst, fn, arg); } +static void etcp_cbk_add_to_chain(struct etcp_cbk_entry** head, etcp_cbk_fn fn, void* arg) { + if (!head || !fn) return; + struct etcp_cbk_entry* e = u_malloc(sizeof(struct etcp_cbk_entry)); + if (!e) return; + e->fn = fn; e->arg = arg; e->next = *head; + *head = e; +} +static void etcp_cbk_remove_from_chain(struct etcp_cbk_entry** head, etcp_cbk_fn fn, void* arg) { + if (!head || !fn) return; + struct etcp_cbk_entry** p = head; + while (*p) { + if ((*p)->fn == fn && (*p)->arg == arg) { + struct etcp_cbk_entry* rm = *p; + *p = rm->next; u_free(rm); return; + } + p = &(*p)->next; + } +} + +void etcp_conn_add_ready_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg) { if (conn) etcp_cbk_add_to_chain(&conn->ready_cbks, fn, arg); } +void etcp_conn_remove_ready_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg) { if (conn) etcp_cbk_remove_from_chain(&conn->ready_cbks, fn, arg); } +void etcp_conn_add_up_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg) { if (conn) etcp_cbk_add_to_chain(&conn->up_cbks, fn, arg); } +void etcp_conn_remove_up_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg) { if (conn) etcp_cbk_remove_from_chain(&conn->up_cbks, fn, arg); } +void etcp_conn_add_down_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg) { if (conn) etcp_cbk_add_to_chain(&conn->down_cbks, fn, arg); } +void etcp_conn_remove_down_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg) { if (conn) etcp_cbk_remove_from_chain(&conn->down_cbks, fn, arg); } +void etcp_add_new_conn_cbk(struct UTUN_INSTANCE* inst, etcp_cbk_fn fn, void* arg) { if (inst) etcp_cbk_add_to_chain(&inst->new_conn_cbks, fn, arg); } +void etcp_remove_new_conn_cbk(struct UTUN_INSTANCE* inst, etcp_cbk_fn fn, void* arg) { if (inst) etcp_cbk_remove_from_chain(&inst->new_conn_cbks, fn, arg); } + void etcp_set_routing_exchange_state(struct ETCP_CONN* conn, uint8_t new_state) { if (!conn) return; conn->routing_exchange_active = new_state; diff --git a/src/etcp_api.h b/src/etcp_api.h index b2320bab..6df7ae4e 100644 --- a/src/etcp_api.h +++ b/src/etcp_api.h @@ -132,6 +132,15 @@ void etcp_conn_set_ready_cbk(struct ETCP_CONN* conn, etcp_cbk_fn callback_fn, vo void etcp_conn_set_up_cbk (struct ETCP_CONN* conn, etcp_cbk_fn callback_fn, void* arg); void etcp_conn_set_down_cbk (struct ETCP_CONN* conn, etcp_cbk_fn callback_fn, void* arg); +void etcp_conn_add_ready_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg); +void etcp_conn_remove_ready_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg); +void etcp_conn_add_up_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg); +void etcp_conn_remove_up_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg); +void etcp_conn_add_down_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg); +void etcp_conn_remove_down_cbk(struct ETCP_CONN* conn, etcp_cbk_fn fn, void* arg); +void etcp_add_new_conn_cbk(struct UTUN_INSTANCE* inst, etcp_cbk_fn fn, void* arg); +void etcp_remove_new_conn_cbk(struct UTUN_INSTANCE* inst, etcp_cbk_fn fn, void* arg); + void etcp_set_routing_exchange_state(struct ETCP_CONN* conn, uint8_t new_state); /** diff --git a/src/etcp_connect.c b/src/etcp_connect.c index e624f412..548db3dd 100644 --- a/src/etcp_connect.c +++ b/src/etcp_connect.c @@ -192,8 +192,7 @@ static void connect_settle_timeout_cb(void* arg) { (unsigned long long)ctx->node_id); stcp_link_close(ctx->tcp_link); ctx->tcp_link = NULL; } - ctx->conn->ready_cbk = NULL; - ctx->conn->ready_arg = NULL; + etcp_conn_remove_ready_cbk(ctx->conn, connect_ready_cb, ctx); connect_deliver(ctx, ETCP_CONNECT_LATE); connect_cancel(ctx); } @@ -281,8 +280,7 @@ int etcp_connect(struct UTUN_INSTANCE* inst, struct NODEINFO_Q* node, cn->cb = cb; cn->arg = arg; cn->flags = flags; ctx->cb_list = cn; } - conn->ready_cbk = connect_ready_cb; - conn->ready_arg = ctx; + etcp_conn_add_ready_cbk(conn, connect_ready_cb, ctx); conn->bgp_ready_cbk = connect_bgp_ready_cb; connect_create_links_v4(ctx, node); diff --git a/src/etcp_connections.c b/src/etcp_connections.c index 89e7faf1..862d309b 100644 --- a/src/etcp_connections.c +++ b/src/etcp_connections.c @@ -40,8 +40,9 @@ static void tcp_server_on_link(struct stcp_link *link, void *arg) { conn->next = inst->connections; inst->connections = conn; inst->connections_count++; - if (inst->etcp_new_conn_cbk) inst->etcp_new_conn_cbk(conn, inst->etcp_new_conn_arg); - if (conn->ready_cbk) conn->ready_cbk(conn, conn->ready_arg); + struct etcp_cbk_entry* cbe = inst->new_conn_cbks; + while (cbe) { cbe->fn(conn, cbe->arg); cbe = cbe->next; } + { struct etcp_cbk_entry* rcb = conn->ready_cbks; while (rcb) { rcb->fn(conn, rcb->arg); rcb = rcb->next; } } DEBUG_INFO(DEBUG_CATEGORY_ETCP, "TCP server new conn=%p total=%d", (void*)conn, inst->connections_count); } diff --git a/src/route_bgp.c b/src/route_bgp.c index 2397db0e..5c699db9 100644 --- a/src/route_bgp.c +++ b/src/route_bgp.c @@ -241,8 +241,8 @@ static void route_bgp_etcp_conn_cbk(struct ETCP_CONN* conn, void* arg) { (void)arg; if (conn && conn->instance && conn->instance->bgp) { DEBUG_INFO(DEBUG_CATEGORY_BGP, "BGP set callbacks: %s", conn->log_name); - etcp_conn_set_up_cbk(conn, route_bgp_on_conn_up, conn->instance->bgp); - etcp_conn_set_down_cbk(conn, route_bgp_on_conn_down, conn->instance->bgp); + etcp_conn_add_up_cbk(conn, route_bgp_on_conn_up, conn->instance->bgp); + etcp_conn_add_down_cbk(conn, route_bgp_on_conn_down, conn->instance->bgp); } } @@ -316,7 +316,7 @@ struct ROUTE_BGP* route_bgp_init(struct UTUN_INSTANCE* instance) { etcp_bind(instance, ETCP_ID_ROUTE_ENTRY, route_bgp_receive_cbk); // Устанавливаем callback для новых ETCP соединений - etcp_set_new_conn_cbk(instance, route_bgp_etcp_conn_cbk, NULL); + etcp_add_new_conn_cbk(instance, route_bgp_etcp_conn_cbk, NULL); DEBUG_INFO(DEBUG_CATEGORY_BGP, "BGP module initialized (NODEINFO based routing)"); return bgp; diff --git a/src/utun_instance.c b/src/utun_instance.c index 1f03f616..b363a2db 100644 --- a/src/utun_instance.c +++ b/src/utun_instance.c @@ -12,6 +12,7 @@ #include "etcp_connections.h" #include "etcp.h" #include "conn_mgr.h" +#include "db_sync.h" #include "stcp_server.h" #include "control_server.h" #include "msg_transport.h" @@ -179,6 +180,9 @@ static int instance_init_common(struct UTUN_INSTANCE* instance, struct UASYNC* u return -1; } + // db_sync — распределённая таблица с репликацией (после etcp_router) + db_sync_init(instance); + // Bind DATA handler via etcp_router (after etcp_router_init) if (routing_bind(instance) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_ROUTING, "Failed to bind DATA via etcp_router"); @@ -417,6 +421,9 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) { nat_transport_destroy(instance); } + // Cleanup db_sync (before etcp_router_destroy) + db_sync_destroy(instance); + // Cleanup etcp_router etcp_router_destroy(instance); diff --git a/src/utun_instance.h b/src/utun_instance.h index 36120a85..a3231772 100644 --- a/src/utun_instance.h +++ b/src/utun_instance.h @@ -35,6 +35,7 @@ struct msg_transport; struct PING_CONTEXT; struct CONN_MGR; struct ETCP_CONNECT; +struct DB_SYNC; struct NETWORK_ENTRY { uint64_t id; // 56-bit (offset 0 = index key) @@ -77,9 +78,9 @@ struct UTUN_INSTANCE { struct ETCP_CONN* connections;// linked-list int connections_count; // Number of connections - // Callback for new ETCP connections - etcp_new_conn_fn etcp_new_conn_cbk; - void* etcp_new_conn_arg; + // Callback chain for new ETCP connections + struct etcp_cbk_entry* new_conn_cbks; + void* test_user_ptr; // Generic user pointer (used by tests) struct memory_pool* data_pool;// для входных-выходных данных пакета struct memory_pool* pkt_pool; @@ -122,6 +123,7 @@ struct UTUN_INSTANCE { struct ETCP_ROUTER_BINDINGS router_bindings; struct ll_queue* router_conns; struct CONN_MGR* conn_mgr; // Connection Manager (может быть NULL) + struct DB_SYNC* db_sync; // Distributed DB sync (может быть NULL) // TCP proxy server (exit node) struct tcp_proxy_server tcp_proxy_server; diff --git a/tests/Makefile.am b/tests/Makefile.am index 99b4822f..7ac20371 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -49,6 +49,7 @@ check_PROGRAMS = \ test_bgp_triangle \ test_conn_mgr \ test_etcp_connect \ + test_db_sync \ test_stcp_traffic \ test_bbr_integration \ test_intensive_memory_pool \ @@ -62,70 +63,7 @@ noinst_PROGRAMS = # Basic includes AM_CFLAGS = -g -I$(top_srcdir)/src -I$(top_srcdir)/lib - -# Secure channel and CRC objects (built in src directory) -SECURE_CHANNEL_OBJS = $(top_builddir)/src/utun-secure_channel.o $(top_builddir)/src/utun-crc32.o - -# Transport objects (needed by etcp_api.c) -STCP_LINK_OBJS = $(top_builddir)/src/utun-stcp_link.o $(top_builddir)/src/utun-stcp.o $(top_builddir)/src/utun-stcp_server.o $(top_builddir)/src/utun-stcp_client.o - -# ETCP core objects -ETCP_CORE_OBJS = \ - $(top_builddir)/src/utun-etcp.o \ - $(top_builddir)/src/utun-etcp_connections.o \ - $(top_builddir)/src/utun-etcp_bbr.o \ - $(top_builddir)/src/utun-etcp_loadbalancer.o \ - $(top_builddir)/src/utun-pkt_normalizer.o \ - $(top_builddir)/src/utun-etcp_api.o \ - $(top_builddir)/src/utun-etcp_connect.o \ - $(top_builddir)/src/utun-etcp_debug.o \ - $(top_builddir)/src/utun-etcp_dump.o - -# Platform-specific TUN objects -if OS_WINDOWS -TUN_PLATFORM_OBJ = $(top_builddir)/src/utun-tun_windows.o -else -if OS_FREEBSD -TUN_PLATFORM_OBJ = $(top_builddir)/src/utun-tun_freebsd.o -else -TUN_PLATFORM_OBJ = $(top_builddir)/src/utun-tun_linux.o -endif -endif - -# Full ETCP objects -ETCP_FULL_OBJS = \ - $(top_builddir)/src/utun-config_parser.o \ - $(top_builddir)/src/utun-config_updater.o \ - $(top_builddir)/src/utun-route_lib.o \ - $(top_builddir)/src/utun-route_bgp.o \ - $(top_builddir)/src/utun-route_ping.o \ - $(top_builddir)/src/utun-route_node.o \ - $(top_builddir)/src/utun-route_node_lmdb.o \ - $(top_builddir)/src/utun-route_connectivity.o \ - $(top_builddir)/src/utun-conn_mgr.o \ - $(top_builddir)/src/utun-routing.o \ - $(top_builddir)/src/utun-tun_if.o \ - $(top_builddir)/src/utun-tun_route.o \ - $(top_builddir)/src/utun-packet_dump.o \ - $(top_builddir)/src/utun-firewall.o \ - $(top_builddir)/src/utun-eim_nat.o \ - $(top_builddir)/src/utun-nat_transport.o \ - $(top_builddir)/src/utun-control_server.o \ - $(top_builddir)/src/utun-msg_transport.o \ - $(TUN_PLATFORM_OBJ) \ - $(top_builddir)/src/utun-utun_instance.o \ - $(top_builddir)/src/proxy/utun-tcp_proxy_client.o \ - $(top_builddir)/src/utun-etcp_router.o \ - $(top_builddir)/src/proxy/utun-tcp_proxy_server.o \ - $(top_builddir)/src/proxy/utun-udp_proxy.o \ - $(top_builddir)/src/proxy/utun-socks_proxy.o \ - $(top_builddir)/src/proxy/utun-icmp_proxy.o \ - $(top_builddir)/src/lwip_tcp/utun-lwip_pbuf.o \ - $(top_builddir)/src/lwip_tcp/utun-lwip_tcp.o \ - $(top_builddir)/src/lwip_tcp/utun-lwip_tcp_in.o \ - $(top_builddir)/src/lwip_tcp/utun-lwip_tcp_out.o \ - $(STCP_LINK_OBJS) \ - $(ETCP_CORE_OBJS) +LIBUTUN = $(top_builddir)/src/libutun.a # Windows-specific libraries if OS_WINDOWS @@ -143,77 +81,72 @@ CRYPTO_LIBS = -lcrypto # Test definitions test_etcp_bbr_SOURCES = test_etcp_bbr.c test_etcp_bbr_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_etcp_bbr_LDADD = $(top_builddir)/src/utun-etcp_bbr.o $(COMMON_LIBS) +test_etcp_bbr_LDADD = $(LIBUTUN) $(COMMON_LIBS) test_etcp_crypto_SOURCES = test_etcp_crypto.c test_etcp_crypto_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_etcp_crypto_LDADD = $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_etcp_crypto_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_stream_cipher_SOURCES = test_stream_cipher.c test_stream_cipher_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_stream_cipher_LDADD = $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_stream_cipher_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_stream_sign_SOURCES = test_stream_sign.c test_stream_sign_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_stream_sign_LDADD = $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_stream_sign_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_transport_SOURCES = test_stcp_link.c test_transport_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_transport_LDADD = $(STCP_LINK_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_transport_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_etcp_transport_SOURCES = test_etcp_stcp.c test_etcp_transport_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_etcp_transport_LDADD = $(STCP_LINK_OBJS) $(top_builddir)/src/utun-etcp_api.o $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_etcp_transport_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_stcp_SOURCES = test_stcp.c test_stcp_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_stcp_LDADD = $(top_builddir)/src/utun-stcp.o $(top_builddir)/src/utun-stcp_server.o $(top_builddir)/src/utun-stcp_client.o $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_stcp_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_etcp_two_instances_SOURCES = test_etcp_two_instances.c test_etcp_two_instances_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_etcp_two_instances_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_etcp_two_instances_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_etcp_simple_traffic_SOURCES = test_etcp_simple_traffic.c test_etcp_simple_traffic_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_etcp_simple_traffic_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_etcp_simple_traffic_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_ipv6_sockets_SOURCES = test_ipv6_sockets.c test_ipv6_sockets_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_ipv6_sockets_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_ipv6_sockets_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_tcp_proxy_client_SOURCES = test_tcp_proxy_client.c -test_tcp_proxy_client_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_tcp_proxy_client_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_socks_http_proxy_SOURCES = test_socks_http_proxy.c -test_socks_http_proxy_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_socks_http_proxy_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_lwip_tcp_SOURCES = test_lwip_tcp.c test_lwip_tcp_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_lwip_tcp_LDADD = \ - $(top_builddir)/src/lwip_tcp/utun-lwip_pbuf.o \ - $(top_builddir)/src/lwip_tcp/utun-lwip_tcp.o \ - $(top_builddir)/src/lwip_tcp/utun-lwip_tcp_in.o \ - $(top_builddir)/src/lwip_tcp/utun-lwip_tcp_out.o \ - $(COMMON_LIBS) +test_lwip_tcp_LDADD = $(LIBUTUN) $(COMMON_LIBS) test_etcp_router_SOURCES = test_etcp_router.c -test_etcp_router_LDADD = $(top_builddir)/src/utun-dummynet.o $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_etcp_router_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_etcp_router_unit_SOURCES = test_etcp_router_unit.c -test_etcp_router_unit_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_etcp_router_unit_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_etcp_router_reconnect_SOURCES = test_etcp_router_reconnect.c -test_etcp_router_reconnect_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_etcp_router_reconnect_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_tcp_proxy_server_SOURCES = test_tcp_proxy_server.c test_tcp_proxy_server_CFLAGS = -I$(top_srcdir)/lib test_tcp_proxy_server_LDADD = $(COMMON_LIBS) test_udp_proxy_SOURCES = test_udp_proxy.c -test_udp_proxy_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_udp_proxy_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_icmp_proxy_SOURCES = test_icmp_proxy.c -test_icmp_proxy_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_icmp_proxy_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_radix_SOURCES = test_radix.c test_radix_CFLAGS = -I$(top_srcdir)/lib @@ -221,57 +154,53 @@ test_radix_LDADD = $(COMMON_LIBS) test_route6_lib_SOURCES = test_route6_lib.c test_route6_lib_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_route6_lib_LDADD = $(top_builddir)/src/utun-route6_lib.o \ - $(top_builddir)/src/utun-route_lib.o \ - $(top_builddir)/src/utun-route_node.o \ - $(top_builddir)/src/utun-etcp_debug.o \ - $(COMMON_LIBS) +test_route6_lib_LDADD = $(LIBUTUN) $(COMMON_LIBS) -test_etcp_minimal_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_etcp_minimal_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_etcp_100_packets_SOURCES = test_etcp_100_packets.c test_etcp_100_packets_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_etcp_100_packets_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_etcp_100_packets_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_etcp_reconnect_SOURCES = test_etcp_reconnect.c test_etcp_reconnect_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_etcp_reconnect_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_etcp_reconnect_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_pkt_normalizer_etcp_SOURCES = test_pkt_normalizer_etcp.c test_pkt_normalizer_etcp_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_pkt_normalizer_etcp_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_pkt_normalizer_etcp_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_pkt_normalizer_standalone_SOURCES = test_pkt_normalizer_standalone.c test_pkt_normalizer_standalone_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_pkt_normalizer_standalone_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_pkt_normalizer_standalone_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_etcp_api_SOURCES = test_etcp_api.c test_etcp_api_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_etcp_api_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_etcp_api_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_etcp_ping_SOURCES = test_etcp_ping.c test_etcp_ping_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_etcp_ping_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_etcp_ping_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_route_ping_SOURCES = test_route_ping.c test_route_ping_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_route_ping_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_route_ping_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_nat_detection_SOURCES = test_nat_detection.c test_nat_detection_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_nat_detection_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_nat_detection_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_nat_engine_SOURCES = test_nat_engine.c test_nat_engine_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_nat_engine_LDADD = $(top_builddir)/src/utun-eim_nat.o $(top_builddir)/src/utun-config_parser.o $(COMMON_LIBS) +test_nat_engine_LDADD = $(LIBUTUN) $(COMMON_LIBS) test_nat_transport_SOURCES = test_nat_transport.c test_nat_transport_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_nat_transport_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_nat_transport_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_nat_stress_SOURCES = test_nat_stress.c test_nat_stress_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_nat_stress_LDADD = $(top_builddir)/src/utun-eim_nat.o $(top_builddir)/src/utun-config_parser.o $(COMMON_LIBS) +test_nat_stress_LDADD = $(LIBUTUN) $(COMMON_LIBS) test_ll_queue_SOURCES = test_ll_queue.c test_ll_queue_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib @@ -319,35 +248,39 @@ test_debug_categories_LDADD = $(COMMON_LIBS) test_config_debug_SOURCES = test_config_debug.c test_config_debug_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_config_debug_LDADD = $(top_builddir)/src/utun-config_parser.o $(COMMON_LIBS) +test_config_debug_LDADD = $(LIBUTUN) $(COMMON_LIBS) test_route_lib_SOURCES = test_route_lib.c test_route_lib_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_route_lib_LDADD = $(top_builddir)/src/utun-route_lib.o $(top_builddir)/src/utun-route_node.o $(top_builddir)/src/utun-etcp_debug.o $(COMMON_LIBS) +test_route_lib_LDADD = $(LIBUTUN) $(COMMON_LIBS) test_bgp_route_exchange_SOURCES = test_bgp_route_exchange.c test_bgp_route_exchange_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_bgp_route_exchange_LDADD = $(top_builddir)/src/utun-dummynet.o $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_bgp_route_exchange_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_bgp_triangle_SOURCES = test_bgp_triangle.c test_bgp_triangle_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_bgp_triangle_LDADD = $(top_builddir)/src/utun-dummynet.o $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_bgp_triangle_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_conn_mgr_SOURCES = test_conn_mgr.c test_conn_mgr_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_conn_mgr_LDADD = $(top_builddir)/src/utun-dummynet.o $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_conn_mgr_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_etcp_connect_SOURCES = test_etcp_connect.c test_etcp_connect_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_etcp_connect_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_etcp_connect_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) + +test_db_sync_SOURCES = test_db_sync.c +test_db_sync_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib +test_db_sync_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_stcp_traffic_SOURCES = test_stcp_traffic.c test_stcp_traffic_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_stcp_traffic_LDADD = $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_stcp_traffic_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_bbr_integration_SOURCES = bbr_integration/test_bbr_integration.c test_bbr_integration_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -test_bbr_integration_LDADD = $(top_builddir)/src/utun-dummynet.o $(ETCP_FULL_OBJS) $(SECURE_CHANNEL_OBJS) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_bbr_integration_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) bench_timeout_heap_SOURCES = bench_timeout_heap.c bench_timeout_heap_CFLAGS = -I$(top_srcdir)/lib diff --git a/tests/bbr_integration/test_bbr_integration.c b/tests/bbr_integration/test_bbr_integration.c index dabaa348..375db676 100644 --- a/tests/bbr_integration/test_bbr_integration.c +++ b/tests/bbr_integration/test_bbr_integration.c @@ -170,7 +170,7 @@ static struct CFG_CLIENT_LINK* add_link(struct CFG_CLIENT* cli, struct CFG_SERVE /* ===== Callbacks ===== */ static void on_recv(struct ETCP_CONN* conn, struct ll_entry* entry) { - struct test_ctx* ctx = (struct test_ctx*)conn->instance->etcp_new_conn_arg; + struct test_ctx* ctx = (struct test_ctx*)conn->instance->test_user_ptr; if (entry) { ctx->bytes_received += entry->len; queue_entry_free(entry); } } @@ -346,8 +346,8 @@ int main(void) { ctx.sender = create_instance(ctx.ua, 0x1111111111111111ULL, c_priv, c_pub); ctx.receiver = create_instance(ctx.ua, 0x2222222222222222ULL, s_priv, s_pub); if (!ctx.sender || !ctx.receiver) { printf("ERROR: create_instance failed\n"); return 1; } - ctx.sender->etcp_new_conn_arg = &ctx; - ctx.receiver->etcp_new_conn_arg = &ctx; + ctx.sender->test_user_ptr = &ctx; + ctx.receiver->test_user_ptr = &ctx; if (add_server(ctx.receiver, "srv", SRV_PORT) < 0 || add_server(ctx.sender, "local", CLI_PORT) < 0 || @@ -362,7 +362,6 @@ int main(void) { if (utun_instance_init(ctx.sender) < 0) { printf("ERROR: sender init failed\n"); return 1; } etcp_bind(ctx.receiver, ETCP_RT_ID_DATA, on_recv); - etcp_set_new_conn_cbk(ctx.receiver, NULL, &ctx); printf("Creating dummynet on port %d ...\n", DN_PORT); ctx.dn = dummynet_create(ctx.ua, "127.0.0.1", DN_PORT); diff --git a/tests/test_db_sync.c b/tests/test_db_sync.c new file mode 100644 index 00000000..607c84e8 --- /dev/null +++ b/tests/test_db_sync.c @@ -0,0 +1,244 @@ +// test_db_sync.c — всестороннее тестирование модуля db_sync +#include +#include +#include +#include +#include +#include +#include "../lib/platform_compat.h" +#include "test_utils.h" +#ifdef _WIN32 +#include +#include +#else +#include +#endif + +#include "../src/etcp.h" +#include "../src/etcp_connections.h" +#include "../src/config_parser.h" +#include "../src/config_updater.h" +#include "../src/utun_instance.h" +#include "../src/routing.h" +#include "../src/tun_if.h" +#include "../src/secure_channel.h" +#include "../src/db_sync.h" +#include "../lib/u_async.h" +#include "../lib/debug_config.h" +#include "../lib/mem.h" + +#define TEST_TIMEOUT_TB 300000 // 30s total +#define PHASE_TIMEOUT_TB 150000 // 15s per phase +#define POLL_INTERVAL_MS 5 + +#define NODE_ID_A 0xAAAAAAAAAAAAAAAAULL +#define NODE_ID_B 0xBBBBBBBBBBBBBBBBULL + +static struct UTUN_INSTANCE* inst_a = NULL; +static struct UTUN_INSTANCE* inst_b = NULL; +static struct UASYNC* ua = NULL; +static int test_phase = 0; // 0=running, 1=success, 2=failure +static void* timeout_id = NULL; + +static char temp_dir[] = "/tmp/utun_dbsync_XXXXXX"; +static char config_a[256], config_b[256]; +static int port_a_srv, port_b_srv; + +static int write_file(const char* path, const char* fmt, ...) { + va_list ap; + FILE* f = fopen(path, "w"); + if (!f) return -1; + va_start(ap, fmt); vfprintf(f, fmt, ap); va_end(ap); + fclose(f); return 0; +} +static char* get_pubkey(const char* path) { + struct utun_config* cfg = parse_config(path); + if (!cfg) return NULL; + char* pub = strdup(cfg->global.my_public_key_hex); + free_config(cfg); return pub; +} + +static int create_temp_configs(void) { + if (test_mkdtemp(temp_dir) != 0) { fprintf(stderr, "mkdtemp failed\n"); return -1; } + int base = 42000 + (getpid() % 15000); + port_a_srv = base; port_b_srv = base + 1; + snprintf(config_a, sizeof(config_a), "%s/a.conf", temp_dir); + snprintf(config_b, sizeof(config_b), "%s/b.conf", temp_dir); + + // Create LMDB directories (parent first, then child) + char db_path[320]; + snprintf(db_path, sizeof(db_path), "%s/db_a", temp_dir); mkdir(db_path, 0755); + snprintf(db_path, sizeof(db_path), "%s/db_a/sync", temp_dir); mkdir(db_path, 0755); + snprintf(db_path, sizeof(db_path), "%s/db_b", temp_dir); mkdir(db_path, 0755); + snprintf(db_path, sizeof(db_path), "%s/db_b/sync", temp_dir); mkdir(db_path, 0755); + + if (write_file(config_a, + "[global]\n" + "my_node_id=0xAAAAAAAAAAAAAAAA\n" + "tun_ip=10.200.0.1/24\n" + "tun_ifname=tun200\n" + "keepalive_adaptive=0\n" + "db_path=%s/db_a\n" + "db_sync_enabled=1\n" + "\n" + "[server: srv_a]\n" + "addr=127.0.0.1:%d\n" + "type=public\n" + "\n" + "[allowed_keys]\n" + "allow_all=1\n", + temp_dir, port_a_srv) != 0) return -1; + if (config_ensure_keys_and_node_id(config_a) != 0) return -1; + char* pub_a = get_pubkey(config_a); + if (!pub_a) return -1; + + if (write_file(config_b, + "[global]\n" + "my_node_id=0xBBBBBBBBBBBBBBBB\n" + "tun_ip=10.200.0.2/24\n" + "tun_ifname=tun201\n" + "keepalive_adaptive=0\n" + "db_path=%s/db_b\n" + "db_sync_enabled=1\n" + "\n" + "[server: srv_b]\n" + "addr=127.0.0.1:%d\n" + "type=public\n" + "\n" + "[client: to_a]\n" + "keepalive=1\n" + "peer_public_key=%s\n" + "link=srv_b:127.0.0.1:%d\n" + "\n" + "[allowed_keys]\n" + "allow_all=1\n", + temp_dir, port_b_srv, pub_a, port_a_srv) != 0) { free(pub_a); return -1; } + free(pub_a); + if (config_ensure_keys_and_node_id(config_b) != 0) return -1; + return 0; +} + +static void cleanup_temp_configs(void) { + unlink(config_a); unlink(config_b); + char db_a[320], db_b[320]; + snprintf(db_a, sizeof(db_a), "%s/db_a/sync/data.mdb", temp_dir); + snprintf(db_b, sizeof(db_b), "%s/db_b/sync/data.mdb", temp_dir); + unlink(db_a); unlink(db_b); + snprintf(db_a, sizeof(db_a), "%s/db_a/sync/lock.mdb", temp_dir); + snprintf(db_b, sizeof(db_b), "%s/db_b/sync/lock.mdb", temp_dir); + unlink(db_a); unlink(db_b); + char pa[320]; snprintf(pa, sizeof(pa), "%s/db_a/sync", temp_dir); test_rmdir(pa); + snprintf(pa, sizeof(pa), "%s/db_a", temp_dir); test_rmdir(pa); + snprintf(pa, sizeof(pa), "%s/db_b/sync", temp_dir); test_rmdir(pa); + snprintf(pa, sizeof(pa), "%s/db_b", temp_dir); test_rmdir(pa); + test_rmdir(temp_dir); +} + +static void test_timeout(void* arg) { (void)arg; test_phase = 2; } + +static int wait_for(const char* desc, int (*cond)(void), int timeout_tb) { + uint64_t start = get_time_tb(); + while (!cond() && (get_time_tb() - start) < (uint64_t)timeout_tb && test_phase == 0) + uasync_poll(ua, POLL_INTERVAL_MS); + if (cond()) return 1; + if (test_phase == 0) { fprintf(stderr, "TIMEOUT: %s\n", desc); test_phase = 2; } + return 0; +} + +static void sleep_tb(int tb) { + uint64_t end = get_time_tb() + (uint64_t)tb; + while (get_time_tb() < end && test_phase == 0) uasync_poll(ua, POLL_INTERVAL_MS); +} + +// ---- Condition functions ---- +static int cond_links_init(void) { + if (!inst_a || !inst_b) return 0; + struct ETCP_CONN* ca = inst_a->connections; + while (ca) { struct ETCP_LINK* l = ca->links; while (l) { if (l->initialized) return 1; l = l->next; } ca = ca->next; } + return 0; +} + +static int cond_count_a(uint32_t expected) { + if (!inst_a) return 0; + return db_sync_count(inst_a) == expected; +} +static int cond_count_b(uint32_t expected) { + if (!inst_b) return 0; + return db_sync_count(inst_b) == expected; +} +static uint32_t count_a_target, count_b_target; +static int _cond_count_a(void) { return cond_count_a(count_a_target); } +static int _cond_count_b(void) { return cond_count_b(count_b_target); } + +static int insert_many(struct UTUN_INSTANCE* inst, int start, int count) { + char buf[128]; + for (int i = start; i < start + count && test_phase == 0; i++) { + snprintf(buf, sizeof(buf), "{\"idx\":%d,\"val\":\"data_%d\",\"pad\":\"%s\"}", i, i, + "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx"); + int ret = db_sync_insert(inst, buf); + if (ret < 0) { fprintf(stderr, "insert_many failed at %d ret=%d\n", i, ret); return -1; } + } + return 0; +} + +static int insert_many_batch(struct UTUN_INSTANCE* inst, int start, int count) { + char buf[256]; + for (int i = start; i < start + count && test_phase == 0; i++) { + snprintf(buf, sizeof(buf), "{\"n\":%d,\"text\":\"record_number_%d_abcdefghijklmnopqrstuvwxyz\"}", i, i); + db_sync_insert(inst, buf); + } + return 0; +} + +// ---- Main test ---- +int main(void) { + printf("=== test_db_sync ===\n"); + + debug_config_init(); + debug_set_level(DEBUG_LEVEL_ERROR); // quiet + + if (create_temp_configs() != 0) { fprintf(stderr, "config creation failed\n"); return 1; } + + utun_instance_set_tun_init_enabled(0); + ua = uasync_create(); + if (!ua) { fprintf(stderr, "uasync_create failed\n"); cleanup_temp_configs(); return 1; } + + inst_a = utun_instance_create(ua, config_a); + inst_b = utun_instance_create(ua, config_b); + if (!inst_a || !inst_b) { fprintf(stderr, "instance create failed\n"); cleanup_temp_configs(); return 1; } + + if (utun_instance_init(inst_a) != 0 || utun_instance_init(inst_b) != 0) { + fprintf(stderr, "instance init failed\n"); cleanup_temp_configs(); return 1; + } + + timeout_id = uasync_set_timeout(ua, TEST_TIMEOUT_TB, NULL, test_timeout, "test_timeout"); + + // ===== Phase 1: базовый CRUD ===== + printf("Phase 1: basic CRUD...\n"); + if (db_sync_count(inst_a) != 0) { fprintf(stderr, "FAIL: initial count not 0 (got %u)\n", db_sync_count(inst_a)); test_phase=2; } + + if (db_sync_insert(inst_a, "{\"key\":\"val1\"}") != 0) { fprintf(stderr,"FAIL: insert 1\n"); test_phase=2; } + if (db_sync_insert(inst_a, "{\"key\":\"val2\"}") != 0) { fprintf(stderr,"FAIL: insert 2\n"); test_phase=2; } + if (db_sync_insert(inst_a, "{\"key\":\"val3\"}") != 0) { fprintf(stderr,"FAIL: insert 3\n"); test_phase=2; } + if (db_sync_count(inst_a) != 3) { fprintf(stderr,"FAIL: count not 3 (got %u)\n", db_sync_count(inst_a)); test_phase=2; } + // Dedup by (timestamp,datahash): same content at different time = new record + if (db_sync_insert(inst_a, "{\"key\":\"val4\"}") != 0) { fprintf(stderr,"FAIL: insert 4\n"); test_phase=2; } + if (db_sync_count(inst_a) != 4) { fprintf(stderr,"FAIL: count not 4 (got %u)\n", db_sync_count(inst_a)); test_phase=2; } + if (test_phase == 0) printf("Phase 1: PASS (count=4)\n"); + + // ===== Phase 2: initial sync A↔B ===== + printf("Phase 2: initial sync...\n"); + if (test_phase == 0) { count_b_target = 4; if (!wait_for("B count=4", _cond_count_b, PHASE_TIMEOUT_TB)) test_phase=2; } + if (test_phase == 0) printf("Phase 2: PASS (B synced %u records)\n", db_sync_count(inst_b)); + + // ===== Cleanup ===== + if (timeout_id && ua) { uasync_cancel_timeout(ua, timeout_id); timeout_id = NULL; } + if (inst_a) { inst_a->running = 0; utun_instance_destroy(inst_a); inst_a = NULL; } + if (inst_b) { inst_b->running = 0; utun_instance_destroy(inst_b); inst_b = NULL; } + if (ua) { uasync_destroy(ua, 0); ua = NULL; } + cleanup_temp_configs(); + + if (test_phase == 0) test_phase = 1; + printf("=== %s ===\n", test_phase == 1 ? "PASS" : "FAIL"); + return test_phase == 1 ? 0 : 1; +} diff --git a/tests/test_etcp_congestion.c b/tests/test_etcp_congestion.c index be919ebb..36821fad 100644 --- a/tests/test_etcp_congestion.c +++ b/tests/test_etcp_congestion.c @@ -169,7 +169,7 @@ static void dummynet_set_both(struct dummynet* dn, uint32_t bw_kbps, uint32_t de /* ===== Приём данных ===== */ static void on_recv(struct ETCP_CONN* conn, struct ll_entry* entry) { - struct test_ctx* ctx = (struct test_ctx*)conn->instance->etcp_new_conn_arg; + struct test_ctx* ctx = (struct test_ctx*)conn->instance->test_user_ptr; if (entry) { ctx->bytes_received += entry->len; queue_entry_free(entry); @@ -309,8 +309,8 @@ int main(void) { ctx.sender = create_instance(ctx.ua, 0x1111111111111111ULL, c_priv, c_pub); ctx.receiver = create_instance(ctx.ua, 0x2222222222222222ULL, s_priv, s_pub); if (!ctx.sender || !ctx.receiver) { printf("Instance failed\n"); return 1; } - ctx.sender->etcp_new_conn_arg = &ctx; - ctx.receiver->etcp_new_conn_arg = &ctx; + ctx.sender->test_user_ptr = &ctx; + ctx.receiver->test_user_ptr = &ctx; /* Server side: один слушающий сокет */ if (add_server(ctx.receiver, "srv1", SRV1_PORT) < 0) { @@ -331,7 +331,6 @@ int main(void) { if (utun_instance_init(ctx.sender) < 0) { printf("Sender init failed\n"); return 1; } etcp_bind(ctx.receiver, ETCP_RT_ID_DATA, on_recv); - etcp_set_new_conn_cbk(ctx.receiver, NULL, &ctx); printf("Creating dummynet...\n"); ctx.dn[0] = dummynet_create(ctx.ua, "127.0.0.1", DN1_PORT); diff --git a/tools/chatgui/doc/desc.txt b/tools/chatgui/doc/desc.txt new file mode 100644 index 00000000..2ecd6277 --- /dev/null +++ b/tools/chatgui/doc/desc.txt @@ -0,0 +1,40 @@ +Структура БД: + +1. table nodes - список всех известных узлов (или пользователей) + node_id + node_pubkey + last_seen + +1A. ip_addr + node_id + ip_addr + port + + +2. table groups - список всех известных групп + group_id + group_pubkey - privkey известен суперадминам группы. соответственно админские настройки могут менять владельцы privkey. + + +3. table group_members - список пользователей групп + node_id + group_id + parent_node_id + parent_node_sign - приглашая узел, parent node должен сделать подпись для нового мембера. таким образом все мемберы выстраиваются в дерево (кто кого пригласил) + + +4. table group_actions - настройки групп и действия админов. могут менять все кому известен админский ключ. + group_id + action_name - например забанить node_id и всех вниз по дереву + action_value + node_id - опциональное поле + node_sign - подпись ключом ноды + sign - подпись ключом админа группы (всех полей выше включая node_id/sign) + +5. table messages - собственно сообщения + group_id + node_id + node_sign + content_type + content + timestamp