Browse Source
- Add libutun.a (noinst) containing all src/*.c except utun.c - utun binary now links libutun.a + libuasync.a - tests/Makefile.am simplified: all tests link libutun.a instead of manual .o file lists (removed ~70 lines of boilerplate) - Add db_sync module: distributed content-addressed DB with LMDB backend, peer-to-peer sync via etcp_router (svc 0x20) - etcp API: etcp_send_in_order(), etcp_clear_inflight_queues(), verify_packet_signature(), BBR getters - BBR: adaptive threshold (1/4 of pipe-min instead of hardcoded) - route_bgp/etcp_connect: STCP transport support in NODEINFO - conn_mgr: NODEINFO_ADDRS doc for nat_type/transport/proto fields - Cleanup: removed AUDIT_FINDINGS.txt, filelist.txt, nat.txtchatgui
23 changed files with 1508 additions and 786 deletions
@ -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 (дубликат секции) |
||||
@ -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 |
||||
@ -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] |
||||
@ -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) то узел отправляет сразу содержимое записей. |
||||
(надо додумать алгоритм, чтобы оптимизировать количество итераций - лучше передать больше данных за раз чем много итераций с ожиданием ответной стороны) |
||||
@ -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 <string.h> |
||||
|
||||
#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 && i<pos; i++) rc = mdb_cursor_get(cursor, &key, &val, MDB_NEXT); |
||||
if (rc==0) { |
||||
*datahash = be64toh(*(uint64_t*)((uint8_t*)key.mv_data+8)); |
||||
mdb_cursor_close(cursor); mdb_txn_abort(txn); return 0; |
||||
} |
||||
mdb_cursor_close(cursor); mdb_txn_abort(txn); return -1; |
||||
} |
||||
|
||||
// ---- SHA256 helpers ----
|
||||
static void db_sha256(const uint8_t* data, size_t len, uint8_t hash[32]) { |
||||
SC_SHA256_CTX ctx; sc_sha256_init(&ctx); sc_sha256_update(&ctx, data, len); sc_sha256_final(&ctx, hash); |
||||
} |
||||
|
||||
// Compute datahash: first 64 bits of SHA256(json_data)
|
||||
static uint64_t db_datahash(const uint8_t* data, size_t len) { |
||||
uint8_t hash[32]; db_sha256(data, len, hash); |
||||
uint64_t dh; memcpy(&dh, hash, 8); return dh; |
||||
} |
||||
|
||||
// Compute chain_hash for a record
|
||||
static void db_chain_hash_compute(const uint8_t prev_chain_hash[32], uint64_t id, uint64_t timestamp, uint64_t datahash, uint8_t out[32]) { |
||||
uint8_t buf[32+8+8+8]; // prev_chain_hash:32 + id:8 + timestamp:8 + datahash:8
|
||||
memcpy(buf, prev_chain_hash, 32); |
||||
memcpy(buf+32, &id, 8); memcpy(buf+40, ×tamp, 8); memcpy(buf+48, &datahash, 8); |
||||
db_sha256(buf, 56, out); |
||||
} |
||||
|
||||
// Get previous record's chain_hash (record before the given key in sorted order)
|
||||
static int db_prev_chain_hash(struct DB_SYNC* db, const uint8_t key[16], uint8_t chain_hash[32]) { |
||||
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 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; i<db->peer_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 && rc<send_count && off+val.mv_size+16 < 8000) { |
||||
uint64_t rts = be64toh(*(uint64_t*)key.mv_data); |
||||
uint64_t rdh = be64toh(*(uint64_t*)((uint8_t*)key.mv_data+8)); |
||||
uint64_t rrid = *(uint64_t*)((uint8_t*)val.mv_data+DB_VAL_OFF_ID); |
||||
uint32_t rdlen = *(uint32_t*)((uint8_t*)val.mv_data+DB_VAL_OFF_DATA_LEN); |
||||
const uint8_t* rdptr = (uint8_t*)val.mv_data+DB_VAL_OFF_DATA; |
||||
memcpy(sdbuf+off, &rrid,8); off+=8; memcpy(sdbuf+off, &rts,8); off+=8; |
||||
memcpy(sdbuf+off, &rdh,8); off+=8; memcpy(sdbuf+off, &rdlen,4); off+=4; |
||||
if (rdlen>0) { 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); |
||||
} |
||||
@ -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 <stdint.h> |
||||
#include <stddef.h> |
||||
|
||||
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
|
||||
@ -0,0 +1,244 @@
|
||||
// test_db_sync.c — всестороннее тестирование модуля db_sync
|
||||
#include <stdio.h> |
||||
#include <stdlib.h> |
||||
#include <string.h> |
||||
#include <stdarg.h> |
||||
#include <time.h> |
||||
#include <sys/stat.h> |
||||
#include "../lib/platform_compat.h" |
||||
#include "test_utils.h" |
||||
#ifdef _WIN32 |
||||
#include <windows.h> |
||||
#include <direct.h> |
||||
#else |
||||
#include <unistd.h> |
||||
#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; |
||||
} |
||||
@ -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 |
||||
Loading…
Reference in new issue