diff --git a/doc/route_p2pconn.txt b/doc/route_p2pconn.txt deleted file mode 100644 index 90c801b0..00000000 --- a/doc/route_p2pconn.txt +++ /dev/null @@ -1,344 +0,0 @@ -Архитектура P2P Direct Connection (установка оптимального подключения между узлами) - -Цель: в большой сети узлов каждый узел умеет находить быстрые подключения -к любому другому узлу — либо прямое (direct), либо через минимальное число -промежуточных узлов с минимальным суммарным RTT. Оба узла пробуют доступные -варианты пингом и выбирают лучший. - - -=== Что уже есть в коде (используем, не дублируем) === - -1. NODEINFO уже содержит поле tranzit_nodes + структуру NODEINFO_TRANZIT_NODE - (node_id, rtt, link_q) — "лучшие транзитные узлы, выбирается/обновляется узлом", - но пока не заполняется и не используется. - -2. Пинг-инфраструктура: route_ping_send_req_addr (запрос удалённого пинга через - третий узел), route_ping_handle_resp (приём результата). Также есть - etcp_send_ping_to_socket — прямой пинг с конкретного сокета на конкретный адрес - с pubkey для шифрования. - -3. NAT-детекция уже работает: проверка типа NAT через третий узел, обмен NAT_INFO. - -4. Path system: NODEINFO_PATH с hop_count, routes → ROUTE_ENTRY → NODEINFO_Q → paths — - уже выбирается путь для маршрутизации пакетов. - -5. Константа ROUTE_SUBCMD_CONN_REQ (0x0B) зарезервирована в route_bgp.h, но не реализована. - -6. NODEINFO_Q уже имеет best_socket, last_ping_time, last_rtt. - - -=== Общий алгоритм (на примере узлов A и B) === - -A ──(ETCP через релейные узлы)──> B [текущий путь, hop_count > 1] - -1. A замечает большой объём трафика к B → запускает CONN_REQ -2. A шлёт B через существующий ETCP-путь: свои адреса-кандидаты + лучших транзитных соседей -3. B получает CONN_REQ, сразу начинает пинговать адреса A, собирает свои данные -4. B шлёт CONN_RESP: свои адреса + транзитных соседей + результаты пробных пингов -5. A получает CONN_RESP, пингует все адреса B -6. Оба вычисляют лучший вариант (прямой или релейный с минимальным RTT) -7. Оба шлют друг другу CONN_RESULT с выбором -8. Оба пытаются создать ETCP link к выбранному адресу (если прямой) -9. Лишние/худшие линки закрываются позже — все линки в рамках одного ETCP_CONN - - -=== Фаза 1: Обмен кандидатами (CONN_REQ → CONN_RESP) === - -Узел-инициатор A собирает и отправляет: - - Свои адреса-кандидаты: из своих ETCP_SOCKET (interface_addr + nat_addr). - Включаются только сокеты с NAT_VERIFIED_* / EIM / PUBLIC (type >= NAT_VERIFIED_UNKNOWN). - Для каждого: ip, port, type, socket_id. - - Лучших транзитных соседей: top-N directly-connected узлов (hop_count == 1) - отсортированных по NODEINFO_Q.last_rtt (лучший RTT первым). - -Узел B при получении CONN_REQ: - - Сохраняет кандидатов A - - Сразу запускает пробные пинги ко всем адресам A (etcp_send_ping_to_socket с pubkey) - - Собирает свои адреса-кандидаты - - Собирает своих транзитных соседей - - Шлёт CONN_RESP со своими данными + уже готовыми результатами проб - - -=== Фаза 2: Пробы (CONN_RESP → Пинги) === - -Узел A при получении CONN_RESP: - - Запускает пинги ко всем адресам B: - каждый адрес B пингуется с каждого своего сокета (N×M проб) - - Использует etcp_send_ping_to_socket с pubkey узла B (из NODEINFO) - - Пинг идёт напрямую по UDP (не через ETCP) — проверяет реальную достижимость - включая NAT traversal - - Ждёт все результаты либо таймаут - -Параллельно B тоже завершает свои пробы (запущенные на шаге 3). - - -=== Фаза 3: Выбор лучшего пути (→ CONN_RESULT) === - -Каждая сторона независимо вычисляет: - - best = {mode: none, rtt: 65535} - - Прямые варианты: - for each (мой_сокет, peer_addr) где probe.ok: - if probe.rtt < best.rtt: - best = {direct, my_sock_idx, peer_addr_idx, probe.rtt} - - Релейные варианты: - for each общий_транзитный_узел R (доступен и A и B): - // RTT(A↔R) берём из своих замеров (A знает RTT до своих соседей) - // RTT(R↔B) берём из транзитных соседей B (присланы в CONN_RESP) - total_rtt = A→R_rtt + R→B_rtt - if total_rtt < best.rtt: - best = {relay, R, total_rtt} - -Обе стороны приходят к одинаковому выводу (информация симметрична). - -Шлют CONN_RESULT с выбором. -Если прямой вариант — оба пытаются создать ETCP link к адресу пира. -Если релейный — используют существующий путь (уже работает). - - -=== Триггер запуска negotiation === - -В route_pkt (routing.c) при отправке пакета узлу с hop_count > 1: - - Накапливаем счётчик байт к этому узлу (per-node counter в bgp) - - При превышении порога (например 64KB) и если negotiation ещё не запущен — стартуем - - Cooldown: не чаще чем раз в 30 секунд для одной пары узлов - - Если negotiation уже в процессе для этой пары — не дублируем - - -=== Tranzit nodes — автоматическое заполнение === - -В route_bgp_update_my_nodeinfo (route_node.c) при каждом обновлении NODEINFO: - - Сканируем directly-connected соседей (из bgp->nodes, где hop_count == 1) - - Сортируем по NODEINFO_Q.last_rtt (лучший RTT первым) - - Берём top-N (до 4) и упаковываем в NODEINFO.tranzit_nodes как массив - NODEINFO_TRANZIT_NODE (node_id, rtt, link_q) - -Плюс периодический refill (раз в ~10 сек) для свежести RTT замеров. - - -=== Структуры данных (новые, route_p2pconn.h) === - -#define P2P_MAX_ADDRS 8 -#define P2P_MAX_RELAYS 8 -#define P2P_MAX_PROBES (P2P_MAX_ADDRS * P2P_MAX_ADDRS) // до 64 - -#define P2P_PHASE_WAIT_RESP 0 // ждём CONN_RESP от пира -#define P2P_PHASE_PROBING 1 // пингуем адреса пира -#define P2P_PHASE_SELECTING 2 // выбор лучшего, отправка CONN_RESULT -#define P2P_PHASE_DONE 3 // завершено - -// Один адрес-кандидат -struct P2P_ADDR { - uint32_t ip; // network byte order - uint16_t port; - uint8_t type; // NAT_VERIFIED_* - uint8_t socket_id; -}; - -// Один транзитный узел -struct P2P_RELAY { - uint64_t node_id; - uint16_t rtt; // x0.1ms - uint16_t link_q; -}; - -// Результат одного пробного пинга -struct P2P_PROBE_RESULT { - uint8_t my_sock_idx; // индекс в my_addrs[] - uint8_t peer_addr_idx; // индекс в peer_addrs[] - uint16_t rtt; // x0.1ms, 0 = fail -}; - -// Состояние одних переговоров (хранится в хеш-таблице bgp->p2p_negotiations) -struct P2P_NEGOTIATION { - struct ll_entry ll; - uint32_t request_id; - uint64_t peer_node_id; - uint8_t phase; - - // Мои данные - struct P2P_ADDR my_addrs[P2P_MAX_ADDRS]; - uint8_t my_addr_count; - struct P2P_RELAY my_relays[P2P_MAX_RELAYS]; - uint8_t my_relay_count; - - // Данные пира (заполняются из CONN_RESP) - struct P2P_ADDR peer_addrs[P2P_MAX_ADDRS]; - uint8_t peer_addr_count; - struct P2P_RELAY peer_relays[P2P_MAX_RELAYS]; - uint8_t peer_relay_count; - - // Результаты проб - struct P2P_PROBE_RESULT probes[P2P_MAX_PROBES]; - uint8_t probe_total; // сколько всего запланировано - uint8_t probe_done; // сколько завершилось (ok + fail) - - // Лучший выбор - uint8_t best_mode; // 1=direct, 2=relay - uint64_t best_relay; // node_id релея (если mode=relay) - uint16_t best_rtt; // x0.1ms - uint8_t best_my_sock_idx; // индекс в my_addrs (если direct) - uint8_t best_peer_addr_idx;// индекс в peer_addrs (если direct) - - void* timeout_timer; // общий таймаут на всю negotiation -}; - -// Счётчик трафика per-node (для триггера, хранится в bgp) -struct P2P_TRAFFIC_COUNTER { - uint64_t node_id; - uint64_t bytes_sent; // накоплено байт - uint64_t last_negotiation_time; // время последней попытки (0 = не было) -}; - - -=== Протокольные пакеты (добавляются в route_bgp.h) === - -ROUTE_SUBCMD_CONN_REQ 0x0B // запрос прямого подключения (уже зарезервирован) -ROUTE_SUBCMD_CONN_RESP 0x0C // ответ с адресами + результаты проб -ROUTE_SUBCMD_CONN_RESULT 0x0E // финальный выбор - -// Пакет CONN_REQ (A → B) -struct BGP_CONN_REQ { - uint8_t cmd; // ETCP_ID_ROUTE_ENTRY - uint8_t subcmd; // ROUTE_SUBCMD_CONN_REQ - uint32_t request_id; // для корреляции - uint8_t addr_count; // число адресов-кандидатов - uint8_t relay_count; // число транзитных узлов - uint8_t reserved[2]; - // далее динамически (размер = addr_count*8 + relay_count*12): - // [addr_count × {ip[4] port[2] type[1] socket_id[1]}] - // [relay_count × {node_id[8] rtt[2] link_q[2]}] -}; - -// Пакет CONN_RESP (B → A) -struct BGP_CONN_RESP { - uint8_t cmd; // ETCP_ID_ROUTE_ENTRY - uint8_t subcmd; // ROUTE_SUBCMD_CONN_RESP - uint32_t request_id; - uint8_t addr_count; // адреса B - uint8_t relay_count; // транзитные узлы B - uint8_t probe_count; // готовые результаты проб B→A - uint8_t reserved; - // [addr_count × {ip[4] port[2] type[1] socket_id[1]}] - // [relay_count × {node_id[8] rtt[2] link_q[2]}] - // [probe_count × {peer_addr_idx[1] my_sock_idx[1] rtt[2] ok[1]}] -}; - -// Пакет CONN_RESULT (A ↔ B, финальный) -struct BGP_CONN_RESULT { - uint8_t cmd; - uint8_t subcmd; // ROUTE_SUBCMD_CONN_RESULT - uint32_t request_id; - uint8_t chosen_mode; // 1=direct, 2=relay - uint8_t reserved; - uint16_t chosen_rtt; // x0.1ms - // Если direct: - uint32_t peer_ip; // IP пира к которому подключаемся - uint16_t peer_port; - uint8_t my_socket_id; // свой сокет для подключения - uint8_t peer_socket_id; // сокет пира - // Если relay (оверлей тех же байт): - // uint64_t relay_node_id; -}; - - -=== API модуля route_p2pconn === - -// Запуск negotiation к узлу peer_node_id. -// Вызывается из route_pkt при превышении порога трафика. -int p2p_start_negotiation(struct ROUTE_BGP* bgp, uint64_t peer_node_id); - -// Проверка: запущена ли уже negotiation для этой пары -int p2p_is_negotiating(struct ROUTE_BGP* bgp, uint64_t peer_node_id); - -// Сбор своих адресов-кандидатов (из etcp_sockets) -int p2p_collect_my_addrs(struct ROUTE_BGP* bgp, struct P2P_ADDR* out, uint8_t max); - -// Сбор лучших транзитных соседей (из directly-connected nodes) -int p2p_collect_my_relays(struct ROUTE_BGP* bgp, struct P2P_RELAY* out, uint8_t max); - -// Обработчики входящих пакетов (вызываются из route_bgp_receive_cbk): -void p2p_handle_conn_req(struct ROUTE_BGP* bgp, struct ETCP_CONN* from, - const uint8_t* data, size_t len); -void p2p_handle_conn_resp(struct ROUTE_BGP* bgp, struct ETCP_CONN* from, - const uint8_t* data, size_t len); -void p2p_handle_conn_result(struct ROUTE_BGP* bgp, struct ETCP_CONN* from, - const uint8_t* data, size_t len); - -// Отмена negotiation при удалении conn -void p2p_cancel_for_conn(struct ROUTE_BGP* bgp, struct ETCP_CONN* conn); - -// Очистка всех negotiation (при destroy bgp) -void p2p_destroy_all(struct ROUTE_BGP* bgp); - - -=== Схема состояний negotiation === - - IDLE ──(трафик > порог)──→ WAIT_RESP ──(CONN_RESP получен)──→ PROBING - ↑ ↑ │ - │ (таймаут) (все пробы готовы) - │ ↓ │ - └──────────────────────────── DONE ←──(CONN_RESULT отправлен)── SELECTING - - -=== Интеграция (какие файлы меняются) === - -Новые файлы: - src/route_p2pconn.h — структуры, константы, прототипы - src/route_p2pconn.c — вся логика negotiation - -Изменения в существующих: - src/route_bgp.h — добавить ROUTE_SUBCMD_CONN_RESP 0x0C, - ROUTE_SUBCMD_CONN_RESULT 0x0E, - структуры пакетов BGP_CONN_REQ/RESP/RESULT - src/route_bgp.c — в route_bgp_receive_cbk добавить обработку - новых subcmd. В struct ROUTE_BGP добавить - поле p2p_negotiations (очередь/хеш negotiation) - и p2p_traffic_counters (per-node counters) - src/route_node.c — в route_bgp_update_my_nodeinfo заполнять - tranzit_nodes из RTT directly-connected соседей - src/routing.c — в route_pkt добавить накопление счётчика трафика - и вызов p2p_start_negotiation при превышении порога - src/Makefile.am — добавить route_p2pconn.c в сборку - - -=== Внутренняя логика route_p2pconn.c === - -p2p_start_negotiation(bgp, peer_node_id): - 1. Проверить что нет активной negotiation для этой пары - 2. Проверить cooldown (30 сек) - 3. Создать struct P2P_NEGOTIATION, заполнить my_addrs, my_relays - 4. Сформировать BGP_CONN_REQ и отправить через существующий путь к пиру - 5. Поставить таймаут на всю negotiation (например 5 сек) - -p2p_handle_conn_req(bgp, from_conn, data, len): - 1. Распарсить BGP_CONN_REQ - 2. Сохранить адреса и релеи инициатора в P2P_NEGOTIATION - 3. Собрать свои my_addrs, my_relays - 4. Запустить пробные пинги к адресам инициатора (p2p_start_probes) - 5. Когда пробы готовы — отправить CONN_RESP с результатами - -p2p_start_probes(neg): - 1. Для каждой пары (мой_сокет, peer_addr) запланировать пинг - 2. Вызвать etcp_send_ping_to_socket с pubkey пира - 3. В коллбэке сохранить результат в probes[], инкрементировать probe_done - 4. Когда probe_done == probe_total → вызвать p2p_select_best - -p2p_select_best(neg): - 1. Прямые: найти пару (my_sock, peer_addr) с минимальным rtt > 0 - 2. Релейные: для каждого общего транзитного узла посчитать суммарный rtt, - найти минимум - 3. Сравнить прямой vs релейный лучший rtt - 4. Записать выбор в neg->best_* - 5. Отправить CONN_RESULT пиру - 6. Вызвать p2p_establish_link(neg) - -p2p_establish_link(neg): - 1. Если best_mode == direct: - - Найти/создать ETCP_CONN к пиру (по node_id) - - Вызвать etcp_link_new с chosen address - 2. Если best_mode == relay: - - Ничего не делаем, существующий путь уже работает - 3. Пометить negotiation как DONE diff --git a/src/_db_arch.txt b/src/_db_arch.txt deleted file mode 100644 index ccc6b83f..00000000 --- a/src/_db_arch.txt +++ /dev/null @@ -1,22 +0,0 @@ -Архитектура таблицы 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/chat/chat_headless_control.c b/src/chat/chat_headless_control.c index 4593e406..85e0e7d4 100644 --- a/src/chat/chat_headless_control.c +++ b/src/chat/chat_headless_control.c @@ -13,6 +13,7 @@ #include "../ntp_time.h" #include "../transport_layer/etcp.h" #include "../transport_layer/etcp_connections.h" +#include "../transport_layer/secure_channel.h" #include "../../lib/u_async.h" #include "../../lib/debug_config.h" #include "../../lib/mem.h" @@ -409,6 +410,54 @@ static void hc_handle_create_channel(struct headless_client* cli, int id, const send_response(cli, id, "{\"created\":true}", NULL); } +static void hc_handle_invite_to(struct headless_client* cli, int id, const char* json) { + char ch[64], node_id_str[32], pubkey_hex[128], addr[128]; + if (json_get_str(json, "ch", ch, sizeof(ch)) < 0) { send_response(cli, id, NULL, "missing 'ch' param"); return; } + if (json_get_str(json, "node_id", node_id_str, sizeof(node_id_str)) < 0) { send_response(cli, id, NULL, "missing 'node_id' param"); return; } + if (json_get_str(json, "pubkey", pubkey_hex, sizeof(pubkey_hex)) < 0) { send_response(cli, id, NULL, "missing 'pubkey' param"); return; } + if (json_get_str(json, "addr", addr, sizeof(addr)) < 0) { send_response(cli, id, NULL, "missing 'addr' param"); return; } + int proto = json_get_int(json, "proto", 1); + if (!g_hc.inst) { send_response(cli, id, NULL, "no instance"); return; } + + uint64_t target_node_id = strtoull(node_id_str, NULL, 16); + if (target_node_id == 0 || target_node_id == g_hc.inst->node_id) { send_response(cli, id, NULL, "invalid node_id"); return; } + + uint8_t pubkey_bin[SC_PUBKEY_SIZE]; + if (sc_hex_to_binary(pubkey_hex, pubkey_bin, SC_PUBKEY_SIZE) != 0) { send_response(cli, id, NULL, "invalid pubkey hex"); return; } + + char ip[128]; snprintf(ip, sizeof(ip), "%s", addr); + char* colon = strrchr(ip, ':'); + if (!colon) { send_response(cli, id, NULL, "invalid addr (no port)"); return; } + *colon = '\0'; + int port = (int)strtol(colon + 1, NULL, 10); + if (port <= 0 || port > 65535) { send_response(cli, id, NULL, "invalid port"); return; } + + uint8_t addrs_data[21]; int addrs_len = 0; + struct in_addr a4; struct in6_addr a6; + if (inet_pton(AF_INET, ip, &a4) == 1) { + addrs_data[0] = 4; addrs_data[1] = 0; addrs_data[2] = (uint8_t)proto; + memcpy(addrs_data + 3, &a4.s_addr, 4); + addrs_data[7] = (uint8_t)(port >> 8); addrs_data[8] = (uint8_t)(port & 0xFF); + addrs_len = 9; + } else if (inet_pton(AF_INET6, ip, &a6) == 1) { + addrs_data[0] = 6; addrs_data[1] = 0; addrs_data[2] = (uint8_t)proto; + memcpy(addrs_data + 3, &a6, 16); + addrs_data[19] = (uint8_t)(port >> 8); addrs_data[20] = (uint8_t)(port & 0xFF); + addrs_len = 21; + } else { send_response(cli, id, NULL, "invalid addr ip"); return; } + + DEBUG_INFO((int)DEBUG_CATEGORY_HEADLESS, "headless: invite_to ch=%s node=0x%016llx addr=%s:%d proto=%d", + ch, (unsigned long long)target_node_id, ip, port, proto); + + chat_sync_invite_to_channel_with_addrs(g_hc.inst, ch, target_node_id, + pubkey_bin, addrs_data, 1, addrs_len); + + char resp[256]; snprintf(resp, sizeof(resp), + "{\"channel_id\":\"%s\",\"node_id\":\"0x%016llx\",\"inviting\":true}", + ch, (unsigned long long)target_node_id); + send_response(cli, id, resp, NULL); +} + static void hc_handle_subscribe(struct headless_client* cli, int id, const char* json) { int enable = json_get_int(json, "enable", 1); cli->subscribed = enable; @@ -427,8 +476,8 @@ static void hc_handle_quit(struct headless_client* cli, int id, const char* json typedef void (*hc_cmd_fn)(struct headless_client* cli, int id, const char* json); -static const char* cmd_names[] = { "ping", "status", "channels", "members", "messages", "send", "invite", "connect", "create_channel", "subscribe", "quit", NULL }; -static hc_cmd_fn cmd_handlers[] = { hc_handle_ping, hc_handle_status, hc_handle_channels, hc_handle_members, hc_handle_messages, hc_handle_send, hc_handle_invite, hc_handle_connect, hc_handle_create_channel, hc_handle_subscribe, hc_handle_quit }; +static const char* cmd_names[] = { "ping", "status", "channels", "members", "messages", "send", "invite", "connect", "create_channel", "invite_to", "subscribe", "quit", NULL }; +static hc_cmd_fn cmd_handlers[] = { hc_handle_ping, hc_handle_status, hc_handle_channels, hc_handle_members, hc_handle_messages, hc_handle_send, hc_handle_invite, hc_handle_connect, hc_handle_create_channel, hc_handle_invite_to, hc_handle_subscribe, hc_handle_quit }; void hc_handle_command(struct headless_client* cli, const char* json) { if (!cli || !json) return; diff --git a/src/db_sync_doc.md b/src/db_sync_doc.md deleted file mode 100644 index c49d2e7d..00000000 --- a/src/db_sync_doc.md +++ /dev/null @@ -1,224 +0,0 @@ -# db_sync — Distributed append-only table with SQLite + P2P sync - -## 1. Назначение - -Реплицировать append-only таблицу JSON-записей между всеми пирами P2P-сети. Каждый пир в итоге должен иметь идентичный набор записей. Модуль поддерживает несколько независимых инстансов (таблиц), каждый идентифицируется хешем `SHA256(name || id_be)[0:8]`. - -## 2. Ключевые свойства - -- **Криптографическая цепь.** `chain_hash[N] = SHA256(chain_hash[N-1] || id || timestamp || author || author_signature)`. Записи упорядочены `ORDER BY timestamp, author_signature`. Первые 8 байт — `chain_hash8` — используется для быстрого сравнения в протоколе. -- **Ed25519-подписи.** Каждая запись подписана автором. Записи без подписи или с неверной подписью отвергаются. -- **Append-only.** Записи не редактируются и не удаляются явно, только TTL-очистка собственных неотправленных записей. -- **PUSH — мгновенная доставка.** При локальном `insert_signed` запись немедленно шлётся всем пирам с `sync_state >= 1` (steady state). -- **Sync — сравнение цепей.** При старте/реконнекте стороны обмениваются хешами цепей для поиска расхождений. -- **Multi-instance.** Несколько независимых таблиц внутри одного процесса (разные hash). - -## 3. Архитектура протокола - -### 3.1. На что оптимизирован - -99% времени новые записи просто дописываются в конец. Самый частый сценарий — пир A добавил сообщение, пир B получил его и вставил в конец своей цепи. Протокол должен: -- Доставлять новые записи немедленно (PUSH) -- При реконнекте быстро понять "у нас всё совпадает до позиции N, добрось хвост" (hash_MATCH) -- При реальном расхождении бинарным поиском найти точку и слить (divergence + REFINE) - -### 3.2. Общие идеи реализации - -**Криптографическая цепь** — гарантия целостности. Если `chain_hash8(N)` совпадает на двух пирах, цепь идентична до позиции N (вероятность коллизии 2^-64). При вставке не в конец — каскадный пересчёт `chain_hash` всех последующих записей через `db_cascade_from`. - -**PUSH** — рабочий механизм доставки в steady state. Запись, вставленная локально, немедленно уходит всем синхронизированным пирам. ACK_PUSH подтверждает получение и обновляет delivery_chain. - -**Sync** — полное сравнение цепей при старте/реконнекте. Три ветки: peer_empty (пир пуст — отдать всё), hash_MATCH (цепи совпали до позиции N — отдать хвост), divergence (цепи разошлись — найти точку расхождения). - -**Sparse checkpoints** — бинарный поиск расхождения. INIT_RESP возвращает хеши в степенях двойки от tp (tp-1, tp-2, tp-4, ..., до 16 шт). Это позволяет за O(log N) сравнений сузить диапазон, не передавая хеш каждой записи. - -**REFINE** — финальное сужение. Если sparse-хешей INIT_RESP недостаточно (диапазон >1), REFINE запрашивает до 16 дополнительных хешей. Когда диапазон ≤1 — сразу SEND_DATA. - -**Каскадное уведомление.** При вставке записи не в конец `synced_pos` всех остальных пиров сбрасывается до позиции вставки — им потребуется пересинхронизация. - -### 3.3. Фазы протокола - -**Фаза А — Инициализация.** `db_sync_instance_add` читает таблицу из SQLite, проверяет целостность цепи (`db_verify_chain` — автофикс при расхождении), и немедленно шлёт `INIT_SYNC(my_count)` всем подключённым пирам с `sync_state == 0`. Если связь появилась позже — `conn_up` делает то же самое. - -**Фаза Б — Сравнение цепей.** Получатель INIT_SYNC вычисляет `tp = min(my_count, peer_count)`, tp-- если >0 (последняя гарантированно общая позиция), и возвращает INIT_RESP: `peer_ch8` на позиции tp, `sc` (количество sparse-хешей), sparse-хеши на позициях tp-2^k. - -**Фаза В — Три ветки:** - -| Ветка | Условие | Сценарий | Действие | -|-------|---------|----------|----------| -| peer_empty | `peer_ch8==0 && sc==0` | Пир пуст (0 записей) | Отправить все свои записи через SEND_DATA | -| hash_MATCH | `my_ch8 == peer_ch8` | Цепи идентичны до tp | Отправить хвост [tp+1..mc) через SEND_DATA | -| divergence | `my_ch8 != peer_ch8` | Разные истории | Анализ sparse-хешей → REFINE → SEND_DATA | - -**Фаза Г — REFINE.** Инициатор анализирует sparse-хеши: `ds` — последняя совпавшая позиция, `de` — первая разошедшаяся. Если `de-ds ≤ 1` → сразу SEND_DATA (hc=0). Иначе → REFINE с до 16 своих хешей, равномерно распределённых в [ds..de]. Получатель сравнивает со своей цепью, находит точку совпадения, шлёт SEND_DATA от этой точки. - -**Фаза Д — SEND_DATA.** Получатель вставляет записи с проверкой Ed25519-подписи, делает `db_cascade_from(fix_from)`, уведомляет остальных пиров о сдвиге цепи (сброс их synced_pos). Если у получателя после вставки записей больше чем у отправителя — proactive push-back (шлёт свой хвост). Когда все записи получены → SYNC_DONE с итоговым count и chain_hash8. - -**Фаза Е — SYNC_DONE.** Сравнение итогового count и chain_hash8. Не совпало — ретрай всей процедуры с начала (до 3 раз, потом give up с partial sync). Совпало — `sync_state=2`, синхронизация завершена. - -### 3.4. PUSH — отдельный от sync механизм - -PUSH матчится по `hash` — если у пира нет si с таким же hash, PUSH не доставляется. Это позволяет изолировать тестирование sync-протокола от PUSH: вставлять данные через si с уникальным hash (PUSH не уходит — нет получателя), затем удалять tmp si (данные в SQLite сохраняются), создавать si с общим hash — instance_add запускает чистый sync. - -## 4. Peer management - -- При поднятии ETCP-соединения для каждого инстанса добавляется `SI_PEER` и запускается синхронизация. -- При разрыве соединения `sync_state` пира сбрасывается в 0. -- `peer_check` таймер (каждые 5с) перебирает активных пиров и запускает синхронизацию для тех, у кого `sync_state == 0`. Выбирается пир с минимальным `synced_pos` — двигаемся от самого старого несинхронизированного участка. -- PUSH рассылается только пирам в состоянии `sync_state >= 1`. - -## 5. Сообщения протокола - -| Сообщение | Wire-формат | Описание | -|-----------|-------------|----------| -| `INIT_SYNC (0x01)` | `[type:1][count:4]` | Инициатор шлёт количество записей | -| `INIT_RESP (0x02)` | `[type:1][tp:4][ch8:8][sc:1][(pos:4,ch8:8)*sc]` | tp, chain_hash8 на tp, sc sparse-хешей | -| `REFINE (0x03)` | `[type:1][from:4][to:4][hc:1][(pos:4,ch8:8)*hc]` | hc=0 → сразу SEND_DATA; hc>0 → до 16 хешей в диапазоне | -| `SEND_DATA (0x04)` | `[type:1][from:4][count:2][(id:8,ts:8,author:8,dlen:4,data,sig_len:1,sig)*count]` | Пакет до 32 записей | -| `PUSH (0x05)` | `[type:1][id:8][ts:8][author:8][dlen:4][data][sig_len:1=64][sig:64]` | Рассылка одной записи synced-пирам | -| `ACK_PUSH (0x06)` | `[type:1][ts:8][author:8]` | Подтверждение PUSH, обновление delivery_chain | -| `SYNC_DONE (0x07)` | `[type:1][count:4][ch8:8]` | Финальный count + chain_hash8 | -| `ERROR (0x08)` | `[type:1][code:1]` | Коды: 0x01=NOT_FOUND, 0x02=DISABLED | - -Все сообщения маршрутизируются через ETCP service `0x20` с префиксом `[svc:1][hash_be:8][payload]`. - -## 6. Wire-формат записи - -`[id:8][ts:8][author_node_id:8][dlen:4][data:json][sig_len:1=64][sig:ed25519:64]` - -Дубликаты определяются по `(timestamp, author_signature)`. - -## 7. Структуры данных - -```c -struct DB_SYNC { - struct UTUN_INSTANCE* inst; - sqlite3* db; - uint8_t shared_db; - uint64_t last_connected_tb; - struct DB_SYNC_INSTANCE* instances; - int instance_count, instance_capacity; - void* peer_check_timer; - uint8_t enabled; -}; - -struct DB_SYNC_INSTANCE { - struct DB_SYNC* db_sync; - uint64_t hash; - char table_name[64]; - uint64_t next_id; - uint64_t last_timestamp_ms; - uint8_t enabled; - void* ttl_timer; - struct SI_PEER* peers; - int peer_count, peer_capacity; - db_sync_insert_cb on_insert; - void* on_insert_arg; -}; - -struct SI_PEER { - uint64_t node_id; - uint32_t synced_pos; // последняя общая позиция (0-based) - uint8_t sync_state; // 0=idle, 1=syncing, 2=synced - uint8_t sync_retry_count; // счётчик retry SYNC_DONE mismatch - uint64_t sync_start_tb; // время начала sync (для timeout) -}; -``` - -## 8. API - -### Глобальный жизненный цикл - -```c -int db_sync_init(struct UTUN_INSTANCE* inst); -void db_sync_destroy(struct UTUN_INSTANCE* inst); -``` - -`db_sync_init` — открывает SQLite по пути `/chats.db`, биндит ETCP service `0x20`, вешает коллбэки соединений, стартует `peer_check` таймер (5с). Если `db_sync_enabled = 0` — disabled-режим. - -`db_sync_destroy` — отменяет таймеры, анбиндит сервис, снимает коллбэки, закрывает SQLite. - -### Управление инстансами - -```c -struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const char* table_name, uint64_t hash); -void db_sync_instance_remove(struct DB_SYNC_INSTANCE* si); -``` - -`db_sync_instance_add` — создаёт/регистрирует инстанс. Создаёт SQLite-таблицу (если нет), читает цепь, `db_verify_chain` (автофикс), стартует TTL-таймер. Если есть подключённые пиры с `sync_state == 0` — немедленно шлёт `INIT_SYNC`. - -`db_sync_instance_remove` — деактивирует: `enabled=0`, отменяет TTL-таймер, освобождает peers, удаляет из массива. Таблица БД не удаляется. - -### Операции с данными - -```c -int db_sync_insert_signed(struct DB_SYNC_INSTANCE* si, const char* json_data, size_t len, - const uint8_t* sig, size_t sig_len, uint64_t ts); -uint32_t db_sync_count(struct DB_SYNC_INSTANCE* si); -uint64_t db_sync_get_last_timestamp(struct DB_SYNC_INSTANCE* si); -uint64_t db_sync_next_timestamp(struct DB_SYNC_INSTANCE* si); -``` - -`db_sync_insert_signed` — вставляет подписанную запись. `sig` — 64 байта Ed25519. Проверяет подпись, проверяет дубликат, вычисляет chain_hash, вставляет с каскадным пересчётом, рассылает PUSH synced-пирам. Возвращает: 0=успех, 1=дубликат, -1=ошибка, -2=неверная подпись. - -### Чтение - -```c -typedef void (*db_sync_select_cb)(void* arg, uint64_t id, uint64_t timestamp, - const char* data, size_t data_len, uint64_t author, - const uint8_t* author_sig, size_t sig_len, - int delivered_peers, const char* delivery_chain); -int db_sync_select(struct DB_SYNC_INSTANCE* si, uint32_t offset, uint32_t limit, - db_sync_select_cb cb, void* arg); -``` - -### Callback на вставку - -```c -typedef void (*db_sync_insert_cb)(struct DB_SYNC_INSTANCE* si, - uint64_t record_timestamp, - const char* json_data, size_t len, - uint64_t author_node_id, void* arg); -void db_sync_set_insert_cb(struct DB_SYNC_INSTANCE* si, db_sync_insert_cb cb, void* arg); -``` - -### Верификация цепи - -```c -int db_sync_chain_verify(struct DB_SYNC_INSTANCE* si); -``` - -Возвращает 0 если все chain_hash корректны, 1 если расхождение. - -## 9. Конфигурация - -| Параметр | По умолчанию | Описание | -|----------|-------------|----------| -| `db_sync_enabled` | `0` | Включить модуль синхронизации | -| `db_sync_ttl` | `86400` | TTL неотправленных собственных записей (сек) | -| `db_path` | — | Путь к БД; файл `/chats.db` | - -## 10. Константы - -| Константа | Значение | Описание | -|-----------|---------|----------| -| `ETCP_RT_ID_DB_SYNC` | `0x20` | ETCP service ID | -| `DB_REFINE_HASHES` | `16` | Макс. хешей в REFINE | -| `DB_SEND_DATA_MAX` | `32` | Макс. записей в SEND_DATA | -| `DB_SIG_SIZE` | `64` | Размер Ed25519 подписи | -| `DB_SYNC_PEER_CHECK_INTERVAL` | `5` | Интервал проверки пиров (сек) | -| `DB_SYNC_SYNC_TIMEOUT` | `15` | Таймаут ожидания ответа при sync (сек) | -| `DB_SYNC_TTL_INTERVAL` | `3600` | Интервал TTL-очистки (сек) | - -## 11. Тестирование sync-протокола - -PUSH доставляет записи немедленно и независимо от sync. Чтобы тестировать чистый sync-протокол, нужно чтобы PUSH не вмешивался: - -1. Вставить данные через si с **уникальным hash** — PUSH уходит, но не доставляется (нет получателя с таким hash) -2. `remove_si()` — данные сохраняются в SQLite, si деактивирован -3. Создать si с **общим hash** на обеих сторонах — `instance_add` запускает чистый sync-протокол - -Этот паттерн используется для тестирования всех трёх веток INIT_RESP: -- **peer_empty**: одна сторона с данными, другая пустая -- **hash_MATCH**: обе имеют общий префикс, но у одной больше записей -- **divergence**: независимые вставки на обеих сторонах до синхронизации diff --git a/src/req.txt b/src/req.txt deleted file mode 100644 index 8bd041c7..00000000 --- a/src/req.txt +++ /dev/null @@ -1,29 +0,0 @@ -Возможности utun: - -- роутинг между нодами - напрямую (предпочтительный вариант), или через наилучшего кандидата. -Наилучший кандидат - узел с минимальным overall score. -overall score = правило выбираемок пользователем - % loss multiplier: - <0.3% - x1 - <1% - x0.5 - <5% - x0.2 - <10% - x0.1 - > - x0.03 - - либо minimum delay - -Формат данных: -todo... - -tun -> (src/dst ip) -> routing table -> next hop -> transmit to -> etcp -> routing -> tun - -routing: выбирает по dst ip next hop. -next hop = {type + conrol struct}. type= {1-tun, 2-etcp} - - - - -control dgram = connect request - если нет прямого соединения -при установке соединения сервер отправляет ip:port клиенту - -routing dgram: - -> diff --git a/src/routing_layer/_route_tz.txt b/src/routing_layer/_route_tz.txt deleted file mode 100644 index 33a1b530..00000000 --- a/src/routing_layer/_route_tz.txt +++ /dev/null @@ -1,9 +0,0 @@ -Задача сделать механизм маршрутизации между узлами для групповых чатов. - -Узлы для всех чатов находятся в одной таблице - это nodeinfo (информация о подключении) -В другой таблице есть список узлов конкретного чата. - -надо сделать возможность хранить одновременно несколько групп узлов. в группу узлов добавить тип - utun или чатовая. -если utun - в группе используется роутинг (ipv4/ipv6). -если чатовая - то только связность между узлами. - diff --git a/src/routing_layer/_todo.txt b/src/routing_layer/_todo.txt deleted file mode 100644 index 7ebb4661..00000000 --- a/src/routing_layer/_todo.txt +++ /dev/null @@ -1,3 +0,0 @@ -При добавлении узла группы в БД дополнительно -- проверяем текущие подключения -- если узел найден в активных подключениях то добавляем его в группу topo_group diff --git a/src/routing_layer/conn_mgr_core.c b/src/routing_layer/conn_mgr_core.c index 24628bba..f5c4be24 100644 --- a/src/routing_layer/conn_mgr_core.c +++ b/src/routing_layer/conn_mgr_core.c @@ -347,6 +347,13 @@ int cm_open(struct CONN_MGR* mgr, uint64_t node_id, uint32_t idle_timeout_ms, if (!mgr) { if (cb) cb(NULL, node_id, 0, CONN_EVENT_TIMEOUT, cb_arg); return -1; } struct TOPO_GROUP* group = mgr->group; if (!group) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "BGP not initialized"); return -1; } + if (node_id == 0) { + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "conn_mgr: cm_open node_id=0 — rejected grp=%016llx type=%d ch=%s", + (unsigned long long)group->group_id, group->group_type, + group->channel_id[0] ? group->channel_id : "-"); + if (cb) cb(NULL, node_id, 0, CONN_EVENT_TIMEOUT, cb_arg); + return -1; + } uint64_t gid = group->group_id; struct TOPO_GROUP_NODE* target = topo_node_find_by_id(group, node_id); int loaded_from_db = 0; diff --git a/src/routing_layer/route6_lib_doc.md b/src/routing_layer/route6_lib_doc.md deleted file mode 100644 index 30562a7e..00000000 --- a/src/routing_layer/route6_lib_doc.md +++ /dev/null @@ -1,50 +0,0 @@ -# route6_lib — библиотека IPv6-маршрутизации на radix tree - -## 1. Назначение -Локальная таблица маршрутизации IPv6 для поиска longest-prefix match по заданному адресу. -Используется при пробросе трафика в TUN-интерфейс — по destination IP определяем узел-владелец -подсети и отправляем пакет ему через ETCP. - -Индексируется по `node_id` через `ll_queue`-хеш — это позволяет эффективно удалять все маршруты -узла при его withdraw из BGP-топологии. - -## 2. Как пользоваться -```c -struct ROUTE6_TABLE *rt = route6_table_create(ua); - -// При получении BGP-апдейта с новыми подсетями узла: -route6_insert(rt, nodeq); - -// При withdraw узла: -route6_delete(rt, nodeq); - -// При пробросе IPv6-пакета: -const struct ROUTE6_DATA *rd = route6_lookup(rt, dst_addr); -if (rd) forward_to_etcp(rd->v_node_info, pkt); -``` - -Ключевой нюанс: `make_key6` формирует 17-байтовый ключ в формате sockaddr-like -(первый байт = длина = 17, остальные 16 — адрес). Radix tree сравнивает ключи с отступом 1 байт -(RADIX_OFF), пропуская заголовок длины. - -## 3. API - -### Структуры -| Структура | Назначение | -|-----------|------------| -| `ROUTE6_DATA` | Запись маршрута: ключ+маска (17 байт sockaddr-like), `prefix_length`, `node_id` (для хеша), `v_node_info` (указатель на узел-владелец), `rn_nodes[2]` (scratch для radix) | -| `ROUTE6_TABLE` | Таблица: `rnh` (radix tree head), `queue` (ll_queue с хеш-индексом по `node_id`) | - -### Функции -| Функция | Назначение | -|---------|------------| -| `route6_table_create(ua)` | Создать таблицу: `rn_inithead` + `queue_new` с хеш-индексом по `node_id` | -| `route6_table_destroy(table)` | Удалить все маршруты (radix+queue), освободить таблицу | -| `route6_insert(table, node)` | Вставить все `v6_subnets` узла в radix tree и queue (с хеш-индексом) | -| `route6_delete(table, node)` | Удалить из radix tree все маршруты с заданным `node_id` (поиск через хеш) | -| `route6_lookup(table, addr)` | Longest-prefix match по 16-байтовому адресу, возвращает `ROUTE6_DATA*` или NULL | - -### Зависимости -- `lib/radix.h` — radix tree (longest-match lookup) -- `lib/ll_queue.h` — lock-free очередь с хеш-индексом для поиска по `node_id` -- `src/topo_node.h` — `TOPO_GROUP_NODE`, `TOPO_SUBNET6` diff --git a/src/routing_layer/route_bgp.txt b/src/routing_layer/route_bgp.txt deleted file mode 100644 index 1175869d..00000000 --- a/src/routing_layer/route_bgp.txt +++ /dev/null @@ -1,42 +0,0 @@ -Ключевые изменения в понимании - -paths - -route_bgp_process_nodeinfo(bgp, from_conn, data, entry->len): - проверяет текущую версию nodeinfo, если совпадает - только bgp_update local nodelist - если версия новая или нет узла - spread: - bgp_update local nodelist - realloc: добавляет next hop: hoplist += prev_node_id, hop_count++ - bgp_spread - -1. функция распространения маршрута -bgp_spread(struct ROUTE_BGP bgp, struct NODEINFO* n) - send to: - - все подключения, если не найден uid подключения в hoplist - - -bgp_update local nodelist: - - обновляем саму node - - удалеям conn где lash hop=prev_node_id - - добавляем новый conn с новым hoplist - -2. withdraw: -bgp_withdraw(struct ROUTE_BGP bgp, uint64_t node_to_del, uint64_t wd_source) - wd_source это узел который захотел withdraw. - - находим у себя node to del. удаляем если в hoplist найден wd_node (или мы = wd_node) - если удалили - распространяем по всем линкам с этими же аргументами - - - -=============================== -определение nat - -определяем что оба узла в локальной сети: - -если обе ноды имеют локальные адреса сокетов пробуем пинг по локальным адресам - -если узел подключился и его сокетовый ip:port отличается от connection ip:port - устанавливаем флаг nat. и инициируем проверку пингом третьим узлом: -проверка считается успешной если по результату пинга ip источника пинга отличается от ip запросившего пинг узла. - -проверка типа nat узла который подключился: -- пингуем узел с целью зафиксировать свой ip (в сокете может быть например 0.0.0.0) -- даём команду третьему узлу пингануть этот узел. -- по результатам двух пингов определяем тип nat diff --git a/src/routing_layer/route_lib_doc.md b/src/routing_layer/route_lib_doc.md deleted file mode 100644 index 00751cdb..00000000 --- a/src/routing_layer/route_lib_doc.md +++ /dev/null @@ -1,68 +0,0 @@ -# route_lib – библиотека таблицы IPv4-маршрутизации - -## 1. Назначение -Ведёт таблицу IPv4-маршрутов в памяти: динамический массив `ROUTE_ENTRY`, отсортированный по network-адресу с поддержкой вставки/удаления и LPM-поиска (Longest Prefix Match). Используется модулем `routing` для принятия решения, куда слать пакет — какому узлу через ETCP или локально в TUN. - -Таблица не содержит логики синхронизации — наполнение маршрутами идёт через BGP (`topo_node`). - -## 2. Как пользоваться -```c -struct ROUTE_TABLE* rt = route_table_create(); - -// Вставка подсетей узла (из BGP): -route_insert(rt, topo_nodeq); - -// Поиск маршрута (LPM): -struct ROUTE_ENTRY* e = route_lookup(rt, dest_ip); -if (e && e->v_node_info && e->v_node_info->hop_count > 0) { - // Отправить через ETCP узлу e->v_node_info->node->node_id -} else { - // Локальный маршрут — отдать в TUN -} - -// Удаление маршрутов узла: -route_delete(rt, topo_nodeq); - -route_table_destroy(rt); -``` - -**Ключевые нюансы:** -- Адреса в `ROUTE_ENTRY.network` хранятся в **big-endian** (сетевой порядок) -- Записи отсортированы по `network` ASC, при равенстве — по `prefix_length` DESC (более специфичные раньше) -- Алгоритм LPM: бинарный поиск с fallback-проверкой соседней записи -- `v_node_info == NULL` или `hop_count == 0` означает **локальный маршрут** -- Пересечение маршрутов запрещено — `route_insert` проверяет `check_route_overlap_in_table` -- Ёмкость массива динамически расширяется (×2) при переполнении - -## 3. API - -### Структуры - -| Структура | Поля | Описание | -|-----------|------|----------| -| `ROUTE_ENTRY` | `network` (uint32_t, BE), `prefix_length` (uint8_t), `v_node_info` (TOPO_GROUP_NODE*) | Запись маршрута. `v_node_info==NULL` → локальный | -| `ROUTE_TABLE` | `entries`, `count`, `capacity`, `dynamic_subnets`, `local_subnets`, `stats` | Таблица маршрутизации. Динамический массив + статистика | - -### Функции - -| Функция | Описание | -|---------|----------| -| `route_table_create()` | Выделяет `ROUTE_TABLE`, `entries` на `INITIAL_CAPACITY=100`, счётчики в 0 | -| `route_table_destroy(table)` | Освобождает `entries`, `dynamic_subnets`, `local_subnets`, саму таблицу | -| `route_insert(table, node)` | Вставляет все v4-подсети узла (`TOPO_GROUP_NODE`) с сортировкой. Проверяет overlap. Расширяет capacity при необходимости | -| `route_delete(table, node)` | Удаляет все записи, где `v_node_info == node`, со сдвигом массива | -| `route_lookup(table, dest_ip)` | LPM-поиск. Возвращает `ROUTE_ENTRY*` или NULL. Инкрементит `stats` (hits/misses) | -| `route_table_print(table)` | Печатает содержимое таблицы в лог | -| `parse_subnet(str, network, prefix_length)` | Парсит `"a.b.c.d/n"` → network (BE) + prefix_length. Возвращает 0/-1 | -| `is_local_subnet(ip)` | Проверяет, является ли IP локальным/мультикастовым/зарезервированным (0=true). RFC1918 + 0.0.0.0/8 + 127.0.0.0/8 + АPIPA + multicast + broadcast + CGNAT (100.64.0.0/10) | -| `route_add_local_subnet(table, network, prefix_length)` | **Объявлена в .h, но не реализована** | - -### Внутренние функции - -| Функция | Описание | -|---------|----------| -| `prefix_to_mask(prefix)` | Преобразует длину префикса в маску: `/24` → `0xFFFFFF00` | -| `routes_overlap(n1, p1, n2, p2)` | Проверяет, совпадают ли две подсети (одинаковые network+prefix) | -| `check_route_overlap_in_table(net, pre, entries, count)` | Проверяет overlap со всеми существующими записями | -| `binary_search_insert_pos(entries, count, net, pre)` | Находит позицию для вставки с сохранением сортировки | -| `binary_search_lpm(table, dest_ip)` | Бинарный LPM-поиск с fallback на соседнюю запись | diff --git a/src/routing_layer/route_p2pconn.txt b/src/routing_layer/route_p2pconn.txt deleted file mode 100644 index 98ccf871..00000000 --- a/src/routing_layer/route_p2pconn.txt +++ /dev/null @@ -1,11 +0,0 @@ -Установка прямого подключения между узлами - -1. узел собирает адреса - кандидаты для подключения: -список пар адрес-порт: - - адреса и порты из конфига сокетов (если не 0.0.0.0) - - адреса и порты из сокетов (nat) где nat = EIM -список ретрансляторов узлов-кандидатов: - - rtt, nodeid, адрес-порт - - -пары ip:port передаются на peer в команде "direct connection request" ROUTE_SUBCMD_CONN_REQ 0x0B diff --git a/src/routing_layer/topo_group.c b/src/routing_layer/topo_group.c index b9bf13a4..a6ee7696 100644 --- a/src/routing_layer/topo_group.c +++ b/src/routing_layer/topo_group.c @@ -655,6 +655,11 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from uint64_t node_id = ni->node_id; if (ni->hop_count >= MAX_HOPS) { DEBUG_WARN(DEBUG_CATEGORY_BGP, "NODEINFO from %s dropped: too many hops (%d)", from->log_name, ni->hop_count); return -1; } + if (node_id == 0) { + DEBUG_ERROR(DEBUG_CATEGORY_BGP, "NODEINFO node_id=0 from %s grp=%016llx — rejected (malformed)", + from->log_name, (unsigned long long)group->group_id); + return -1; + } if (node_id == group->instance->node_id) return 0; struct TOPO_GROUP_NODE* nodeinfo1 = topo_node_find_by_id(group, node_id); @@ -871,7 +876,18 @@ void topo_group_send_nodeinfo(struct TOPO_GROUP* group, struct TOPO_GROUP_NODE* p[0] = ETCP_ID_TOPO_ENTRY; p[1] = TOPO_SUBCMD_NODEINFO; struct TOPO_NODE* sni = topo_node_registry_find(group->instance->topo_groups, node->node_id); - if (!sni) { DEBUG_WARN(DEBUG_CATEGORY_BGP, "send_nodeinfo: node %016llx NOT in registry — skip forward to %s", (unsigned long long)node->node_id, conn->log_name); u_free(p); return; } + if (!sni) { + DEBUG_WARN(DEBUG_CATEGORY_BGP, + "send_nodeinfo: node %016llx NOT in registry — skip forward to %s grp=%016llx type=%d ch=%s is_local=%d paths=%d nodes_total=%d", + (unsigned long long)node->node_id, conn->log_name, + (unsigned long long)group->group_id, group->group_type, + group->channel_id[0] ? group->channel_id : "-", + (node == group->local_node) ? 1 : 0, + node->paths ? queue_entry_count(node->paths) : -1, + group->nodes ? queue_entry_count(group->nodes) : -1); + u_free(p); + return; + } DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "send_nodeinfo: node %016llx ver=%d grp=%016llx to conn=%s", (unsigned long long)node->node_id, sni->ver, (unsigned long long)sni->group_id, conn->log_name); int ser_len = topo_node_serialize(sni, node, p + 2, max_sz - 2, cumulative_rtt); if (ser_len < 0) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "send_nodeinfo: serialize failed for node %016llx", (unsigned long long)node->node_id); u_free(p); return; } diff --git a/src/routing_layer/topo_group_doc.md b/src/routing_layer/topo_group_doc.md deleted file mode 100644 index 73f9f462..00000000 --- a/src/routing_layer/topo_group_doc.md +++ /dev/null @@ -1,315 +0,0 @@ -# topo_group - -## 1. Назначение - -BGP-подобный обмен топологией узлов между пирами uTun через ETCP. Каждый узел хранит таблицу -известных узлов и при подключении нового пира синхронизирует её с ним. Изменения информации об -узле и обнаружение недоступности распространяются всем пирам (NODEINFO/WITHDRAW). - -Модуль поддерживает **изолированные группы** — разные «пространства имён» узлов с разной -семантикой: -- **TOPO_GROUP_TYPE_UTUN (1)** — VPN-сеть: обмен маршрутами (подсети) между узлами, NAT-детекция, - проверка связности. -- **TOPO_GROUP_TYPE_CHAT (2)** — чат-группа: только узлы, без подсетей, персистентность в SQLite. - -Группа по умолчанию (UTUN, `group_id = TOPO_GROUP_UTUN`) создаётся автоматически при -`topo_groups_init()`. Дополнительные группы создаются через `topo_groups_create_group()` (для -chatgui). Узлы разных групп изолированы, маршрутизация пакетов внутри группы ведётся по -`group_id` из wire-заголовков. - -Помимо обмена топологией, модуль выполняет: -- **Сборку local_node** (`topo_group_update_my_nodeinfo`) — при старте и при изменениях сокетов/active-mode. -- **NAT-детекцию** — выделена в отдельный модуль `nat_detection.h` (свой ETCP_ID 0x02). -- **Поиск оптимального маршрута** (`topo_group_find_conn_for_node`) — по min hop_count среди live-путей. -- **Зондирование связности** — `route_connectivity_probe_node()` для новых узлов. -- **Авто-подключение к узлам CHAT-группы** — `topo_group_connect` (бесконечный цикл до успеха). -- **Восстановление каскадно отвалившихся узлов** — `topo_recovery` (после разрыва соединения). -- **Invite/join** — `topo_group_invite` (прямое подключение + проверка членства). -- **Персистентность в SQLite** — `topo_node_sqlite` (таблицы nodes, node_addresses, peers_*). -- **Per-group broadcast** — `broadcast` (дёшевая рассылка данных с TTL внутри группы). - -## 2. Как пользоваться - -### 2.1. Инициализация (при старте utun) - -```c -struct TOPO_GROUPS* g = topo_groups_init(instance); // контейнер + utun-группа по умолчанию -``` -Здесь же: -- Открывается SQLite (`{db_path}/chats.db`), если задан `db_path`; БД — `instance->topo_sqlite_db`. -- `etcp_bind(instance, ETCP_ID_TOPO_ENTRY, topo_group_receive_cbk)` — приёмник пакетов топологии. -- `etcp_add_conn_status_cbk(instance, topo_group_conn_status, g)` — события UP/DOWN/DELETE всех соединений. -- `etcp_add_socket_cbk(...)` — изменения адресов/статуса сокетов. -- `utun_add_activity_cbk(...)` — смена active/standby режима. - -### 2.2. Новое ETCP-соединение - -При ETCP UP (`topo_group_conn_status` → `topo_group_new_conn` для каждой группы): -1. Дедупликация (тот же conn стреляет UP дважды — UDP- и TCP-линк). -2. `topo_recovery_cancel_for_node` — узел появился в сети, отменяем его recovery. -3. Обновление peer-узла: `conn_presence/conn_up |= NCONN_DIRECT`. -4. `topo_group_add_to_senders` + `topo_group_send_table_request` (запрос полной таблицы). -5. `topo_group_connect_on_up` (авто-подключение CHAT). - -На `REQUEST_TABLE` пир отвечает: NODEINFO о себе (local_node) + full table sync (все узлы, кроме -тех, чей путь уже включает этого пира — защита от петель) + `TABLE_COMPLETE`. При получении -`TABLE_COMPLETE` вызывается `etcp_set_routing_exchange_state(from_conn, 3)`. - -### 2.3. NODEINFO (обновление/создание узла) - -`topo_group_process_nodeinfo(group, from, data, len)`: -1. Проверка `group_id` пакета (`ni->group_id == group->group_id`). -2. Проверка соответствия типа группы по флагу `TOPO_FLAG_SEND_SUBNETS` (CHAT → 0, UTUN → 1); - при несовпадении — `TOPO_SUBCMD_ERR_GROUP_MISMATCH`. -3. Проверка `hop_count < MAX_HOPS`; пропуск собственного `node_id`. -4. Проверка версии: если `(int8_t)(last_ver - new_ver) >= 0` (stale) — только обновляется путь. -5. Десериализация wire-формата → `TOPO_NODE` (глобальная идентичность) + подсети + hop_list. -6. Верификация идентичности: `node_id == sc_derive_node_id_from_pubkey(public_key)` и - Ed25519 self-signature (`x25519_self_sig`) над каноническим сообщением. -7. Сохранение в глобальный реестр `node_registry` (дедупликация по node_id, ref_count). -8. Обновление/создание `TOPO_GROUP_NODE` (per-group): paths, last_ver, subnets. -9. `topo_group_remove_path_by_hop` + `topo_group_add_path` (через текущий conn). -10. `conn_presence/conn_up |= NCONN_BGP` (+ `NCONN_DIRECT` если прямой). -11. `route_insert` — только для UTUN-групп. -12. SQLite: `topo_node_sqlite_member_exists` → `addrs_put`; для CHAT — `member_placeholder_put` - (source=1, вне merkle-хеша); `nodeinfo_updated`. -13. `control_server_notify_node_change`. -14. Callback `node_updated_cb` (legacy, в текущем коде потребителей нет). -15. Fire BGP node event callbacks (`TOPO_NODE_EVENT_NEW`/`UPDATE`) — используют member_sync, media_delivery. -16. `topo_fire_nodeinfo_cbk` — глобальные nodeinfo-подписчики. -17. Broadcast NODEINFO всем пирам, кроме самого узла и тех, чей node_id есть в hop_list (антипетля). -18. `route_connectivity_probe_node` для новых/изменённых узлов с адресами. - -### 2.4. WITHDRAW (удаление узла) - -`topo_group_process_withdraw(group, sender, data, len)`: -1. Поиск `TOPO_GROUP_NODE` по `node_id`. -2. `topo_group_remove_path_by_hop(nq, wd_source)` — удаление путей, содержащих wd_source в hop_list. -3. Если путей не осталось: сброс `conn_presence/conn_up`, `route_delete` (UTUN), уведомление - control_server, отмена connectivity probe, освобождение paths и group-полей, удаление из - `nodes`, `etcp_router_conn_close_all_for_node`, broadcast withdraw (кроме отправителя), - fire `TOPO_NODE_EVENT_REMOVE`. -4. **BGP не удаляет member-запись из SQLite** — удаление из `peers_*` только через битую подпись - при верификации (member_sync). - -### 2.5. Отключение соединения - -При ETCP DOWN/DELETE (`topo_group_remove_conn`): -1. Проверка наличия conn в senders_list (если нет — выход). -2. Для каждого узла с путями через conn: вычисляется next_hop из hop_list, `remove_path`. Если - узел стал недоступен — `topo_recovery_add_node` (для каскадных, кроме самого пира), полная - очистка (как в WITHDRAW) + broadcast withdraw. -3. Удаление conn из `senders_list`. -4. Если были каскадные узлы — `topo_recovery_start`. -5. `topo_group_connect_on_down`. - -### 2.6. Поиск пути до узла - -`topo_group_find_conn_for_node(group, node_id)` — перебирает paths узла, выбирает min `hop_count`, -предпочитая live-пути (`conn->links_up > 0`); fallback на любой путь. Интенсивно используется -`etcp_router` и `conn_mgr` (next-hop для маршрутов и multihop-соединений). - -## 3. API - -### 3.1. Структуры - -| Структура | Назначение | -|---|---| -| `TOPO_GROUPS` | Контейнер всех групп. `group_list` (hash), `node_registry` (глобальный реестр TOPO_NODE*, hash), memory_pool'ы (v4/v6 sock_meta, addr, subnet), `node_updated_cb` (legacy) | -| `TOPO_GROUP` | Одна группа. `group_id`, `group_type`, `nodes` (queue TOPO_GROUP_NODE), `senders_list` (queue conn'ов), `local_node`, `ed25519_public_key`, `channel_id` (CHAT), `conn_mgr`, `recovery_list`, `connect`, `broadcast`, `node_cbks` | -| `TOPO_NODE` | Глобальная идентичность узла (в `node_registry`): pubkeys, `x25519_self_sig`, name, v4/v6 sock_meta+addrs, `client_type`, `client_activity`, `ver`, `group_ref_count` | -| `TOPO_GROUP_NODE` | Per-group запись узла: `node_id`, `subnets`, `paths`, `last_ver`, `conn_mgr_*`, `connectivity`, `conn_presence`, `conn_up`, `handle` (ncd) | -| `TOPO_NODEPATH` | Путь до узла через ETCP_CONN: `{conn, hop_count, cumulative_rtt}` + inline `hop_list[hop_count]` | -| `TOPO_NODESUBNETS` | Подсети узла: списки `TOPO_SUBNET4`/`TOPO_SUBNET6` | -| `TOPO_CONNECTIVITY` | Связность: probe_status, rtt-метрики, статусы (interface/nat/real), `probe_deferred_wait` (standby) | -| `TOPOMSG_NODEINFO_PKT` | Wire NODEINFO: cmd + subcmd + TOPOMSG_NODE + переменная часть | -| `TOPOMSG_WITHDRAW_PKT` | Wire WITHDRAW: `group_id`, `node_id` (удаляемый), `wd_source` (инициатор) | -| `TOPOMSG_TABLE_REQ` | Wire запроса/завершения таблицы: `group_id` | -| `TOPOMSG_ERR_GROUP_MISMATCH` | Wire ошибки типа группы: `group_id`, `expected_type`, `received_flags` | -| `NAT_DETECTION` | Контекст NAT-детекции (отдельный модуль) | - -### 3.2. Константы - -| Константа | Значение | Назначение | -|---|---|---| -| `ETCP_ID_TOPO_ENTRY` | 0x01 | ETCP ID пакетов топологии | -| `TOPO_SUBCMD_NODEINFO` | 0x04 | Полная информация об узле | -| `TOPO_SUBCMD_REQUEST_TABLE` | 0x05 | Запрос полной таблицы | -| `TOPO_SUBCMD_WITHDRAW` | 0x06 | Узел стал недоступен | -| `TOPO_SUBCMD_TABLE_COMPLETE` | 0x0B | Завершение начальной синхронизации | -| `TOPO_SUBCMD_ERR_GROUP_MISMATCH` | 0x0C | Ошибка несоответствия типа группы | -| `MAX_HOPS` | 16 | Максимальная длина hop-листа | -| `BGP_NODES_HASH_SIZE` / `TOPO_NODE_REGISTRY_HASH_SIZE` | 256 | Хеш-таблицы nodes / node_registry | -| `TOPO_GROUP_UTUN` | 0x8000000000000000ULL | group_id utun-группы по умолчанию | -| `TOPO_GROUP_TYPE_UTUN` / `TOPO_GROUP_TYPE_CHAT` | 1 / 2 | Тип группы | -| `TOPO_FLAG_SEND_SUBNETS` | 0x01 | Флаг NODEINFO: узел передаёт подсети | -| `NCONN_DIRECT` / `NCONN_INDIRECT` / `NCONN_BGP` | 1/2/4 | Биты наличия соединений (presence/up) | -| `CONN_TYPE_DIRECT` / `REVERSE` / `INDIRECT` / `NONE` | 1/2/3/0 | Стратегия conn_mgr | -| `CLIENT_TYPE_SERVER` / `DESKTOP` / `MOBILE` | 0/1/2 | Тип клиента | -| `CLIENT_ACTIVITY_STANDBY` / `ACTIVE` | 0/1 | Режим (standby/active) | -| `PROBE_STATUS_NONE/IN_PROGRESS/DONE/DEFERRED` | 0/1/2/3 | Состояние пробы | -| `PROBE_RESULT_UNKNOWN/REACHABLE/UNREACHABLE` | 0/1/2 | Результат пробы | -| `TOPO_NODE_EVENT_NEW/UPDATE/REMOVE` | 0/1/2 | События BGP node callbacks | -| `CONN_MGR_MAX_INTERMEDIARIES` | 3 | Максимум посредников для INDIRECT | - -### 3.3. Инициализация/завершение - -| Функция | Назначение | -|---|---| -| `topo_groups_init(instance)` | Контейнер + node_registry + пулы + SQLite + utun-группа + ETCP/activity коллбэки | -| `topo_groups_destroy(instance)` | Отвязка коллбэков, destroy групп, node_registry, пулов, закрытие SQLite | -| `topo_group_create(instance, group_id, group_type)` | Создать TOPO_GROUP: nodes/senders/local_node/conn_mgr/broadcast | -| `topo_group_destroy(group)` | recovery_cancel_all, connect_destroy, broadcast_destroy, conn_mgr_destroy, очистка nodes/local_node | -| `topo_groups_create_group(g, group_id, group_type, channel_id)` | Создать группу + добавить в group_list + new_conn для существующих соединений + connect_init (CHAT) | -| `topo_groups_find(g, group_id)` / `topo_groups_get_default(g)` | Поиск группы / utun-группа | -| `topo_group_on_activity_change(instance, active, arg)` | Обновить local_node и broadcast его по всем группам/пирам | - -### 3.4. Жизненный цикл соединения - -| Функция | Назначение | -|---|---| -| `topo_group_new_conn(group, conn)` | Дедуп, отмена recovery, обновить peer, add_to_senders, request_table, connect_on_up | -| `topo_group_remove_conn(group, conn)` | Удалить пути, удалить недоступные узлы (recovery/broadcast), убрать conn, connect_on_down | -| `topo_group_conn_status(conn, status, arg)` | Callback status: UP→new_conn, DOWN/DELETE→remove_conn для каждой группы | - -### 3.5. Обработка пакетов - -| Функция | Назначение | -|---|---| -| `topo_group_receive_cbk(from_conn, entry)` | Точка входа: извлекает `group_id`, находит группу, диспетчеризация по subcmd | -| `topo_group_process_nodeinfo(group, from, data, len)` | Полный цикл NODEINFO (см. 2.3) | -| `topo_group_process_withdraw(group, sender, data, len)` | Удаление путей/узла по WITHDRAW (см. 2.4) | -| `topo_group_handle_request_table(group, conn)` | local_node + full table + add_to_senders + TABLE_COMPLETE | - -### 3.6. Поиск и пути - -| Функция | Назначение | -|---|---| -| `topo_group_find_conn_for_node(group, node_id)` | Оптимальный ETCP_CONN (min hop_count, предпочтение live) | -| `topo_node_find_by_id(group, node_id)` | Поиск TOPO_GROUP_NODE по node_id (hash) | -| `topo_group_add_path(nq, conn, hop_list, hop_count, cumulative_rtt)` | Создать TOPO_NODEPATH + inline hop_list | -| `topo_group_remove_path(nq, conn)` | Удалить путь; возвращает 1 если путей не осталось | -| `topo_group_remove_path_by_hop(nq, wd_source)` | Удалить пути, содержащие wd_source в hop_list | -| `topo_group_should_send_to(nq, target_id)` | Проверка антипетли (target_id не во всех путях) | - -### 3.7. Broadcast - -| Функция | Назначение | -|---|---| -| `topo_group_send_nodeinfo(group, node, conn, cumulative_rtt)` | Сериализовать и отправить NODEINFO одному conn | -| `topo_group_send_withdraw(group, node_id)` | Broadcast WITHDRAW (wd_source = self) | -| `topo_group_broadcast_withdraw(group, node_id, wd_source, exclude)` | Broadcast WITHDRAW всем, кроме exclude | -| `topo_group_send_full_table(group, conn)` | Все известные узлы (антипетля) указанному conn | -| `topo_group_send_table_request/complete(group, conn)` | REQUEST_TABLE / TABLE_COMPLETE | -| `topo_group_send_err_group_mismatch(group, conn, ...)` | ERR_GROUP_MISMATCH | - -### 3.8. Коллбэки (подписки) - -| Функция | Назначение | -|---|---| -| `topo_group_add_node_cbk/remove_node_cbk(group, fn, arg)` | BGP node events (NEW/UPDATE/REMOVE) — consumer: member_sync, media_delivery | -| `utun_add_nodeinfo_cbk/remove_nodeinfo_cbk(instance, fn, arg)` | Глобальная подписка на изменения nodeinfo (per-instance) | -| `topo_fire_nodeinfo_cbk(instance, group, node)` | Вызвать глобальные nodeinfo-подписчики | -| `topo_groups_set_node_updated_cb(groups, fn)` | Legacy callback (в текущем коде потребителей нет) | - -### 3.9. NAT (делегаты в nat_detection) - -| Функция | Назначение | -|---|---| -| `topo_group_set_nat_check_local(group, allow)` | Делегат в `nat_detection_set_allow_local` | - -### 3.10. Связанные модули - -| Модуль | Назначение | -|---|---| -| `topo_group_connect` | Авто-подключение к узлам CHAT-группы (3 фазы, бесконечный цикл, standby-aware) | -| `topo_group_invite` | Invite/join канала через прямое ncd-подключение (INVITE_INFO_REQ/RESP, ETCP_ID 0x13) | -| `topo_recovery` | Восстановление каскадно отвалившихся узлов после разрыва соединения | -| `broadcast` | Per-group рассылка данных с TTL (ETCP_ID_BROADCAST), `broadcast_init/сброс` | -| `topo_node` | Модель данных: TOPO_NODE/TOPO_GROUP_NODE, ser/deserialize, node_registry, sig-msg | -| `topo_node_sqlite` | Персистентность: nodes, node_addresses, peers_*, member block/owner/placeholder | - -## 4. Архитектурные связи - -### 4.1. Внешние зависимости - -| Модуль | Как используется | -|---|---| -| `topo_node.c/h` | Модель узла, ser/deserialize, node_registry, update_my_nodeinfo, sign_self | -| `topo_node_sqlite.h` | Персистентность (node_put, member_placeholder_put, addrs_put, nodeinfo_updated) | -| `nat_detection.h` | NAT-детекция (отдельный модуль, ETCP_ID 0x02) | -| `route_lib.h` | route_insert / route_delete | -| `route_connectivity.h` | probe_node / cancel_node / cancel_all | -| `etcp_api.h` | etcp_bind, etcp_add_conn_status_cbk, etcp_add_socket_cbk, etcp_unbind | -| `etcp_connections.h` | ETCP_CONN, conn->peer_node_id, links_up | -| `etcp.h` | etcp_send, etcp_set_routing_exchange_state | -| `etcp_router.h` | etcp_router_conn_close_all_for_node | -| `control_server.h` | notify_node_change / notify_node_removed | -| `conn_mgr.h` | conn_mgr_init/destroy (per-group), поиск путей | -| `secure_channel.h` | sc_derive_ed25519_pubkey, sc_derive_node_id_from_pubkey, sc_ed25519_verify | -| `broadcast.h` | broadcast_init/destroy (per-group) | -| `memory_pool.h` | Пулы для сериализации | -| `debug_config.h` | Логирование, категория BGP | -| `mem.h` | u_malloc/u_free/u_calloc | - -### 4.2. Кто вызывает topo_group - -| Вызывающий | Что вызывает | Для чего | -|---|---|---| -| `utun_instance.c` | `topo_groups_init/destroy` | Инициализация/завершение | -| `etcp_connections.c` | `topo_group_send_nodeinfo` | Переотправка local_node при смене сокетов | -| `etcp_router.c` | `topo_group_find_conn_for_node` | next-hop для маршрутов | -| `conn_mgr_*.c` | `topo_group_find_conn_for_node`, `topo_fire_nodeinfo_cbk` | Пути для соединений | -| `nat_detection.c` | `topo_group_send_nodeinfo`, `topo_group_update_my_nodeinfo` | Обновление local_node после NAT-чека | -| `member_sync.c` | `topo_group_add_node_cbk` | BGP → member_sync → node_props_changed | -| `media_delivery.c` | `topo_group_add_node_cbk` | События узлов для медиа | -| `route_connectivity.c` | `topo_fire_nodeinfo_cbk` | Уведомление об изменении узла | - -## 5. Форматы пакетов (wire) - -### 5.1. NODEINFO -``` -cmd(1) | subcmd(1) | TOPOMSG_NODE | [node_name] | [v4_sock_meta[]] | [v4_addrs[]] | [v6_sock_meta[]] | [v6_addrs[]] | [v4_subnets[]] | [v6_subnets[]] | [hop_list[]] -``` -`TOPOMSG_NODE` содержит `group_id`, `node_id`, `ver`, pubkeys, `x25519_self_sig(64)`, счётчики, -`hop_count`, `client_type`, `client_activity`, `cumulative_rtt`. - -### 5.2. WITHDRAW -``` -cmd(1) | subcmd(1) | group_id(8) | node_id(8) | wd_source(8) = 26 байт -``` - -### 5.3. REQUEST_TABLE / TABLE_COMPLETE -``` -cmd(1) | subcmd(1) | group_id(8) = 10 байт -``` - -### 5.4. ERR_GROUP_MISMATCH -``` -cmd(1) | subcmd(1) | group_id(8) | expected_type(1) | received_flags(1) = 12 байт -``` - -## 6. Примечания - -- **Глобальный реестр идентичности**: `TOPO_NODE` (pubkeys, адреса, имя) хранится один раз в - `node_registry` и разделяется между группами через `group_ref_count`. Per-group данные — в - `TOPO_GROUP_NODE` (paths, subnets, connectivity, conn_presence/up). -- **Версионирование NODEINFO**: `ver` (uint8_t) сравнивается циклически `(int8_t)(cur - new) >= 0`; - stale-версия только обновляет путь, не заменяя данные узла. -- **Верификация идентичности**: `node_id == derive(x25519)` + Ed25519 self-signature над каноническим - сообщением (`topo_node_build_sig_msg`). Несовпадение → пакет отвергается как forgery. -- **Защита от петель**: NODEINFO рассылается пирам, чей node_id отсутствует в hop_list всех путей; - WITHDRAW broadcast всем, кроме отправителя. -- **Hop-лист хранится в пути**: `TOPO_NODEPATH` содержит inline `hop_list[hop_count]` сразу после - структуры; `next_hop` для recovery вычисляется из него. -- **Связность**: для новых узлов с адресами запускается `route_connectivity_probe_node`; при - удалении — `route_connectivity_cancel_node`. В standby пробы откладываются (`PROBE_STATUS_DEFERRED`). -- **Виртуальные каналы**: при удалении узла закрываются `etcp_router_conn_close_all_for_node`. -- **CHAT-группа пишет плейсхолдер** `source=1` (pubkeys из BGP, без подписей) — вне merkle-хеша; - полная запись приходит позже через merkle-путь (member_sync). -- **BGP не удаляет member из SQLite**: работает только со своей RAM-таблицей; удаление из `peers_*` - — только битая подпись (member_sync). -- **NAT check для локальных сетей**: по умолчанию запрещён (`allow_nat_check_local=0`), включается - для тестов через `topo_group_set_nat_check_local`. -- **Обновление local_node при NAT/active-изменениях**: `nat_detection` и `topo_group_on_activity_change` - вызывают `topo_group_update_my_nodeinfo` → новая версия → NODEINFO broadcast всем пирам. diff --git a/src/routing_layer/topo_strategy.txt b/src/routing_layer/topo_strategy.txt deleted file mode 100644 index 941ce1e7..00000000 --- a/src/routing_layer/topo_strategy.txt +++ /dev/null @@ -1,89 +0,0 @@ -++ OOP: -class Animal { -public: // доступен всем -//private: // только внутри класса -//protected: // внутри класса и наследников - virtual void sound() {play(...);}// virtual — разрешает переопределение. - - virtual void sound2() = 0;// если =0 - метод обязательно должен быть реализован в наследнике, а сам класс становится абстрактным и его экземпляры создавать нельзя. - - Animal() {// конструктор - совпадает с именем класса - cout << "Created"; - } - ~Animal() {// деструктор - cout << "Deleted"; - } - -}; - -class Dog : public Animal {// наследование -public: - void sound() override {}// override — говорит компилятору, что метод действительно переопределяет базовый. -}; - - - - -Стратегия обмена - -узлы которые online живут только в пемяти. и синхронизируемся с ними используя topo_group. - -routing layer: - - topo_group - обеспечивает корректное построение оптимальной таблицы роутинга (per group) между узлами в рамках установленных подключений. И корректировку этой таблицы при перестроении подключений. - - - topo_group_discovery - обеспечивает оптимизацию подключений, а именно: - - периодически пингует узлы проверяя их доступность - - -для каждой группы админ назначает суперузлы: a b c ... стать усуперузлом - получить разрешение с админской подписью. - -суперузел активирует bgp обмен с остаьлными суперузлами. -как: - -суперузел - это по сути сервер для группы. который также синхронизируется с другими суперузлами - -клиенты приоритетно подключаются к суперузлам: - - подключаются к кому-то - - запрашивают таблицу маршрутизации - - - - - -1. Узлы за nat сами выбирают подходящие для себя точки подключения к сети. - 2 шт с минимальным rtt - 2 шт с минимальной загруженностью - - в чате (gui) можно выбрать предпочитаемые узлы для подключения - - -2. узлы с eim nat - ??? если не хватает узлов прямым адресом - -3. узлы с прямыми ип - суть - выстраивают линки между собой для надежного обмена конфигурацией сети - как: - - в таблице узлов для каждой группы - -У каждого узла проставляется доступность (можем ли мы каким-либо образом подключиться): acessible, reverse (при подключении с другой стороны), no - - -Единое обновление инфы об узле (per-group): - - обновляем тип узла в каждой группе где есть узел при обовлении типа нат узла. прямая доступность узла: yes/no. - - - - -======================================= -функционал группы. - -Есть суперадмины группы - владеют приватным ключом суперадмина (один на группу). -Есть админы - владеют сгенерированным приватным ключом админа и подписанным суперадмином. -суперадмины могут делать revoke админу. - -админские действия также распространяются по чату как и обычные сообщения. -ид группы = пубкей группы - - -юзер - join по ссылке (приватный ключ пира). получает ключи других пиров. Может только читать. -админ подписывает разрешение на писать (как и другие разрешения). подпись храним в таблице юзеров diff --git a/src/transport_layer/etcp_doc.md b/src/transport_layer/etcp_doc.md deleted file mode 100644 index 920e77d6..00000000 --- a/src/transport_layer/etcp_doc.md +++ /dev/null @@ -1,289 +0,0 @@ -# etcp - -## 1. Назначение - -Ядро протокола ETCP — TCP-подобный надёжный транспорт поверх UDP с шифрованием (AES-CCM + X25519), multi-link (несколько каналов между двумя узлами), фрагментацией/сборкой пакетов, ретрансмиссией с экспоненциальным backoff, BBR congestion control и burst-измерением пропускной способности. - -Модуль управляет **одним ETCP-соединением** (struct `ETCP_CONN`) между локальным и удалённым узлом. Соединение может иметь несколько линков (struct `ETCP_LINK`), работающих через разные сокеты / сетевые пути. - -Этапы жизни соединения: -``` -создание (pending, state=0) → первый линк проинициализирован → ready (state=1) → close (state=2, detach) → ref_count=0 → полное освобождение ресурсов -``` - -## 2. Как пользоваться - -### 2.1. Типовой сценарий - -1. **Создание:** `etcp_connection_create(instance, name)` — создаёт `ETCP_CONN` в состоянии `pending` (state=0), размещает в очереди `instance->connections` с индексом peer_node_id=0. -2. **Добавление линков:** `etcp_link_new(etcp, socket, remote_addr, is_server)` — создаёт `ETCP_LINK`, добавляет в список `etcp->links`. -3. **Обмен ключами / INIT:** через `etcp_connections.c` происходит обмен INIT_REQUEST/INIT_RESPONSE, установка `peer_node_id`, инициализация `secure_channel` (X25519 + AES-CCM). -4. **Готовность:** `etcp_conn_ready()` → `etcp_conn_queue_set_ready()` — переиндексирует соединение в очереди `instance->connections` с реальным `peer_node_id`, переводит `state=1`, запускает callback'и `ready_cbks`, таймер метрик. -5. **Отправка данных:** `etcp_int_send(etcp, data, len)` — аллоцирует память из `data_pool`, создаёт `ETCP_FRAGMENT`, помещает в `input_queue`. -6. **Приём данных:** потребитель читает из `etcp->output_queue` — там лежат собранные `ETCP_FRAGMENT` с полезной нагрузкой. -7. **Закрытие:** `etcp_connection_close(etcp)` — 2-фазное: Phase 1 — detach от внешнего мира (линки, таймеры, роутинг), Phase 2 — отложенное освобождение ресурсов (очереди, пулы памяти, callback'и, struct) когда ref_count=0. - -### 2.2. Ключевые концепции - -#### Очереди данных (data flow) - -``` - ┌──────────────┐ - et cp_int_send ──► │ input_queue │ входная очередь (ETCP_FRAGMENT из io_pool) - └──────┬───────┘ - │ input_queue_cb: перенос во inflight с присвоением seq - ┌──────▼───────┐ - │ input_send_q │ очередь на отправку (INFLIGHT_PACKET из inflight_pool) - └──────┬───────┘ - │ etcp_request_pkt: формирует ETCP_DGRAM с секциями ACK+PAYLOAD - ┌──────▼───────┐ - │input_wait_ack│ ожидание ACK (INFLIGHT_PACKET, index по seq) - └──────┬───────┘ - │ etcp_ack_recv: ACK получен → освобождение inflight - [удалён] -``` - -``` -Приём: - ┌──────────┐ - etcp_conn_input ──► │ recv_q │ очередь сборки (ETCP_FRAGMENT, index по seq) - └────┬─────┘ - │ etcp_output_try_assembly: поиск непрерывной последовательности - ┌────▼──────┐ - │output_queue│ собранные пакеты для потребителя - └───────────┘ -``` - -- **`input_queue`** (ETCP_FRAGMENT) — пакеты от приложения к отправке. Callback `input_queue_cb` переносит их в `input_send_q`, создавая `INFLIGHT_PACKET` с seq. -- **`input_send_q`** (INFLIGHT_PACKET, hash-индекс по seq) — ожидающие отправки (новые или ретрансмиссии). Callback `input_send_q_cb` вызывает `etcp_conn_process_send_queue`. -- **`input_wait_ack`** (INFLIGHT_PACKET, hash-индекс по seq) — отправленные пакеты, ожидающие ACK. Callback `wait_ack_cb` запускает retrans-таймер. -- **`recv_q`** (ETCP_FRAGMENT, hash-индекс по seq) — принятые фрагменты, ждущие сборки. -- **`output_queue`** — собранные непрерывные пакеты для потребителя. -- **`ack_q`** (ACK_PACKET, hash-индекс по seq) — неотправленные ACK подтверждения, ожидающие piggyback. -- **`transit_queues`** — транзитные очереди для роутинга (хеш по `src_node_id:dst_node_id`). - -#### Ретрансмиссия с экспоненциальным backoff - -- Таймаут ретрансмиссии: `timeout = (rtt_avg_10 * K1 + jitter * K2) / 16`, где `K1=32`, `K2=32`. -- Минимальный timeout: 50 tb (5ms), максимальный: 10000 tb (1000ms). -- Пакеты в `input_wait_ack` упорядочены по времени отправки (FIFO). Проверка (`ack_timeout_check`) идёт с головы: как только найден пакет с неистекшим таймаутом — взводится таймер до его истечения, дальше не сканируем. -- Ретрансмит: перенос из `input_wait_ack` в `input_send_q`, инкремент `send_count`, добавление в `send_hist`. -- При ретрансмите старый линк (`last_link`) получает `total_retransmissions++` и вычитание из `inflight_bytes/packets`. - -#### RTT/Jitter/Bandwidth метрики - -- **RTT:** измеряется через секцию `ETCP_SECTION_TIMESTAMP` — результирующий `rtt_last` на соединение берётся как среднее по всем линкам (rtt_sum/cnt), `rtt_avg_10` — максимум RTT по линкам. -- **Jitter:** per-link экспоненциальное сглаживание разности RTT: `jitter += (|new_rtt - prev_rtt| * 65536 - jitter) / 32`. На соединении — среднее по линкам. -- **TT (transmission time):** время доставки отправленных пакетов, вычисляется из timestamp-секции. -- **RT (recv time):** время доставки принятых пакетов. -- **Bandwidth:** per-link, обновляется из BBR pacing rate и burst-измерений. - -#### Burst bandwidth measurement - -- Отправитель формирует пачку из `BURST_PACKET_COUNT` (12) пакетов, каждый содержит секцию `MEAS_TS` (type=0x07) с burst_id, flags, порядковым номером и µs-таймстемпом. -- Флаги: `IS_FIRST`, `IS_LAST`, `IS_FILLER` (если нет данных — пакет заполняется FILLER-секцией). -- Пропускаются первые `BURST_SKIP_COUNT` (3) пакета для исключения переходных процессов при замере inter-packet gap. -- Приёмник собирает времена прихода, вычисляет min/avg gap между пакетами, отправляет ответ `MEAS_RESP` (type=0x08) с gap_avg, gap_min, pkt_count. -- Bandwidth вычисляется как: `BW = pkt_size * 8000 / gap_min` (Kbps), BDP = `BW * 1000 / 8 * min_rtt_sec`. -- Минимальный интервал между burst: 500ms. Таймаут ожидания ответа: 2s. -- Ответ может быть piggyback'нут в обычный пакет (вне активного burst). - -#### Channel timestamp handling - -- Секция `ETCP_SECTION_TIMESTAMP` (type=0x06, 5 байт): добавляется когда есть новые данные о времени приёма от пира. -- Содержит: ret_ts (текущее время пира на момент получения нашего пакета) и dt (разница local_time - timestamp пакета). -- Позволяет вычислять RTT, jitter, tt, rt на всех линках, а также `recv_dt_avg_rx/tx` — экспоненциально сглаженные оценки односторонних задержек. - -#### 2-phase deferred cleanup - -- **Phase 1** (синхронно в `etcp_connection_close`): detach от внешнего мира — закрытие линков, отмена таймеров, остановка метрик, удаление из роутинга, удаление из очереди `instance->connections`, установка `state=2`. -- **Phase 2** (отложено через `uasync_call_soon`): освобождение очередей, memory_pool'ов, callback-цепочек, struct. Выполняется когда `ref_count == 0`. Если есть внешние ссылки — ждёт `etcp_conn_ref_free()`. -- `ref_count` нужен для безопасного доступа к conn из асинхронных контекстов (коллбэки, роутинг, BGP). - -#### Backpressure (пороговое ожидание) - -- `input_queue` имеет порог 0 (ждёт полного освобождения) — новые пакеты добавляются только когда очередь пуста. -- `etcp_conn_process_send_queue` пробует протолкнуть данные из `input_queue` → `input_send_q` когда `wait_ack_bytes <= optimal_inflight`. -- При ACK вызывается `input_queue_try_resume` — если `wait_ack_bytes <= optimal_inflight` и input_send_q пуста, возобновляет `input_queue_cb`. - -#### Normalizer (фрагментация/сборка) - -- `struct PKTNORM` — per-connection фрагментатор/дефрагментатор. Разбивает большие пакеты на части ≤ `frag_size` (MTU - overhead). -- При reinit/reset несобранные фрагменты возвращаются в normalizer для повторной обработки (`etcp_return_inflight_to_normalizer`). - -### 2.3. State machine - -| Поле | Значения | Описание | -|------|----------|----------| -| `state` | 0=pending, 1=ready, 2=deleted | Внешнее состояние conn | -| `initialized` | 0/1 | Хотя бы один линк прошёл обмен ключами | -| `reset_done` | 0/1 | 0=reinit разрешён, 1=соединение стабильно | -| `tx_state` | DATA_WAIT(1)/LINK_WAIT(2) | DATA_WAIT=можно отправлять, LINK_WAIT=все линки busy | -| `links_up` | 0/1 | Хотя бы один линк в статусе up | -| `got_initial_pkt` | 0/1 | Получен ли первый пакет (seq=1) после reset | -| `routing_exchange_active` | 0-4 | Статус BGP-обмена | - -## 3. API - -### 3.1. Основные структуры - -#### `struct ETCP_CONN` (`etcp.h:140`) - -Главная управляющая структура соединения. Содержит: -- **Состояние:** `state`, `ref_count`, `initialized`, `reset_done`, `links_up`, `tx_state`. -- **Очереди:** `input_queue`, `input_send_q`, `input_wait_ack`, `recv_q`, `output_queue`, `ack_q`, `transit_queues`. -- **Пулы памяти:** `inflight_pool` (для INFLIGHT_PACKET), `io_pool` (для ETCP_FRAGMENT). -- **Линки:** связный список `links`, `last_rr_link` (round-robin для load balancer). -- **Крипто:** `crypto_ctx` (secure_channel), `peer_ed25519_pubkey`. -- **Нормалайзер:** `normalizer` (PKTNORM) — фрагментация/сборка. -- **Метрики:** `rtt_last/avg_10/avg_100`, `jitter`, `tt_last`, окна (unacked_bytes, max_inflight, optimal_inflight). -- **Таймеры:** `retrans_timer`, `ack_resp_timer`. -- **Callback-цепочки:** `ready_cbks`, `up_cbks`, `down_cbks`, `bgp_ready_cbk`. -- **Счётчики:** ACK hit/miss, дубликаты, ретрансмиссии, `debug[8]` для live watch. -- **Метрики качества:** `metrics` (гистограммы RTT/потерь с 10-мин снимками, тотальные счётчики). -- **Идентификация:** `log_name` (формат `XXXX→YYYY [name]`), `name` (из конфига), `session_id`. - -#### `struct INFLIGHT_PACKET` (`etcp.h:108`) - -Пакет в состоянии inflight (отправлен или ждёт отправки). Поля: -- `seq` — sequence number (ID по протоколу). -- `state` — `WAIT_SEND` (в input_send_q) или `WAIT_ACK` (в input_wait_ack). -- `last_link` — последний линк, через который отправлен (для retrans/loss detection). -- `last_timestamp` — время последней отправки (для retrans timeout). -- `send_count` — количество попыток отправки. -- `retrans_req_count` — количество запросов ретрансмиссии. -- `send_hist[8]` — номера линков через которые передавался (кольцевой буфер, send_count — head). -- `delivered_at_send`, `inflight_at_send`, `is_app_limited` — снэпшоты для BBR rate estimation. - -#### `struct ETCP_FRAGMENT` (`etcp.h:123`) - -Фрагмент данных в очередях ввода/вывода. Поля: -- `seq` — sequence number. -- `timestamp` — timestamp пакета (от удалённой стороны при приёме). -- `ll` — `ll_entry` с dgram (указатель на data_pool) и len. - -#### `struct ACK_PACKET` (`etcp.h:129`) - -Неотправленное подтверждение приёма. Поля: -- `seq` — sequence number подтверждаемого пакета. -- `pkt_timestamp` — timestamp подтверждаемого пакета (часы удалённой стороны). -- `recv_timestamp` — локальное время приёма (для вычисления задержки ACK). - -#### `struct etcp_metrics` (`etcp.h:70`) - -Гистограммы и счётчики качества соединения: -- **Рабочие гистограммы:** `rtt_hist[12]`, `loss_hist[11]` — обновляются в реальном времени. -- **Финальные снимки:** `rtt_hist_final[12]`, `loss_hist_final[11]` — копируются раз в 10 мин. -- **Счётчики окна:** `work_samples/sent/rcvd/bytes_sent/bytes_rcvd/lost` — сбрасываются при снимке. -- **Тотальные:** `total_sent/rcvd/bytes_sent/bytes_rcvd/lost` — кумулятивные счётчики. -- **RTT-бакеты:** 0-1, 1-2, 2-3, 3-5, 5-8, 8-13, 13-21, 21-34, 34-55, 55-90, 90-150, 150+ ms (в 0.1ms). -- **Loss-бакеты:** 0%, 1%, ..., 10%+. -- Интервал снимка: 10 мин. Минимум пакетов для обновления: 200. - -### 3.2. Основные функции - -#### Жизненный цикл - -| Функция | Описание | -|---------|----------| -| `etcp_connection_create(instance, name)` | Создаёт ETCP_CONN, очереди, пулы, нормалайзер, помещает в `instance->connections` с key=0 (pending). | -| `etcp_connection_close(etcp)` | Phase 1: detach (линки, таймеры, роутинг). Phase 2: deferred free когда ref_count=0. | -| `etcp_conn_ref_take(conn)` | Увеличивает ref_count. Возвращает -1 если conn удалён. | -| `etcp_conn_ref_free(conn)` | Уменьшает ref_count. Если 0 и state==2 — планирует Phase 2. | -| `etcp_conn_reset(etcp)` | Сброс состояния после reinit: seq, метрики, очистка очередей, возврат данных в нормалайзер. | -| `etcp_conn_reinit(etcp)` | Полный reinit: сброс флагов, сброс NAT, сброс роутера, вызов `etcp_conn_reset`. | -| `etcp_links_reset(etcp)` | Сброс флага initialized у всех линков (после сбоя). | -| `etcp_conn_ready(conn)` | Вызывается когда первый линк проинициализирован. Устанавливает `initialized=1`, `reset_done=1`, запускает `conn_queue_set_ready`. | -| `etcp_conn_queue_set_ready(conn)` | Переиндексирует соединение в `instance->connections` с реальным `peer_node_id`, state=1, запускает ready/up callback'и, стартует таймер метрик. | -| `etcp_conn_set_peer_node_id(conn, id)` | Устанавливает peer_node_id. Если state=1 — переиндексирует в очереди. | -| `etcp_conn_on_inflight_lim_changed(etcp)` | Пересчитывает `optimal_inflight` как сумму `inflight_lim_bytes` всех линков. | -| `etcp_update_mtu(etcp)` | Пересчитывает MTU соединения как минимум MTU всех линков. Обновляет `frag_size` нормалайзера. | - -#### Отправка - -| Функция | Описание | -|---------|----------| -| `etcp_int_send(etcp, data, len)` | Отправка данных через ETCP: аллокация из data_pool, создание ETCP_FRAGMENT, помещение в input_queue. | -| `etcp_request_pkt(etcp)` | Формирует ETCP_DGRAM для отправки: выбирает линк через load balancer, берёт INFLIGHT_PACKET из input_send_q, добавляет ACK-секцию, опциональные секции (MEAS_TS, MEAS_RESP, TIMESTAMP, FILLER), payload. Перемещает inflight в wait_ack. | -| `etcp_encrypt_send(dgram)` | (etcp_connections.c) Шифрует ETCP_DGRAM (AES-CCM), добавляет nonce/tag/crc32/pubkey, отправляет через udp_send. | - -#### Приём - -| Функция | Описание | -|---------|----------| -| `etcp_conn_input(pkt)` | Обработка расшифрованного пакета: разбор секций (ACK, TIMESTAMP, PAYLOAD, MEAS_TS, MEAS_RESP, FILLER, METRICS), вызов `etcp_ack_recv`, добавление данных в `recv_q`/`ack_q`. | -| `etcp_ack_recv(etcp, seq, ts, dts)` | Обработка ACK: поиск пакета в `input_wait_ack`/`input_send_q`, вычитание inflight с линка, BBR rate estimation, освобождение памяти. | -| `etcp_output_try_assembly(etcp)` | Сборка выходной очереди: ищет непрерывную последовательность от `last_delivered_id+1` в `recv_q`, переносит в `output_queue`. | -| `etcp_stats(etcp)` | Вывод статистики в debug-лог: размеры очередей, RTT, счётчики, ID. | - -#### Callback'и / События - -| Функция | Описание | -|---------|----------| -| `etcp_on_link_down(etcp)` | Вызывается при падении линка. Проверяет все линки, если все down — вызывает `etcp_on_down`. | -| `etcp_on_up(etcp)` | Запуск цепочки `up_cbks`. | -| `etcp_on_down(etcp)` | Запуск цепочки `down_cbks`. | - -#### Метрики качества - -| Функция | Описание | -|---------|----------| -| `etcp_metrics_init(m)` | Обнуление структуры метрик. | -| `etcp_metrics_add_rtt(etcp, rtt_tb)` | Добавление RTT-замера в гистограмму. | -| `etcp_metrics_add_sent(etcp, len)` | Учёт отправленного пакета (working + total). | -| `etcp_metrics_add_rcvd(etcp, len)` | Учёт принятого пакета. | -| `etcp_metrics_add_loss(etcp, count)` | Учёт потерь. | -| `etcp_metrics_start_timer(etcp)` | Запуск 10-минутного таймера снимков. | -| `etcp_metrics_stop_timer(etcp)` | Остановка таймера метрик. | - -### 3.3. Константы - -| Константа | Значение | Описание | -|-----------|----------|----------| -| `MAX_INFLIGHT_BYTES` | 65536 | Начальное окно | -| `RETRANS_K1` | 32 | Множитель RTT для retrans timeout | -| `RETRANS_K2` | 32 | Множитель jitter для retrans timeout | -| `ACK_DELAY_TB` | 20 | Задержка отправки ACK (2ms) | -| `RTT_HISTORY_SIZE` | 10 | Размер истории RTT для jitter | -| `INFLIGHT_INITIAL_HASH_SIZE` | 1024 | Начальный размер хеш-таблиц | -| `MAX_INFLIGHT_SIZE` | 16384 | Максимум элементов в recv_q (защита от атак) | -| `ASM_BUF_MAX_SIZE` | 65536 | Максимальный размер буфера сборки | -| `BURST_PACKET_COUNT` | 12 | Пакетов в burst | -| `BURST_SKIP_COUNT` | 3 | Пропускаемых пакетов при замере gap | -| `MIN_BURST_INTERVAL_TB` | 5000 | Мин. интервал между burst (500ms) | -| `BURST_RESP_TIMEOUT_TB` | 20000 | Таймаут ответа на burst (2s) | -| `ETCP_METRICS_RTT_BUCKETS` | 12 | Число RTT-бакетов в гистограмме | -| `ETCP_METRICS_LOSS_BUCKETS` | 11 | Число loss-бакетов | -| `ETCP_METRICS_INTERVAL_TB` | 6000000 | Интервал снимков (10 мин) | -| `ETCP_METRICS_MIN_PACKETS` | 200 | Мин. пакетов для снимка | - -### 3.4. Секции пакета - -| Секция | Type | Размер | Описание | -|--------|------|--------|----------| -| `PAYLOAD` | 0x00 | 5 + data | Данные: type(1) + seq(4) + payload | -| `ACK` | 0x01 | 8 + N*8 | ACK: type(1) + count(1) + last_delivered_id(4) + rx_dup_count(2) + N×{seq(4) + recv_ts(2) + delay(2)} | -| `TIMESTAMP` | 0x06 | 5 | Channel timestamp: type(1) + ret_ts(2) + dt(2) | -| `MEAS_TS` | 0x07 | 9 | Burst measurement: type(1) + burst_id(2) + flags(1) + seq(1) + ts_us(2) + pkt_sz(2) | -| `MEAS_RESP` | 0x08 | 13 | Burst response: type(1) + burst_id(2) + valid(1) + gap_avg(4) + gap_min(4) + pkt_count(1) | -| `FILLER` | 0x09 | 3 + N | Filler padding: type(1) + len(2) + zeros(N) | -| `METRICS` | 0x0A | 2 + 64 + N | Metrics CSV: type(1) + csv_len(2) + ed25519_sig(64) + csv_data(N) | - -### 3.5. Зависимости - -| Модуль | Как используется | -|--------|-----------------| -| `etcp_connections.h` | ETCP_LINK, ETCP_DGRAM, ETCP_SOCKET, etcp_link_new/close, etcp_encrypt_send | -| `etcp_loadbalancer.h` | etcp_loadbalancer_select_link, etcp_loadbalancer_send | -| `etcp_router.h` | etcp_router_transit_queues_destroy, etcp_router_pause_retrans_for_node | -| `etcp_debug.h` | etcp_dump_pkt_sections | -| `etcp_connect.h` | etcp_connect_cancel_for_conn | -| `pkt_normalizer.h` | PKTNORM (pn_init, pn_deinit, pn_reset) | -| `secure_channel.h` | sc_context_t для шифрования | -| `routing.h` | routing_del_conn | -| `topo_group.h` | route_ping_cancel_for_conn | -| `../lib/ll_queue.h` | Все очереди | -| `../lib/u_async.h` | uasync_set_timeout, uasync_call_soon | -| `../lib/memory_pool.h` | memory_pool_alloc/free/init/destroy | -| `../lib/mem.h` | u_malloc, u_calloc, u_realloc, u_free, u_strdup | -| `../lib/debug_config.h` | DEBUG_ERROR/WARN/INFO/DEBUG/TRACE | diff --git a/src/transport_layer/etcp_loadbalancer_doc.md b/src/transport_layer/etcp_loadbalancer_doc.md deleted file mode 100644 index db033bb1..00000000 --- a/src/transport_layer/etcp_loadbalancer_doc.md +++ /dev/null @@ -1,51 +0,0 @@ -# ETCP Load Balancer (etcp_loadbalancer) - -## 1. Назначение -Балансировка исходящего трафика ETCP между несколькими линками (UDP-сокетами) внутри одного ETCP_CONN. Выбирает оптимальный линк по минимальному `inflight_bytes`, при равенстве — round-robin. Для каждого линка работает traffic shaper — ограничение полосы пропускания (bandwidth pacing) с таймерной задержкой и burst-режимом. - -## 2. Как пользоваться -Типовой сценарий: перед отправкой пакета вызывается `etcp_loadbalancer_send(dgram)`. Функция сама выбирает линк (если он ещё не задан), отправляет и обновляет shaper-нагрузку линка. - -```c -// Отправка пакета через loadbalancer -struct ETCP_DGRAM* dgram = memory_pool_alloc(inst->pkt_pool); -// ... заполнение dgram ... -etcp_loadbalancer_send(dgram); // сам выберет линк, зашифрует, отправит, освободит dgram -``` - -Если все линки заняты (shaper-таймер или inflight превышает лимит), пакет дропается. При освобождении линка shaper-таймер вызывает `loadbalancer_link_ready()`, который нотифицирует `ETCP_CONN->link_ready_for_send_fn`. - -```c -// Проверка статуса связи -if (etcp_loadbalancer_get_link_status(etcp)) { - // есть хотя бы один живой линк -} -``` - -**Нюансы:** -- Алгоритм выбора: min `inflight_bytes`, среди равных — round-robin (поле `last_rr_link` в ETCP_CONN) -- Shaper работает через `shaper_load_time_tb` / `shaper_sub_nanotime` — виртуальное время передачи, сравнивается с `now_tb + SHAPER_BURST_DELAY_TB` (1ms) -- Burst-пакеты (`link->burst_active`) обходят shaper -- При безлимитном bandwidth (`link->bandwidth == 0`) shaper пропускается -- `etcp_loadbalancer_send()` освобождает dgram через `memory_pool_free()` — после вызова dgram недействителен - -## 3. API - -### Структуры -Поля ETCP_LINK, используемые балансировщиком: -- `initialized` — линк готов к работе -- `link_status` — 1 = линк жив -- `inflight_bytes` — байт в полёте (минимизируемый критерий) -- `burst_active` — 1 = burst bypass всех ограничений -- `shaper_timer` — активный таймер задержки shaper'а -- `bandwidth` — Kbps, 0 = безлимит -- `shaper_load_time_tb`, `shaper_sub_nanotime` — аккумулированное время передачи (0.1ms timebase) - -### Функции -| Функция | Описание | -|---|---| -| `etcp_loadbalancer_select_link(etcp)` | Выбрать линк по min inflight_bytes с round-robin. NULL если все заняты | -| `etcp_loadbalancer_send(dgram)` | Выбрать линк → зашифровать/отправить → обновить shaper → освободить dgram | -| `loadbalancer_link_can_send(link)` | 1 если линк не заблокирован shaper'ом и burst не активен | -| `loadbalancer_link_ready(link)` | Уведомить ETCP_CONN о готовности линка (вызывает `link_ready_for_send_fn`) | -| `etcp_loadbalancer_get_link_status(etcp)` | 1 = есть живой линк, 0 = все недоступны | diff --git a/src/transport_layer/etcp_send_test.txt b/src/transport_layer/etcp_send_test.txt deleted file mode 100644 index 8a4fd95e..00000000 --- a/src/transport_layer/etcp_send_test.txt +++ /dev/null @@ -1,32 +0,0 @@ -сделай тест проверки отправкив в udp (etcp_loadbalancer) - - -1. инициализируешь etcp без линков, и один сокет для работы -2. вписываешь в link_ready_for_send_fn свою функцию -3. в тесте открываешь ответный сокет -4. добавляешь линк к etcp -5. - - - - -link initialization: - -client: - -1. init timer: - - шлём init with reset - -Получили response: - - отменяем init, переходим в обычный режим + keepalive - -link не отвечает: - - шлём init without reset - -server: - - просто ждёт init request. - как получил init, делает линк инициализированным. - -общая логика keepalive (client+server): - далее на init линках шлём keepalive со статусом remote keepalive. - если remote keepalive =0 - значит на другой стороне всё плохо и не шлём данные (только keepalive) diff --git a/src/transport_layer/stcp_doc.md b/src/transport_layer/stcp_doc.md deleted file mode 100644 index 36a72b5a..00000000 --- a/src/transport_layer/stcp_doc.md +++ /dev/null @@ -1,52 +0,0 @@ -# STCP (Secure TCP) — общая часть соединения - -## 1. Назначение - -STCP — потоковый TCP-протокол с X25519-ключеобменом и потоковым AES-CTR шифрованием (stream cipher). Обеспечивает безопасное TCP-соединение: клиент подключается, сервер принимает, выполняется ECDH-рукопожатие с обфускацией публичных ключей, после чего трафик шифруется AES-CTR с контрольными суммами CRC32. - -Модуль `stcp.c/stcp.h` содержит общие структуры, константы и базовый жизненный цикл соединения `struct stcp_conn`, используемый как клиентской (`stcp_client`), так и серверной (`stcp_server`) сторонами. - -## 2. Как пользоваться - -Типовая схема: вышестоящий код создаёт `struct stcp_conn` (сервер через `stcp_server_create`, клиент — через `stcp_client_connect`), устанавливает очереди приёма/передачи (`stcp_conn_set_tx_queue`, `stcp_conn_set_rx_queue`), коллбэк закрытия (`stcp_conn_set_on_close`). После рукопожатия (`STCP_STATE_DATA`) сообщения передаются через `rx_queue`/`tx_queue`: отправка — через ll_queue с коллбэком `tx_cb`, приём — данные раскладываются в `rx_queue`. - -Потоковое шифрование (`sc_stream_state`) не требует буферизации целых сообщений — XOR применяется побайтово к потоку, поэтому порядок отправки/приёма критичен. Для каждого направления создаётся отдельный stream (`stream_send`, `stream_recv`). - -Ключевые нюансы: -- Размер сообщения — `uint16_t` в префиксе, макс. 65535 байт. -- Каждое сообщение шифруется: 2 байта длины + данные + 4 байта CRC32. -- CRC32 проверяется после расшифровки для детектирования повреждений. -- При разрыве соединения или ошибке вызывается `on_close`. - -## 3. API - -### Состояния -- `STCP_STATE_INIT` — начальное (на клиенте) -- `STCP_STATE_HS_SERVER_WAIT` — сервер ждёт рукопожатия от клиента -- `STCP_STATE_DATA` — активный обмен данными (шифрованный) -- `STCP_STATE_CLOSED` / `STCP_STATE_ERROR` — завершение - -### stcp_conn — структура соединения -| Поле | Описание | -|------|----------| -| `sock` | TCP-сокет | -| `ua` | UASYNC event loop | -| `state` | Текущее состояние (enum stcp_state) | -| `is_server` | 1 = серверная сторона | -| `session_key[SC_SESSION_KEY_SIZE]` | Ключ X25519 ECDH (AES-128) | -| `stream_send` / `stream_recv` | Состояния потокового AES-CTR для отправки и приёма | -| `my_keys` | Локальные ключи X25519 | -| `peer_pubkey` | Публичный ключ пира (после рукопожатия) | -| `rx_queue` / `tx_queue` | ll_queue для приёма/отправки сообщений | -| `tx_cb` | Коллбэк очереди tx_queue | -| `recv_buf` | Буфер приёма TCP (динамический, до 128KB) | -| `send_buf` | Буфер отправки (для досыла при EAGAIN) | -| `hs_expected_len` / `hs_key_processed` | Состояние рукопожатия | -| `on_ready` | Вызывается после завершения рукопожатия | -| `on_close` | Вызывается при закрытии соединения | - -### Функции -- `stcp_conn_set_tx_queue(c, q)` — установить очередь отправки, привязывает `tx_cb` как коллбэк -- `stcp_conn_set_rx_queue(c, q)` — установить очередь приёма (в неё кладутся расшифрованные сообщения) -- `stcp_conn_set_on_close(c, cb, arg)` — установить коллбэк закрытия (err=0 — норма, иначе код ошибки) -- `stcp_conn_free(c)` — освободить сокет, буферы, очистить stream-ы. Если `allocated=1` — освободить и саму структуру diff --git a/tools/chat_client.py b/tools/chat_client.py index 5781cfe4..399b3053 100644 --- a/tools/chat_client.py +++ b/tools/chat_client.py @@ -107,3 +107,9 @@ class ChatClient: async def create_channel(self, name): return await self._request("create_channel", name=str(name)) + + async def invite_to(self, ch_id, node_id, pubkey, addr, proto=1): + return await self._request( + "invite_to", ch=str(ch_id), node_id=str(node_id), + pubkey=str(pubkey), addr=str(addr), proto=int(proto), + ) diff --git a/tools/chat_invite_test.py b/tools/chat_invite_test.py new file mode 100644 index 00000000..8ece84be --- /dev/null +++ b/tools/chat_invite_test.py @@ -0,0 +1,271 @@ +#!/usr/bin/env python3 +""" +chat_invite_test.py — headless test of the inviter-side invite flow. + +Сценарий: A (owner) говорит B (серверу) «добавь себе мою группу» +(CS_MSG_CHANNEL_INVITE), B авто-принимает (join_policy=0) и join+member_sync +завершается. Проверяем, что канал и мемберы появились на обеих сторонах. + +Usage: + python3 tools/chat_invite_test.py +""" + +import asyncio +import json +import os +import socket +import sys +import tempfile +import time + +sys.path.insert(0, os.path.dirname(__file__)) +from chat_client import ChatClient, ChatClientError + +UTUN_BIN = os.path.join(os.path.dirname(__file__), "..", "src", "utun") +READY_TIMEOUT = 3.0 # max wait for utun ready +SYNC_TIMEOUT = 10.0 # max wait for join + sync +REQUEST_TIMEOUT = 5.0 # per-request timeout +TOTAL_TIMEOUT = 30.0 # overall test timeout + + +def find_free_port(): + s = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + s.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + s.bind(("127.0.0.1", 0)) + port = s.getsockname()[1] + s.close() + return port + + +def write_config(path, etcp_port, ctrl_port, db_subdir, join_policy=None): + chat_section = "" + if join_policy is not None: + chat_section = f"\n[chat]\njoin_policy={join_policy}\n" + content = f"""[global] +db_path={db_subdir} + +[server: srv] +addr=127.0.0.1:{etcp_port} +type=public + +[chatserver] +db_path={db_subdir} +headless_control_bind=127.0.0.1:{ctrl_port} +{chat_section} +[allowed_keys] +allow_all=1 +""" + with open(path, "w") as f: + f.write(content) + + +def read_config_key(path, key): + try: + with open(path, "r") as f: + for line in f: + line = line.strip() + if line.startswith(key + "="): + return line[len(key) + 1:] + except OSError: + return None + return None + + +async def wait_config_key(path, key, timeout=READY_TIMEOUT): + deadline = time.monotonic() + timeout + while time.monotonic() < deadline: + val = read_config_key(path, key) + if val: + return val + await asyncio.sleep(0.1) + return None + + +async def kill_proc(proc, label): + if proc is None or proc.returncode is not None: + return + try: + proc.terminate() + try: + await asyncio.wait_for(proc.wait(), timeout=3.0) + except asyncio.TimeoutError: + proc.kill() + await proc.wait() + except ProcessLookupError: + pass + + +async def wait_ready(cli, timeout=READY_TIMEOUT): + deadline = time.monotonic() + timeout + while time.monotonic() < deadline: + try: + await cli.ping() + return True + except ChatClientError: + await asyncio.sleep(0.1) + return False + + +def check(name, expr, detail=""): + if not expr: + detail = f" ({detail})" if detail else "" + raise AssertionError(f"FAIL: {name}{detail}") + print(f" OK: {name}") + + +async def main(): + etcp_a = find_free_port() + etcp_b = find_free_port() + ctrl_a = find_free_port() + ctrl_b = find_free_port() + + print(f"ports: etcp={etcp_a},{etcp_b} ctrl={ctrl_a},{ctrl_b}") + + tmpdir = tempfile.mkdtemp(prefix="utun_invite_test_") + db_a = os.path.join(tmpdir, "db_a") + db_b = os.path.join(tmpdir, "db_b") + os.makedirs(db_a, exist_ok=True) + os.makedirs(db_b, exist_ok=True) + + config_a = os.path.join(tmpdir, "a.conf") + config_b = os.path.join(tmpdir, "b.conf") + write_config(config_a, etcp_a, ctrl_a, db_a) # A: owner + write_config(config_b, etcp_b, ctrl_b, db_b, join_policy=0) # B: server, autojoin + + proc_a = None + proc_b = None + + try: + print("\n--- Starting utun ---") + log_a = os.path.join(tmpdir, "utun_a.log") + log_b = os.path.join(tmpdir, "utun_b.log") + pid_a = os.path.join(tmpdir, "utun_a.pid") + pid_b = os.path.join(tmpdir, "utun_b.pid") + + proc_a = await asyncio.create_subprocess_exec( + UTUN_BIN, "-f", "-p", pid_a, "-l", log_a, "-c", config_a, + stdout=asyncio.subprocess.DEVNULL, stderr=asyncio.subprocess.DEVNULL, + ) + proc_b = await asyncio.create_subprocess_exec( + UTUN_BIN, "-f", "-p", pid_b, "-l", log_b, "-c", config_b, + stdout=asyncio.subprocess.DEVNULL, stderr=asyncio.subprocess.DEVNULL, + ) + + print(f" proc_a pid={proc_a.pid} proc_b pid={proc_b.pid}") + await asyncio.sleep(0.5) + + print("\n--- Connecting to headless control ---") + async with ChatClient(port=ctrl_a, timeout=REQUEST_TIMEOUT) as cli_a: + if not await wait_ready(cli_a): + raise RuntimeError(f"node A not ready after {READY_TIMEOUT}s") + print(f" node A ready") + + async with ChatClient(port=ctrl_b, timeout=REQUEST_TIMEOUT) as cli_b: + if not await wait_ready(cli_b): + raise RuntimeError(f"node B not ready after {READY_TIMEOUT}s") + print(f" node B ready") + + # ── Create channel on A ── + print("\n--- Create channel ---") + await cli_a.create_channel("TestGroup") + channels = await cli_a.channels() + check("channel created", isinstance(channels, list) and len(channels) == 1, + f"channels={json.dumps(channels)}") + ch_id = str(channels[0]["id"]) + print(f" channel_id={ch_id} name={channels[0]['name']}") + + # ── Read B's node_id + pubkey from config (written back by utun) ── + print("\n--- Read B identity ---") + b_node_id = await wait_config_key(config_b, "my_node_id") + b_pubkey = await wait_config_key(config_b, "my_public_key") + if not b_node_id or not b_pubkey: + raise RuntimeError("failed to read B's node_id/pubkey from config") + b_node_id = "0x" + b_node_id + print(f" B node_id={b_node_id} pubkey={b_pubkey[:16]}...") + + # ── A invites B to add the channel (inviter side) ── + print("\n--- Invite B (inviter → CHANNEL_INVITE) ---") + resp = await cli_a.invite_to(ch_id, b_node_id, b_pubkey, f"127.0.0.1:{etcp_b}", proto=1) + check("invite_to accepted", isinstance(resp, dict) and resp.get("inviting"), + f"resp={json.dumps(resp)}") + + # ── Wait for B to join + member_sync ── + print(f"\n--- Wait join + sync (max {SYNC_TIMEOUT}s) ---") + deadline = time.monotonic() + SYNC_TIMEOUT + member_count_b = 0 + members_b = [] + while time.monotonic() < deadline: + try: + members_b = await cli_b.members(ch_id) + member_count_b = len(members_b) if isinstance(members_b, list) else 0 + if member_count_b >= 2: + break + except ChatClientError: + pass + await asyncio.sleep(0.15) + + print(f"\n--- Members on B ({member_count_b}) ---") + for m in (members_b if isinstance(members_b, list) else []): + print(f" {m['node_id']} name={m.get('name','?')} online={m.get('online')} connected={m.get('connected')}") + check("2 members on B (B joined A)", member_count_b >= 2, + f"got {member_count_b} after {SYNC_TIMEOUT}s") + + await asyncio.sleep(0.5) + + members_a = await cli_a.members(ch_id) + print(f"\n--- Members on A ({len(members_a) if isinstance(members_a, list) else '?'}) ---") + for m in members_a: + print(f" {m['node_id']} name={m.get('name','?')} online={m.get('online')} connected={m.get('connected')}") + check("2 members on A (B visible to A)", isinstance(members_a, list) and len(members_a) >= 2, + f"got {len(members_a) if isinstance(members_a, list) else '?'}") + + channels_b = await cli_b.channels() + check("channel present on B", isinstance(channels_b, list) and len(channels_b) >= 1, + f"channels_b={json.dumps(channels_b)}") + + print("\n=== TEST PASSED ===") + return 0 + + except Exception as e: + print(f"\n=== TEST FAILED: {e} ===", file=sys.stderr) + import traceback + traceback.print_exc() + print(f"\n--- logs ---", file=sys.stderr) + for lp in ["utun_a.log", "utun_b.log"]: + lp = os.path.join(tmpdir, lp) + if os.path.exists(lp): + print(f"\n===== {lp} (tail) =====", file=sys.stderr) + try: + with open(lp, "r") as f: + lines = f.readlines() + for ln in lines[-80:]: + sys.stderr.write(ln) + except OSError: + pass + return 1 + + finally: + print("\n--- Cleanup ---") + await kill_proc(proc_a, "proc_a") + await kill_proc(proc_b, "proc_b") + + for f in [config_a, config_b]: + try: os.unlink(f) + except OSError: pass + try: os.rmdir(db_a) + except OSError: pass + try: os.rmdir(db_b) + except OSError: pass + try: os.rmdir(tmpdir) + except OSError: pass + print(f" temp dir cleaned: {tmpdir}") + + +if __name__ == "__main__": + async def _run(): + try: + return await asyncio.wait_for(main(), timeout=TOTAL_TIMEOUT) + except asyncio.TimeoutError: + print(f"\n=== TEST FAILED: total timeout {TOTAL_TIMEOUT}s ===", file=sys.stderr) + return 1 + sys.exit(asyncio.run(_run())) diff --git a/tools/chatcli b/tools/chatcli index 839e78c0..c24b4f1e 100755 --- a/tools/chatcli +++ b/tools/chatcli @@ -17,6 +17,8 @@ # invite create invite link # connect join channel via invite link # create create new channel +# invite_to [proto] +# invite a node to add our channel # listen interactive event listener import sys, os, json, socket, struct @@ -104,6 +106,12 @@ def cmd_connect(link): def cmd_create(name): _req("create_channel", name=name) +def cmd_invite_to(ch_id, node_id, pubkey, addr, proto=None): + params = {"ch": ch_id, "node_id": node_id, "pubkey": pubkey, "addr": addr} + if proto: + params["proto"] = int(proto) + _req("invite_to", **params) + def cmd_listen(): print(f"Listening on {DEFAULT_HOST}:{DEFAULT_PORT} (Ctrl+C to quit)") s = _connect() @@ -161,6 +169,7 @@ def main(): elif cmd == "invite": cmd_invite(*args[1:2] if len(args) > 1 else (_die("usage: invite "),)) elif cmd in ("connect","join"): cmd_connect(*args[1:2] if len(args) > 1 else (_die("usage: connect "),)) elif cmd == "create": cmd_create(*args[1:2] if len(args) > 1 else (_die("usage: create "),)) + elif cmd == "invite_to": cmd_invite_to(*args[1:6] if len(args) > 4 else (_die("usage: invite_to [proto]"),)) elif cmd == "listen": cmd_listen() else: print(f"Unknown command: {cmd}\nUse --help for usage", file=sys.stderr); sys.exit(1) except socket.timeout: