diff --git a/BUG_bgp_reconnect.md b/BUG_bgp_reconnect.md new file mode 100644 index 00000000..86499074 --- /dev/null +++ b/BUG_bgp_reconnect.md @@ -0,0 +1,89 @@ +# Задача: BGP NODEINFO теряется при флапающем реконнекте (chatgui member online не возвращается) + +## Кратко + +После рестарта Android-пира и реконнекта локальный chatgui показывает мембера **offline** и не возвращает **online**. +Корневая причина — **не в chat/member коде**, а в ETCP: BGP-пакет NODEINFO (347 байт) расшифровывается, но **не доходит до BGP-обработчика** `topo_group_receive_cbk()`. + +Нужно: **точно локализовать точку потери NODEINFO в receive-пути ETCP и починить её.** + +## Что уже сделано (исправлено и работает) + +1. `tools/chatgui/src/mainwindow.cpp` — **исправлен** `onMemberUpdatedCallback`/`onMemberRemovedCallback`: проверка `len < 1+64+1+79` отбрасывала все события с channel_id короче 64 симв. Теперь `len < 1` + точная проверка. (Это была причина «не обновляется вообще».) +2. `src/chat/chat_member.c/h` — добавлен `chat_core_member_online(ch_id,node_id)` = присутствие в `topo_group` (`topo_node_find_by_id`). Используется в `chat_core_get_member_list`/`get_single_member`. Убран `nodes.online` из SQL и блок `connected→2`. +3. `src/chat/chat_core.c`, `src/chat/chat_sync.c`, `src/chat/member_sync.c/h` — удалён онлайн-статус из БД + merkle-online gossip (`member_sync_set_online`, `_ms_apply_update`, `push_update`). +4. `src/chat/chat_status.c` + `tools/chatgui/src/accountlist.cpp` — флаг «bgp present» (0x08) для detail-панели. +5. `tools/chatgui/src/nodespage.cpp` — убрана колонка Online из диагностики. + +Эти правки корректны; статус работает на чистом connect/disconnect/стабильном реконнекте. + +## Найденная причина (подтверждено логами) + +- Android **шлёт** свой NODEINFO: `send_nodeinfo: node 1657b3ea8281c0a8 ver=2` (bgp-debug на Android). +- Локально пакет **расшифровывается**: `decrypt: code=01 dlen=347 plen=350` (347-байтный BGP-пакет, `code=0x01` = `ETCP_ID_TOPO_ENTRY`). +- Но `BGP recv NODEINFO nid=1657` (строка 193 в `topo_group.c`) **не появляется** — до обработчика пакет не доходит. +- При этом маленький `TABLE_COMPLETE` (10 байт) **доходит** (виден в логе как «mislabeled» `nid=aaaaaaaaaadeadbe` — это строка 193 логирует не-NODEINFO пакет как NODEINFO, garbage-поля). +- Признак: в стабильных рестартах (18:18/18:20) узел **добавляется** каждый раз, в флапающем (18:07) — нет. Т.е. потеря привязана к флапу (TCP-линки + ETCP reinit). + +## Receive-путь ETCP (где искать потерю) + +``` +decrypt (etcp_connections.c:2262, лог "decrypt: code=.. dlen=..") + → if (link_state==CONNECTED) etcp_conn_input(pkt) + else memory_pool_free(pkt) ← добавлен WARN "decrypted pkt DROPPED" + → etcp_conn_input (etcp.c:1506, "RX pkt dlen=%d") + → нормализатор (reassembly) → int_queue + → etcp_int_recv (etcp_api.c) ← добавлен WARN "int_recv BGP id=0x01 len=.." + → dispatch по route id (0x01) → topo_group_receive_cbk (topo_group.c:193 "BGP recv NODEINFO") +``` + +## Уже добавленный дебаг (для локализации) + +1. `src/transport_layer/etcp_connections.c` (после дешифровки): + ```c + } else { + DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] decrypted pkt DROPPED: link_state=%d code=0x%02x dlen=%u", + link->etcp->log_name, link->link_state, pkt_code, pkt->data_len); + memory_pool_free(e_sock->instance->pkt_pool, pkt); + } + ``` +2. `src/transport_layer/etcp_api.c` (`etcp_int_recv`, после чтения `id`): + ```c + uint8_t id = e->dgram[0]; + if (id == 0x01) DEBUG_WARN(DEBUG_CATEGORY_ETCP, "int_recv BGP id=0x01 conn=%s len=%zu", conn->log_name, e->len); + ``` +3. Android: `tools/chatgui-android/libutun_lite/instance_lite.c` добавлено + `debug_set_category_level(DEBUG_CATEGORY_BGP, DEBUG_LEVEL_DEBUG);` (пересобрать APK: `cd tools/chatgui-android && ./build.sh`). +4. chatgui конфиг `tools/chatgui/build/vibechat.cfg` — включены `bgp=debug`, `chat_sync=debug`, `member_sync=debug`, `debug=debug`. + +Как читать (по таймстампу флапа): +- есть `decrypt code=01 dlen=..` + `DROPPED` → потеря в проверке `link_state==CONNECTED`; +- есть `decrypt` + нет `DROPPED` + нет `int_recv BGP` → потеря в `etcp_conn_input`/нормализаторе; +- есть `int_recv BGP` + нет `BGP recv NODEINFO` → потеря в BGP-обработчике. + +## Задача для агента + +1. Пересобрать chatgui (`cd tools/chatgui && cmake --build build -j4`) и Android APK (`cd tools/chatgui-android && ./build.sh`). +2. Запустить chatgui (DISPLAY=:0, `./vibechat`), воспроизвести флап реконнекта: + ```bash + adb shell am force-stop com.utun.chat && adb shell am start -n com.utun.chat/.MainActivity + ``` + (повторить несколько раз; флап ловится не каждый раз). +3. В логах (`tools/chatgui/build/chatgui.log`) найти потерю NODEINFO и **точно определить точку** по таблице выше. +4. Починить первопричину (варианты): + - буферизовать расшифрованный пакет до `LINK_STATE_CONNECTED` вместо `memory_pool_free`; + - либо гарантированный BGP-retry (после conn UP, если узел пира не в `group->nodes` через N мс — повторно `topo_group_send_table_request`). +5. Убрать весь диагностический мусор: + - WARN "DROPPED" (etcp_connections.c), WARN "int_recv BGP" (etcp_api.c); + - DEBUG-логи в `chat_member.c` (`on_member_props_changed`), `chat_sync.c` (`cs_on_peer_status_changed`), `memberlistmodel.cpp`; + - `instance_lite.c` BGP-debug; + - `vibechat.cfg` — вернуть категории в закомментированное состояние. +6. Прогнать: connect→online, disconnect→offline, реконнект (флап)→online. `./check.sh` (или `cd src && make -j4`). + +## Воспроизведение / ключевые артефакты + +- Пир: Android SM-A525F (`com.utun.chat`), node_id `0x1657b3ea8281c0a8`. +- Локальный: chatgui (`tools/chatgui/build/vibechat`), node_id `0x24cd036a6e659b9a`. +- Канал: `ch_id=7206723622466219923`, `group_id=0x64036cf3a783ef93`. +- Лог chatgui: `tools/chatgui/build/chatgui.log` (debug_file из конфига). +- Лог Android: `adb logcat -d | grep -iE "send_nodeinfo|NOT in registry|NODEINFO|DROPPED|int_recv"`. diff --git a/src/chat/chat_core.c b/src/chat/chat_core.c index e5678c83..dccdd45e 100644 --- a/src/chat/chat_core.c +++ b/src/chat/chat_core.c @@ -121,7 +121,6 @@ int chat_core_init(struct UTUN_INSTANCE* inst, const char* db_path) { topo_node_sqlite_node_update_verified(g_cc.db, g_cc.my_node_id, my_name, inst->my_keys.public_key, inst->my_ed25519_pubkey, (uint64_t)now_sec, (time_t)now_sec); - topo_node_sqlite_node_set_online(g_cc.db, g_cc.my_node_id, 1); sqlite3_stmt* st = NULL; sqlite3_prepare_v2(g_cc.db, "INSERT OR REPLACE INTO local_identity(id,node_id,name,x25519_pubkey,ed25519_pubkey,created_at,updated_at)" @@ -311,10 +310,13 @@ int chat_core_get_channels_json(char* buf, size_t buf_size, size_t* out_len) { char peers_tbl[80]; peers_table_name(ch_id, peers_tbl, sizeof(peers_tbl)); int peer_count = 0, online_count = 0; sqlite3_stmt* ps = NULL; - char psql[200]; snprintf(psql, sizeof(psql), - "SELECT COUNT(*), SUM(CASE WHEN n.online=1 THEN 1 ELSE 0 END) FROM \"%s\" p LEFT JOIN nodes n ON p.node_id=n.node_id", peers_tbl); + char psql[200]; snprintf(psql, sizeof(psql), "SELECT node_id FROM \"%s\"", peers_tbl); if (sqlite3_prepare_v2(g_cc.db, psql, -1, &ps, NULL) == SQLITE_OK) { - if (sqlite3_step(ps) == SQLITE_ROW) { peer_count = sqlite3_column_int(ps, 0); online_count = sqlite3_column_int(ps, 1); } + while (sqlite3_step(ps) == SQLITE_ROW) { + uint64_t pid = (uint64_t)sqlite3_column_int64(ps, 0); + peer_count++; + if (chat_core_member_online(ch_id, pid)) online_count++; + } sqlite3_finalize(ps); } @@ -386,19 +388,19 @@ int chat_core_get_members_json(const char* ch_id, char* buf, size_t buf_size, si *w++ = '['; const char* sep = ""; char sql[256]; snprintf(sql, sizeof(sql), - "SELECT p.node_id, n.name, n.online, n.x25519_pubkey, n.ed25519_pubkey, p.adm_tags " + "SELECT p.node_id, n.name, n.x25519_pubkey, n.ed25519_pubkey, p.adm_tags " "FROM \"%s\" p LEFT JOIN nodes n ON p.node_id=n.node_id", peers_tbl); sqlite3_stmt* st = NULL; if (sqlite3_prepare_v2(g_cc.db, sql, -1, &st, NULL) != SQLITE_OK) { *out_len = 0; return -1; } while (sqlite3_step(st) == SQLITE_ROW) { uint64_t nid = (uint64_t)sqlite3_column_int64(st, 0); const char* name = (const char*)sqlite3_column_text(st, 1); - int online = sqlite3_column_int(st, 2); - const uint8_t* x25519 = (const uint8_t*)sqlite3_column_blob(st, 3); - const uint8_t* ed25519 = sqlite3_column_blob(st, 4); - int x25519_len = sqlite3_column_bytes(st, 3); - int ed25519_len = sqlite3_column_bytes(st, 4); - const char* adm_tags = (const char*)sqlite3_column_text(st, 5); + int online = chat_core_member_online(ch_id, nid); + const uint8_t* x25519 = (const uint8_t*)sqlite3_column_blob(st, 2); + const uint8_t* ed25519 = sqlite3_column_blob(st, 3); + int x25519_len = sqlite3_column_bytes(st, 2); + int ed25519_len = sqlite3_column_bytes(st, 3); + const char* adm_tags = (const char*)sqlite3_column_text(st, 4); char esc_name[256], x25519_hex[65], ed25519_hex[65], esc_tags[384]; json_escape(name ? name : "", esc_name, sizeof(esc_name)); x25519_hex[0] = ed25519_hex[0] = '\0'; diff --git a/src/chat/chat_member.c b/src/chat/chat_member.c index 288787f6..1d7c5300 100644 --- a/src/chat/chat_member.c +++ b/src/chat/chat_member.c @@ -11,11 +11,26 @@ #include "member_sync.h" #include "../utun_instance.h" +#include "../routing_layer/topo_group.h" #include "../../lib/mem.h" #include "../../lib/ll_queue.h" #include "../../lib/json_flat.h" #include "../../lib/platform_compat.h" +#include + +/* ─── live-онлайн: узел присутствует в BGP/topo_group канала ─── */ + +int chat_core_member_online(const char* ch_id, uint64_t node_id) { + if (!g_cc.inst || !ch_id || !ch_id[0]) return 0; + if (node_id == g_cc.my_node_id) return 1; /* self всегда онлайн (local_node вне group->nodes) */ + if (!g_cc.inst->topo_groups) return 0; + uint64_t group_id = strtoull(ch_id, NULL, 10); + struct TOPO_GROUP* g = topo_groups_find(g_cc.inst->topo_groups, group_id); + if (!g || g->group_type != TOPO_GROUP_TYPE_CHAT) return 0; + return topo_node_find_by_id(g, node_id) != NULL; +} + /* ─── парсинг adm_tags в flags ─── */ static uint8_t parse_adm_tags_flags(const char* adm_tags) { @@ -91,11 +106,11 @@ int chat_core_get_member_list(const char* ch_id, uint8_t** out, int* count) { sqlite3_stmt* st = NULL; char sql[512]; snprintf(sql, sizeof(sql), - "SELECT p.node_id, COALESCE(n.online,0), COALESCE(n.name,'')," + "SELECT p.node_id, COALESCE(n.name,'')," " COALESCE(p.local_nick,''), COALESCE(p.adm_tags,'')," " COALESCE(p.storage,0)" " FROM \"%s\" p LEFT JOIN nodes n ON p.node_id=n.node_id" - " ORDER BY n.online DESC, p.node_id ASC", peers_tbl); + " ORDER BY p.node_id ASC", peers_tbl); if (sqlite3_prepare_v2(g_cc.db, sql, -1, &st, NULL) != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "chat_member: get_list — query failed ch=%s err=%s", @@ -115,11 +130,11 @@ int chat_core_get_member_list(const char* ch_id, uint8_t** out, int* count) { } struct chat_member_display* m = &members[cnt]; m->node_id = (uint64_t)sqlite3_column_int64(st, 0); - m->online = (uint8_t)sqlite3_column_int(st, 1); - const char* node_name = (const char*)sqlite3_column_text(st, 2); - const char* local_nick = (const char*)sqlite3_column_text(st, 3); - const char* adm_tags = (const char*)sqlite3_column_text(st, 4); - m->storage = (uint8_t)sqlite3_column_int(st, 5); + m->online = (uint8_t)chat_core_member_online(ch_id, m->node_id); + const char* node_name = (const char*)sqlite3_column_text(st, 1); + const char* local_nick = (const char*)sqlite3_column_text(st, 2); + const char* adm_tags = (const char*)sqlite3_column_text(st, 3); + m->storage = (uint8_t)sqlite3_column_int(st, 4); m->is_self = (uint8_t)(m->node_id == my_id ? 1 : 0); m->flags = parse_adm_tags_flags(adm_tags); m->rtt = 0xFFFF; @@ -174,12 +189,16 @@ int chat_core_get_member_list(const char* ch_id, uint8_t** out, int* count) { sqlite3_finalize(st); } - /* 3. проверяем connected (активное ETCP-соединение) */ - for (int i = 0; i < cnt; i++) { - if (!members[i].online) continue; - if (g_cc.inst->connections - && queue_find_data_by_index(g_cc.inst->connections, (const uint8_t*)&members[i].node_id)) - members[i].online = 2; /* 2 = connected (online + active conn) */ + /* 3. сортируем: онлайн первыми, затем по node_id (стабильно) */ + for (int i = 1; i < cnt; i++) { + struct chat_member_display key = members[i]; + int j = i - 1; + while (j >= 0 && ((members[j].online < key.online) + || (members[j].online == key.online && members[j].node_id > key.node_id))) { + members[j + 1] = members[j]; + j--; + } + members[j + 1] = key; } /* 4. сериализуем в wire-формат */ @@ -207,7 +226,7 @@ int chat_core_get_single_member(const char* ch_id, uint64_t node_id, uint8_t* ou sqlite3_stmt* st = NULL; char sql[384]; snprintf(sql, sizeof(sql), - "SELECT COALESCE(n.online,0), COALESCE(n.name,'')," + "SELECT COALESCE(n.name,'')," " COALESCE(p.local_nick,''), COALESCE(p.adm_tags,'')," " COALESCE(p.storage,0)" " FROM \"%s\" p LEFT JOIN nodes n ON p.node_id=n.node_id" @@ -219,11 +238,11 @@ int chat_core_get_single_member(const char* ch_id, uint64_t node_id, uint8_t* ou struct chat_member_display m; m.node_id = node_id; - m.online = (uint8_t)sqlite3_column_int(st, 0); - const char* node_name = (const char*)sqlite3_column_text(st, 1); - const char* local_nick = (const char*)sqlite3_column_text(st, 2); - const char* adm_tags = (const char*)sqlite3_column_text(st, 3); - m.storage = (uint8_t)sqlite3_column_int(st, 4); + m.online = (uint8_t)chat_core_member_online(ch_id, node_id); + const char* node_name = (const char*)sqlite3_column_text(st, 0); + const char* local_nick = (const char*)sqlite3_column_text(st, 1); + const char* adm_tags = (const char*)sqlite3_column_text(st, 2); + m.storage = (uint8_t)sqlite3_column_int(st, 3); m.is_self = (uint8_t)(node_id == my_id ? 1 : 0); m.flags = parse_adm_tags_flags(adm_tags); m.rtt = 0xFFFF; @@ -253,11 +272,6 @@ int chat_core_get_single_member(const char* ch_id, uint64_t node_id, uint8_t* ou sqlite3_finalize(st); } - /* connected */ - if (m.online && g_cc.inst->connections - && queue_find_data_by_index(g_cc.inst->connections, (const uint8_t*)&node_id)) - m.online = 2; - chat_core_serialize_member(&m, out); return 0; } @@ -320,9 +334,9 @@ void chat_core_request_member_rtt_trampoline(void* arg) { u_free(members); } -/* ─── node_props_changed callback (member_sync → chat_event) ─── */ +/* ─── node_props_changed callback (BGP-события узла → MEMBER_UPDATED в GUI) ─── */ -static void on_adm_tags_changed(uint64_t node_id, const char* adm_tags, const char* channel_id, void* arg) { +static void on_member_props_changed(uint64_t node_id, const char* adm_tags, const char* channel_id, void* arg) { (void)adm_tags; (void)arg; if (!g_cc.initialized || !g_cc.inst || !g_cc.db || !channel_id || !channel_id[0]) return; @@ -348,6 +362,6 @@ static void on_adm_tags_changed(uint64_t node_id, const char* adm_tags, const ch } void chat_member_init(void) { - member_sync_add_props_cbk(on_adm_tags_changed, NULL); - DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "chat_member: initialized (adm_tags → MEMBER_UPDATED)"); + member_sync_add_props_cbk(on_member_props_changed, NULL); + DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "chat_member: initialized (BGP node events → MEMBER_UPDATED)"); } diff --git a/src/chat/chat_member.h b/src/chat/chat_member.h index c0a43bef..945db58e 100644 --- a/src/chat/chat_member.h +++ b/src/chat/chat_member.h @@ -38,6 +38,10 @@ int chat_core_get_member_list(const char* ch_id, uint8_t** out, int* count); /* Получить одного мембера по node_id. out = буфер размером CHAT_MEMBER_DISPLAY_SIZE. */ int chat_core_get_single_member(const char* ch_id, uint64_t node_id, uint8_t* out); +/* Live-онлайн: узел присутствует в BGP/topo_group канала (self всегда онлайн). + Не читает nodes.online — статус не хранится в БД. */ +int chat_core_member_online(const char* ch_id, uint64_t node_id); + /* Сериализовать из внутренней структуры в wire-формат */ struct chat_member_display { uint64_t node_id; diff --git a/src/chat/chat_profile.c b/src/chat/chat_profile.c index 7797b2bd..40290476 100644 --- a/src/chat/chat_profile.c +++ b/src/chat/chat_profile.c @@ -9,6 +9,7 @@ #include "chat_member.h" #include "../../lib/json_flat.h" #include "../routing_layer/topo_node_sqlite.h" +#include "../routing_layer/topo_node.h" #include "member_sync.h" #include "../utun_instance.h" @@ -102,8 +103,11 @@ void chat_core_update_my_member(void) { } void chat_core_sync_my_addresses(void) { - /* Адреса теперь синхронизируются BGP/topo_group (topo_node_update_my_addresses). - Здесь оставлено только обновление userinfo/подписи своего мембера. */ + /* Адреса синхронизируются BGP/topo_group через topo_node_update_my_addresses + (реестр + БД node_addresses + классификация node_type/storage + рассылка nodeinfo). + Вызывается после готовности БД, поэтому здесь гарантированно персистит свои адреса. */ + if (g_cc.inst) + topo_node_update_my_addresses(g_cc.inst); chat_core_update_my_member(); } diff --git a/src/chat/chat_status.c b/src/chat/chat_status.c index 6c487405..373067d3 100644 --- a/src/chat/chat_status.c +++ b/src/chat/chat_status.c @@ -345,6 +345,16 @@ static void collect_member_detail(uint64_t node_id) { if (conn->links_up) flags |= 2; if (conn->initialized) flags |= 4; } + /* bit 3 (0x08): узел присутствует в BGP/topo_group — live-онлайн (self всегда онлайн) */ + { + int bgp_present = (node_id == inst->node_id) ? 1 : 0; + struct ll_queue* gl = inst->topo_groups ? inst->topo_groups->group_list : NULL; + for (struct ll_entry* ge = gl ? gl->head : NULL; !bgp_present && ge; ge = ge->next) { + struct TOPO_GROUP* group = (struct TOPO_GROUP*)ge->data; + if (group && topo_node_find_by_id(group, node_id)) bgp_present = 1; + } + if (bgp_present) flags |= 8; + } uint16_t own_tcp_active = 0; if (node_id == inst->node_id) { diff --git a/src/chat/chat_sync.c b/src/chat/chat_sync.c index 03cb87b3..d1c36e55 100644 --- a/src/chat/chat_sync.c +++ b/src/chat/chat_sync.c @@ -370,35 +370,10 @@ static void cs_schedule_sync(struct chat_sync* cs) { DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "%s: sync scheduled in %ums", CS_ID, CS_SYNC_INTERVAL_MS); } -/* ── Unified peer status change (локальная БД + GUI + push_update + schedule_sync) ── */ - -static void cs_on_remote_status_changed(uint64_t peer, int online) { - if (!g_cs) return; - DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "%s: remote status changed peer=%016llx online=%d", - CS_ID, (unsigned long long)peer, online); - for (int i = 0; i < g_cs->channel_count; i++) { - for (int j = 0; j < g_cs->channels[i].peer_count; j++) { - if (g_cs->channels[i].peer_ids[j] == peer) { - const char* ch_id = g_cs->channels[i].channel_id; - size_t cl = strlen(ch_id); - uint8_t evt[1 + 64 + 1 + CHAT_MEMBER_DISPLAY_SIZE]; - evt[0] = (uint8_t)cl; - memcpy(evt + 1, ch_id, cl); - evt[1 + cl] = 1; - if (chat_core_get_single_member(ch_id, peer, evt + 1 + cl + 1) == 0) - chat_event_post(CHAT_EVT_MEMBER_UPDATED, evt, 1 + (int)cl + 1 + CHAT_MEMBER_DISPLAY_SIZE); - cs_post_channel_online(g_cs, ch_id); - break; - } - } - } -} +/* ── Peer status change (GUI + online-бар + schedule_sync; онлайн в БД не хранится) ── */ static void cs_on_peer_status_changed(uint64_t peer, int online) { - /* 1. Локальная БД (онлайн НЕ в merkle sync — только локальная запись) */ - int online_changed = member_sync_set_online(g_cs->inst, peer, online); - - /* 2. GUI + online-бар */ + /* 1. GUI + online-бар */ int found_in_channel = 0; for (int i = 0; i < g_cs->channel_count; i++) { for (int j = 0; j < g_cs->channels[i].peer_count; j++) { @@ -413,12 +388,6 @@ static void cs_on_peer_status_changed(uint64_t peer, int online) { if (chat_core_get_single_member(ch_id, peer, mevt + 1 + cl + 1) == 0) chat_event_post(CHAT_EVT_MEMBER_UPDATED, mevt, 1 + (int)cl + 1 + CHAT_MEMBER_DISPLAY_SIZE); cs_post_channel_online(g_cs, ch_id); - - /* 3. Рассылка дельты всем synced-соседям — только если online реально изменился */ - if (online_changed) { - uint8_t st = (uint8_t)online; - merkle_sync_push_update(g_cs->inst, g_cs->channels[i].channel_id, peer, 0x01, &st, 1); - } break; } } @@ -427,7 +396,7 @@ static void cs_on_peer_status_changed(uint64_t peer, int online) { DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "%s: peer_status_changed peer=0x%016llx online=%d channels=%d found=%d", CS_ID, (unsigned long long)peer, online, g_cs->channel_count, found_in_channel); - /* 4. Запланировать member_sync_start (throttled) */ + /* 2. Запланировать member_sync_start (throttled) */ cs_schedule_sync(g_cs); } @@ -462,12 +431,6 @@ static void _on_invite_sync_done(uint64_t peer, const char* ns, int result, void } else { cs_post_channel_online(g_cs, sa->ch_id); cs_resume_db_sync(sa->ch_id); - if (member_sync_set_online(g_cs->inst, sa->node_id, 1)) { - uint8_t one = 1; - merkle_sync_push_update(g_cs->inst, sa->ch_id, sa->node_id, 0x01, &one, 1); - DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: invite_sync push_update node=0x%016llx online=1 ch=%s", - CS_ID, (unsigned long long)sa->node_id, sa->ch_id); - } } } u_free(sa); @@ -658,7 +621,6 @@ int chat_sync_init(struct UTUN_INSTANCE* inst) { cs->join_timer = NULL; member_sync_init(inst); - member_sync_set_node_updated_cb(cs_on_remote_status_changed); DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: initialized", CS_ID); return 0; diff --git a/src/chat/member_sync.c b/src/chat/member_sync.c index 9f9292fd..3703c7c3 100644 --- a/src/chat/member_sync.c +++ b/src/chat/member_sync.c @@ -721,25 +721,11 @@ static int _member_apply_items(void* ctx, const char* ns, uint64_t from_peer, return 0; } -static member_sync_node_updated_fn g_node_updated_cb = NULL; - -static int _ms_apply_update(void* ctx, const char* ns, uint64_t key, - uint8_t type, const uint8_t* data, size_t len) { - (void)ns; - struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)ctx; - if (type == 0x01 && len >= 1) { - int online = data[0]; - member_sync_set_online(inst, key, online); - if (g_node_updated_cb) g_node_updated_cb(key, online); - } - return 0; -} - static const struct merkle_sync_data_ops g_member_ops = { .update_bucket_hash = _member_update_bucket_hash, .get_items = _member_get_items, .apply_items = _member_apply_items, - .apply_update = _ms_apply_update, + .apply_update = NULL, /* online-статус больше не рассылается через merkle */ }; /* ── Public API ── */ @@ -935,23 +921,6 @@ int member_sync_count(struct UTUN_INSTANCE* inst, const char* ch_id) { return c; } -void member_sync_set_node_updated_cb(member_sync_node_updated_fn cb) { - g_node_updated_cb = cb; -} - -int member_sync_set_online(struct UTUN_INSTANCE* inst, uint64_t node_id, int online) { - /* ВАЖНО: онлайн-статус НЕ участвует в merkle sync — не хешируется и не - * версионируется. Только локальная запись в nodes.online. Рассылка online - * идёт отдельным лёгким push_update (MSG_ITEM_UPDATE), вне дерева. */ - if (!inst) return 0; - DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: set_online nid=%016llx online=%d", MS_ID, (unsigned long long)node_id, online); - sqlite3* db = _db(inst); if (!db) return 0; - int val = online ? 1 : 0; - if (topo_node_sqlite_node_get_online(db, node_id) == val) return 0; - topo_node_sqlite_node_set_online(db, node_id, val); - return 1; -} - const uint8_t* member_sync_get_hash(struct UTUN_INSTANCE* inst, const char* ch_id, uint8_t level, uint64_t prefix64) { return merkle_sync_get_hash(inst, ch_id, level, prefix64); diff --git a/src/chat/member_sync.h b/src/chat/member_sync.h index c99a350d..8fdf300c 100644 --- a/src/chat/member_sync.h +++ b/src/chat/member_sync.h @@ -25,11 +25,6 @@ * adm_tags, adm_tags_sig, storage); * // дерево автоматически пересчитано * - * // онлайн-статус (из cs_on_conn_up / cs_on_conn_down) - * // только локальная запись; возвращает 1 если изменилось — тогда рассылка: - * if (member_sync_set_online(inst, node_id, 1)) - * merkle_sync_push_update(inst, ch_id, node_id, 0x01, (uint8_t[]){1}, 1); - * * // отмена синхронизации (коллбэк НЕ вызывается) * member_sync_cancel(inst, peer, ch_id); * @@ -191,30 +186,10 @@ int member_sync_count(struct UTUN_INSTANCE* inst, const char* ch_id); Вызывается после локального изменения данных (adm_tags, имя, адреса). */ void member_sync_broadcast_one(struct UTUN_INSTANCE* inst, const char* ch_id, uint64_t member_id); -/* - * Установить онлайн-статус узла (nodes.online = 0/1). - * - * ВАЖНО: онлайн-статус НЕ участвует в merkle sync — он не хешируется и не - * версионируется (в _compute_member_hash поля online нет). Это только локальная - * запись в БД. Рассылка online идёт отдельным лёгким push_update (MSG_ITEM_UPDATE), - * вне дерева, и только если значение реально изменилось (возвращает 1). - * - * Возвращает 1 если значение изменилось, 0 если совпало. - */ -int member_sync_set_online(struct UTUN_INSTANCE* inst, uint64_t node_id, int online); - /* Получить хеш бакета из merkle_tree_hash (для тестов). */ const uint8_t* member_sync_get_hash(struct UTUN_INSTANCE* inst, const char* ch_id, uint8_t level, uint64_t prefix64); -/* - * Коллбэк: узел изменился по инициативе удалённого пира (через MSG_ITEM_UPDATE). - * Вызывается из uasync-потока. Потребитель (chat_sync) может из него - * постить GUI-события. - */ -typedef void (*member_sync_node_updated_fn)(uint64_t node_id, int online); -void member_sync_set_node_updated_cb(member_sync_node_updated_fn cb); - /* * Коллбэк: изменились adm_tags любого узла (включая себя). * channel_id — канал (namespace), в котором произошло изменение. diff --git a/src/chat/merkle_sync.h b/src/chat/merkle_sync.h index a63524be..6421adb4 100644 --- a/src/chat/merkle_sync.h +++ b/src/chat/merkle_sync.h @@ -45,7 +45,6 @@ struct UTUN_INSTANCE; * * // ── изменил данные — дерево само пересчиталось ── * member_sync_put(inst, ch_id, node_id, x25519, ed25519, join_sig, addrs, ac); - * member_sync_set_online(inst, node_id, 1); * * // ── отменил — коллбэк не вызовется ── * member_sync_cancel(inst, peer, ch_id); diff --git a/src/routing_layer/topo_node.c b/src/routing_layer/topo_node.c index b396807f..a5777ed1 100644 --- a/src/routing_layer/topo_node.c +++ b/src/routing_layer/topo_node.c @@ -724,63 +724,23 @@ int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GR // Для chat групп подсети не передаём if (group->group_type == TOPO_GROUP_TYPE_CHAT) { vc = 0; vc6 = 0; } - int sock_count = 0, addr_count = 0, sock6_count = 0, addr6_count = 0; - int tcp4_count = 0, tcp6_count = 0; - struct ETCP_SOCKET* e_sock = instance->etcp_sockets; - while (e_sock) { - if (e_sock->is_tcp) { e_sock = e_sock->next; continue; } - if (e_sock->type == CFG_SERVER_TYPE_PRIVATE || e_sock->type == CFG_SERVER_TYPE_LOCAL) { e_sock = e_sock->next; continue; } - if (e_sock->local_addr.ss_family == AF_INET) { sock_count++; addr_count++; } - else if (e_sock->local_addr.ss_family == AF_INET6) { sock6_count++; addr6_count++; } - e_sock = e_sock->next; - } - { struct ETCP_SOCKET* ts = instance->etcp_sockets; - while (ts) { if (!ts->is_tcp) { ts = ts->next; continue; } - if (ts->type == CFG_SERVER_TYPE_PRIVATE || ts->type == CFG_SERVER_TYPE_LOCAL) { ts = ts->next; continue; } - if (ts->interface_addr.ss_family == AF_INET || ts->local_addr.ss_family == AF_INET) { tcp4_count++; sock_count++; addr_count++; } - else if (ts->interface_addr.ss_family == AF_INET6 || ts->local_addr.ss_family == AF_INET6) { tcp6_count++; sock6_count++; addr6_count++; } - ts = ts->next; } - } - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "counts: udp_v4s=%d udp_v4a=%d udp_v6s=%d udp_v6a=%d tcp4=%d tcp6=%d v4a_total=%d v6a_total=%d", - sock_count, addr_count - tcp4_count, sock6_count, addr6_count - tcp6_count, - tcp4_count, tcp6_count, addr_count, addr6_count); - int changed = 1; uint8_t old_ver = 0; if (group->local_node) { struct TOPO_NODE* oni = topo_node_registry_find(instance->topo_groups, group->local_node->node_id); if (oni) { old_ver = oni->ver; - changed = (vc != (group->local_node->subnets ? topo_list_count((struct _topo_head*)group->local_node->subnets->v4_subnets) : 0)) - || (sock_count != topo_list_count((struct _topo_head*)oni->v4_sock_meta)) - || (addr_count != topo_list_count((struct _topo_head*)oni->v4_addrs)) - || (name_len != (oni->node_name ? strlen(oni->node_name) : 0)) - || (memcmp(oni->public_key, instance->my_keys.public_key, SC_PUBKEY_SIZE) != 0) - || (vc6 != (group->local_node->subnets ? topo_list_count((struct _topo_head*)group->local_node->subnets->v6_subnets) : 0)) - || (sock6_count != topo_list_count((struct _topo_head*)oni->v6_sock_meta)) - || (addr6_count != topo_list_count((struct _topo_head*)oni->v6_addrs)) - || (oni->client_type != instance->client_type) - || (oni->client_activity != instance->client_activity); - if (!changed && oni->v4_sock_meta) { - struct ETCP_SOCKET* es = instance->etcp_sockets; - struct TOPO_SOCKMETA4* sm = oni->v4_sock_meta; - while (es && sm) { - if (es->type == CFG_SERVER_TYPE_PRIVATE || es->type == CFG_SERVER_TYPE_LOCAL) { es = es->next; continue; } - if (es->local_addr.ss_family != AF_INET) { es = es->next; continue; } - if (sm->nat_type != es->nat_type) { changed = 1; break; } - sm = sm->next; es = es->next; - } - } - if (group->local_node && group->instance && group->instance->topo_groups) - if (memcmp(oni->ed25519_public_key, group->ed25519_public_key, SC_PUBKEY_SIZE) != 0) changed = 1; + changed = (vc != (group->local_node->subnets ? topo_list_count((struct _topo_head*)group->local_node->subnets->v4_subnets) : 0)) + || (name_len != (oni->node_name ? strlen(oni->node_name) : 0)) + || (memcmp(oni->public_key, instance->my_keys.public_key, SC_PUBKEY_SIZE) != 0) + || (vc6 != (group->local_node->subnets ? topo_list_count((struct _topo_head*)group->local_node->subnets->v6_subnets) : 0)) + || (oni->client_type != instance->client_type) + || (oni->client_activity != instance->client_activity) + || (memcmp(oni->ed25519_public_key, group->ed25519_public_key, SC_PUBKEY_SIZE) != 0); } } - if (!changed && group->local_node) { - struct TOPO_NODE* oni2 = topo_node_registry_find(instance->topo_groups, group->local_node->node_id); - int old_v4a = oni2 ? topo_list_count((struct _topo_head*)oni2->v4_addrs) : -1; - int old_v6a = oni2 ? topo_list_count((struct _topo_head*)oni2->v6_addrs) : -1; - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "nodeinfo unchanged: old_v4a=%d new_v4a=%d old_v6a=%d new_v6a=%d tcp4=%d tcp6=%d", - old_v4a, addr_count, old_v6a, addr6_count, tcp4_count, tcp6_count); - } + if (!changed && group->local_node) + DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "my nodeinfo unchanged (identity), ver=%d grp=%016llx", + group->local_node->last_ver, (unsigned long long)group->group_id); if (changed) { if (group->local_node) { @@ -810,70 +770,7 @@ int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GR lq->node_id = saved_node_id; lq->last_ver = saved_ver; } - e_sock = instance->etcp_sockets; - while (e_sock) { - if (e_sock->is_tcp) { e_sock = e_sock->next; continue; } - if (e_sock->type == CFG_SERVER_TYPE_PRIVATE || e_sock->type == CFG_SERVER_TYPE_LOCAL) { e_sock = e_sock->next; continue; } - if (e_sock->local_addr.ss_family == AF_INET) { - { struct TOPO_SOCKMETA4* sm = memory_pool_alloc(group->instance->topo_groups->v4_sock_meta_pool); - sm->id = e_sock->sock_id; sm->config_type = e_sock->type; sm->nat_type = e_sock->nat_type; - sm->next = ni->v4_sock_meta; ni->v4_sock_meta = sm; } - - struct sockaddr_in* local_sin = (struct sockaddr_in*)&e_sock->local_addr; - struct sockaddr_in* if_sin = (struct sockaddr_in*)&e_sock->interface_addr; - struct sockaddr_in* nat_sin = (struct sockaddr_in*)&e_sock->nat_addr; - int nat_verified = (e_sock->nat_addr.ss_family == AF_INET && nat_sin->sin_addr.s_addr != 0 && e_sock->nat_type != NAT_VERIFIED_STRICT); - { struct TOPO_ADDR4* a = memory_pool_alloc(group->instance->topo_groups->v4_addr_pool); - if (nat_verified) { - memcpy(a->addr, &nat_sin->sin_addr.s_addr, 4); a->port = ntohs(nat_sin->sin_port); - a->type = TOPO_ADDR_NAT; a->socket_id = e_sock->sock_id; - } else { - int use_local = (local_sin->sin_addr.s_addr != 0); - memcpy(a->addr, use_local ? &local_sin->sin_addr.s_addr : &if_sin->sin_addr.s_addr, 4); - a->port = ntohs(use_local ? local_sin->sin_port : if_sin->sin_port); - a->type = TOPO_ADDR_INTERFACE; a->socket_id = e_sock->sock_id; - } - a->protocol = TOPO_PROTO_UDP; - a->next = ni->v4_addrs; ni->v4_addrs = a; } - } else if (e_sock->local_addr.ss_family == AF_INET6) { - { struct TOPO_SOCKMETA6* sm6 = memory_pool_alloc(group->instance->topo_groups->v6_sock_meta_pool); - sm6->id = e_sock->sock_id; sm6->config_type = e_sock->type; sm6->nat_type = e_sock->nat_type; - sm6->next = ni->v6_sock_meta; ni->v6_sock_meta = sm6; } - struct sockaddr_in6* if_sin6 = (struct sockaddr_in6*)&e_sock->interface_addr; - { struct TOPO_ADDR6* a6 = memory_pool_alloc(group->instance->topo_groups->v6_addr_pool); - memcpy(a6->addr, &if_sin6->sin6_addr, 16); a6->port = ntohs(if_sin6->sin6_port); - a6->type = TOPO_ADDR_INTERFACE; a6->socket_id = e_sock->sock_id; a6->protocol = TOPO_PROTO_UDP; - a6->next = ni->v6_addrs; ni->v6_addrs = a6; } - } - e_sock = e_sock->next; - } - { struct ETCP_SOCKET* ts = instance->etcp_sockets; - while (ts) { if (!ts->is_tcp) { ts = ts->next; continue; } - if (ts->type == CFG_SERVER_TYPE_PRIVATE || ts->type == CFG_SERVER_TYPE_LOCAL) { ts = ts->next; continue; } - struct sockaddr_storage* addr = ts->interface_addr.ss_family ? &ts->interface_addr : &ts->local_addr; - if (!addr || !addr->ss_family) { ts = ts->next; continue; } - if (addr->ss_family == AF_INET) { - struct sockaddr_in* sin = (struct sockaddr_in*)addr; - { struct TOPO_SOCKMETA4* sm = memory_pool_alloc(group->instance->topo_groups->v4_sock_meta_pool); - sm->id = ts->sock_id; sm->config_type = ts->type; sm->nat_type = ts->nat_type; - sm->next = ni->v4_sock_meta; ni->v4_sock_meta = sm; } - { struct TOPO_ADDR4* a = memory_pool_alloc(group->instance->topo_groups->v4_addr_pool); - memcpy(a->addr, &sin->sin_addr.s_addr, 4); a->port = ntohs(sin->sin_port); - a->type = TOPO_ADDR_INTERFACE; a->socket_id = ts->sock_id; a->protocol = TOPO_PROTO_TCP; - a->next = ni->v4_addrs; ni->v4_addrs = a; } - } else if (addr->ss_family == AF_INET6) { - struct sockaddr_in6* sin6 = (struct sockaddr_in6*)addr; - { struct TOPO_SOCKMETA6* sm6 = memory_pool_alloc(group->instance->topo_groups->v6_sock_meta_pool); - sm6->id = ts->sock_id; sm6->config_type = ts->type; sm6->nat_type = ts->nat_type; - sm6->next = ni->v6_sock_meta; ni->v6_sock_meta = sm6; } - { struct TOPO_ADDR6* a6 = memory_pool_alloc(group->instance->topo_groups->v6_addr_pool); - memcpy(a6->addr, &sin6->sin6_addr, 16); a6->port = ntohs(sin6->sin6_port); - a6->type = TOPO_ADDR_INTERFACE; a6->socket_id = ts->sock_id; a6->protocol = TOPO_PROTO_TCP; - a6->next = ni->v6_addrs; ni->v6_addrs = a6; } - } - ts = ts->next; } - } - + /* адреса и sock_meta владеет topo_node_update_my_addresses (реестр + БД + подпись + рассылка) */ topo_node_sign_self(instance, ni); if (group->group_type != TOPO_GROUP_TYPE_CHAT && (vc || vc6)) { @@ -895,11 +792,11 @@ int topo_group_update_my_nodeinfo(struct UTUN_INSTANCE* instance, struct TOPO_GR } if (group->group_type == TOPO_GROUP_TYPE_CHAT) - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "my_nodeinfo updated: v4s=%d v4a=%d v6s=%d v6a=%d tcp4=%d tcp6=%d v4sub=%d v6sub=%d ver=%d grp=%016llx", - sock_count, addr_count, sock6_count, addr6_count, tcp4_count, tcp6_count, vc, vc6, ni->ver, (unsigned long long)group->group_id); + DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "my_nodeinfo updated: v4sub=%d v6sub=%d ver=%d grp=%016llx", + vc, vc6, ni->ver, (unsigned long long)group->group_id); else - DEBUG_INFO(DEBUG_CATEGORY_BGP, "my_nodeinfo updated: v4s=%d v4a=%d v6s=%d v6a=%d tcp4=%d tcp6=%d v4sub=%d v6sub=%d ver=%d grp=%016llx", - sock_count, addr_count, sock6_count, addr6_count, tcp4_count, tcp6_count, vc, vc6, ni->ver, (unsigned long long)group->group_id); + DEBUG_INFO(DEBUG_CATEGORY_BGP, "my_nodeinfo updated: v4sub=%d v6sub=%d ver=%d grp=%016llx", + vc, vc6, ni->ver, (unsigned long long)group->group_id); if (group->group_type != TOPO_GROUP_TYPE_CHAT && instance->rt && !route_insert(instance->rt, lq)) DEBUG_WARN(DEBUG_CATEGORY_ROUTING, "failed to insert local routes"); } else { @@ -985,58 +882,59 @@ int topo_node_update_my_addresses(struct UTUN_INSTANCE* instance) { ts = ts->next; } } - /* Без изменений — освобождаем временные списки и выходим. */ - if (sockmeta4_list_equal(ni->v4_sock_meta, m4_head) && addr4_list_equal(ni->v4_addrs, v4_head) && - sockmeta6_list_equal(ni->v6_sock_meta, m6_head) && addr6_list_equal(ni->v6_addrs, v6_head)) { + /* Без изменений — подпись/рассылку гасим, но БД всё равно синхронизируем: + она могла быть не готова (или пуста) при предыдущем вызове. */ + int changed = !(sockmeta4_list_equal(ni->v4_sock_meta, m4_head) && addr4_list_equal(ni->v4_addrs, v4_head) && + sockmeta6_list_equal(ni->v6_sock_meta, m6_head) && addr6_list_equal(ni->v6_addrs, v6_head)); + if (changed) { + free_v4_sock_list(instance->topo_groups->v4_sock_meta_pool, ni->v4_sock_meta); ni->v4_sock_meta = m4_head; + free_v4_addr_list(instance->topo_groups->v4_addr_pool, ni->v4_addrs); ni->v4_addrs = v4_head; + free_v6_sock_list(instance->topo_groups->v6_sock_meta_pool, ni->v6_sock_meta); ni->v6_sock_meta = m6_head; + free_v6_addr_list(instance->topo_groups->v6_addr_pool, ni->v6_addrs); ni->v6_addrs = v6_head; + ni->ver = (ni->ver % 255) + 1; + topo_node_sign_self(instance, ni); + DEBUG_INFO(DEBUG_CATEGORY_BGP, "my addresses updated, new ver=%d", ni->ver); + } else { free_v4_sock_list(instance->topo_groups->v4_sock_meta_pool, m4_head); free_v4_addr_list(instance->topo_groups->v4_addr_pool, v4_head); free_v6_sock_list(instance->topo_groups->v6_sock_meta_pool, m6_head); free_v6_addr_list(instance->topo_groups->v6_addr_pool, v6_head); - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "my addresses unchanged (ver=%d), skip update", ni->ver); - return 0; + DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "my addresses unchanged (ver=%d), persist DB", ni->ver); } - free_v4_sock_list(instance->topo_groups->v4_sock_meta_pool, ni->v4_sock_meta); ni->v4_sock_meta = m4_head; - free_v4_addr_list(instance->topo_groups->v4_addr_pool, ni->v4_addrs); ni->v4_addrs = v4_head; - free_v6_sock_list(instance->topo_groups->v6_sock_meta_pool, ni->v6_sock_meta); ni->v6_sock_meta = m6_head; - free_v6_addr_list(instance->topo_groups->v6_addr_pool, ni->v6_addrs); ni->v6_addrs = v6_head; - - ni->ver = (ni->ver % 255) + 1; - DEBUG_INFO(DEBUG_CATEGORY_BGP, "my addresses updated, new ver=%d", ni->ver); - - topo_node_sign_self(instance, ni); - + /* всегда: адреса (addrs_put) + классификация node_type/storage (nodeinfo_updated) */ if (instance->topo_sqlite_db) { time_t now_sec = ntp_time_get_seconds(instance); topo_node_sqlite_node_update_verified(instance->topo_sqlite_db, instance->node_id, instance->name, instance->my_keys.public_key, instance->my_ed25519_pubkey, (uint64_t)now_sec, now_sec); topo_node_sqlite_addrs_put(instance->topo_sqlite_db, instance->node_id, ni); + topo_node_sqlite_nodeinfo_updated(instance->topo_sqlite_db, instance->node_id); } - struct ll_entry* ge = instance->topo_groups->group_list->head; - while (ge) { - struct TOPO_GROUP* g = (struct TOPO_GROUP*)ge; - if (g->local_node) { - if (g->senders_list) { - struct ll_entry* se = g->senders_list->head; - while (se) { - struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)se->data; - if (item && item->conn) topo_group_send_nodeinfo(g, g->local_node, item->conn, 0); - se = se->next; + if (changed) { + struct ll_entry* ge = instance->topo_groups->group_list->head; + while (ge) { + struct TOPO_GROUP* g = (struct TOPO_GROUP*)ge; + if (g->local_node) { + if (g->senders_list) { + struct ll_entry* se = g->senders_list->head; + while (se) { + struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)se->data; + if (item && item->conn) topo_group_send_nodeinfo(g, g->local_node, item->conn, 0); + se = se->next; + } } } + ge = ge->next; } - ge = ge->next; } - return 1; + return changed ? 1 : 0; } void topo_node_on_socket_changed(struct ETCP_SOCKET* sock, int event, void* arg) { (void)arg; if (!sock || !sock->instance) return; DEBUG_INFO(DEBUG_CATEGORY_BGP, "socket %s changed (event=0x%x), updating nodeinfo", sock->name, event); - int changed = topo_node_update_my_addresses(sock->instance); - if (changed > 0 && sock->instance->topo_sqlite_db) - topo_node_sqlite_nodeinfo_updated(sock->instance->topo_sqlite_db, sock->instance->node_id); + topo_node_update_my_addresses(sock->instance); } diff --git a/src/routing_layer/topo_node_sqlite.c b/src/routing_layer/topo_node_sqlite.c index e0de6657..b2f1c275 100644 --- a/src/routing_layer/topo_node_sqlite.c +++ b/src/routing_layer/topo_node_sqlite.c @@ -84,10 +84,18 @@ int topo_node_sqlite_addrs_put(sqlite3* db, uint64_t node_id, struct TOPO_NODE* sqlite3_exec(db, "BEGIN", NULL, NULL, NULL); sqlite3_stmt* ds = NULL; + int deleted = 0; if (sqlite3_prepare_v2(db, "DELETE FROM node_addresses WHERE node_id=?", -1, &ds, NULL) == SQLITE_OK) { sqlite3_bind_int64(ds, 1, (sqlite3_int64)node_id); - sqlite3_step(ds); + int drc = sqlite3_step(ds); + deleted = sqlite3_changes(db); + if (drc != SQLITE_DONE) + DEBUG_WARN(DEBUG_CATEGORY_BGP, "addrs_put: DELETE rc=%d (not DONE) node=%016llx: %s", + drc, (unsigned long long)node_id, sqlite3_errmsg(db)); sqlite3_finalize(ds); + } else { + DEBUG_WARN(DEBUG_CATEGORY_BGP, "addrs_put: DELETE prepare failed node=%016llx: %s", + (unsigned long long)node_id, sqlite3_errmsg(db)); } sqlite3_stmt* is = NULL; @@ -113,8 +121,9 @@ int topo_node_sqlite_addrs_put(sqlite3* db, uint64_t node_id, struct TOPO_NODE* sqlite3_bind_int(is, 8, (int)nat); sqlite3_bind_int(is, 9, (int)a->socket_id); if (sqlite3_step(is) != SQLITE_DONE) - DEBUG_ERROR(DEBUG_CATEGORY_BGP, "addrs_put: INSERT failed node=%016llx fam=4 sock=%d: %s", - (unsigned long long)node_id, a->socket_id, sqlite3_errmsg(db)); + DEBUG_ERROR(DEBUG_CATEGORY_BGP, "addrs_put: INSERT failed node=%016llx fam=4 sock=%d at=%d cfg=%d nat=%d addr=%d.%d.%d.%d:%d: %s", + (unsigned long long)node_id, a->socket_id, at, cfg, nat, + a->addr[0], a->addr[1], a->addr[2], a->addr[3], a->port, sqlite3_errmsg(db)); sqlite3_reset(is); written++; } @@ -132,16 +141,16 @@ int topo_node_sqlite_addrs_put(sqlite3* db, uint64_t node_id, struct TOPO_NODE* sqlite3_bind_int(is, 8, (int)nat); sqlite3_bind_int(is, 9, (int)a->socket_id); if (sqlite3_step(is) != SQLITE_DONE) - DEBUG_ERROR(DEBUG_CATEGORY_BGP, "addrs_put: INSERT failed node=%016llx fam=6 sock=%d: %s", - (unsigned long long)node_id, a->socket_id, sqlite3_errmsg(db)); + DEBUG_ERROR(DEBUG_CATEGORY_BGP, "addrs_put: INSERT failed node=%016llx fam=6 sock=%d at=%d cfg=%d nat=%d: %s", + (unsigned long long)node_id, a->socket_id, at, cfg, nat, sqlite3_errmsg(db)); sqlite3_reset(is); written++; } sqlite3_finalize(is); sqlite3_exec(db, "COMMIT", NULL, NULL, NULL); - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "addrs_put: node=%016llx v4=%d v6=%d total=%d", - (unsigned long long)node_id, + DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "addrs_put: node=%016llx deleted=%d v4=%d v6=%d written=%d", + (unsigned long long)node_id, deleted, topo_list_count((struct _topo_head*)ni->v4_addrs), topo_list_count((struct _topo_head*)ni->v6_addrs), written); return 0; diff --git a/src/transport_layer/etcp.c b/src/transport_layer/etcp.c index 4af19ca5..89a4f090 100644 --- a/src/transport_layer/etcp.c +++ b/src/transport_layer/etcp.c @@ -241,6 +241,7 @@ struct ETCP_CONN* etcp_connection_create(struct UTUN_INSTANCE* instance, char* n // etcp->window_size = MAX_INFLIGHT_BYTES; // Not used etcp->mtu = ETCP_RFC791_MIN_MTU; // Default MTU per RFC 791 etcp->next_tx_id = 1; + random_bytes((uint8_t*)&etcp->reset_id, 8); // начальная случайная эпоха ресета etcp->rtt_avg_10 = 10; // Initial guess (1ms) etcp->rtt_history_idx = 0; memset(etcp->rtt_history, 0, sizeof(etcp->rtt_history)); @@ -592,12 +593,24 @@ void etcp_links_reset(struct ETCP_CONN* etcp) {// Если сбой в обме } } -void etcp_conn_reinit(struct ETCP_CONN* etcp, const char* reason) {// Если сбой в обмене или ребутнулась одна из сторон -> необходимо заново переинициализировать соединение +void etcp_conn_reinit(struct ETCP_CONN* etcp, const char* reason) { + etcp_conn_reinit_id(etcp, reason, etcp ? etcp->reset_id : 0); +} + +// Фатальный ресет (normalizer desync / bad fragment size): меняем эпоху reset_id +void etcp_conn_fatal_reinit(struct ETCP_CONN* etcp, const char* reason) { + uint64_t new_id; random_bytes((uint8_t*)&new_id, 8); + etcp_conn_reinit_id(etcp, reason, new_id); +} + +void etcp_conn_reinit_id(struct ETCP_CONN* etcp, const char* reason, uint64_t reset_id) {// Если сбой в обмене или ребутнулась одна из сторон -> необходимо заново переинициализировать соединение if (!etcp) return; if (etcp->state == 2) { DEBUG_WARN(DEBUG_CATEGORY_ETCP, "[%s] conn_reinit on deleted conn", etcp->log_name); return; } - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] REINIT: %s (init=%d links=%d reinit=%u tx=%d)", - etcp->log_name, reason, etcp->initialized, etcp->links_up, etcp->reinit_count, etcp->tx_state); + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] REINIT: %s (init=%d links=%d reinit=%u tx=%d) rid=%016llx", + etcp->log_name, reason, etcp->initialized, etcp->links_up, etcp->reinit_count, etcp->tx_state, + (unsigned long long)reset_id); + etcp->reset_id = reset_id; etcp->setup_start_tb = get_time_tb(); etcp->reinit_count++; etcp->reinit_pending = 1; @@ -626,6 +639,32 @@ void etcp_conn_reinit(struct ETCP_CONN* etcp, const char* reason) {// Если DEBUG_TRACE(DEBUG_CATEGORY_ETCP, "end"); } +// Применение чужого reset_id (эпохи ресета пира). +// Если id отличается от уже виденного — пир сделал новый ресет. +// slave (больший node_id) принимает эпоху пира и ресетится; master сохраняет свою. +void etcp_conn_apply_peer_reset_id(struct ETCP_CONN* conn, uint64_t peer_id, uint64_t peer_reset_id) { + if (!conn) return; + int first_time = (conn->peer_reset_id == 0); + if (peer_reset_id == conn->peer_reset_id) return; // ретрансмиссия известного ресета + conn->peer_reset_id = peer_reset_id; + if (peer_reset_id == conn->reset_id) return; // уже в одной эпохе + int i_am_master = conn->instance && conn->instance->node_id < peer_id; + if (first_time) { + // первый обмен: принимаем id как есть БЕЗ ресета (мы уже свежие при init) + if (!i_am_master) conn->reset_id = peer_reset_id; // slave принимает id мастера + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "[%s] first peer reset id=%016llx (no reset)", conn->log_name, (unsigned long long)peer_reset_id); + return; + } + if (i_am_master) { + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] peer reset id=%016llx != mine=%016llx — I'm master, keeping my id", + conn->log_name, (unsigned long long)peer_reset_id, (unsigned long long)conn->reset_id); + return; + } + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] peer reset id=%016llx != mine=%016llx — slave adopts and re-resets", + conn->log_name, (unsigned long long)peer_reset_id, (unsigned long long)conn->reset_id); + etcp_conn_reinit_id(conn, "peer reset epoch", peer_reset_id); +} + // внутренняя функция. Вызывается один раз когда первый линк готов. void etcp_conn_ready(struct ETCP_CONN* conn) { if (!conn) return; diff --git a/src/transport_layer/etcp.h b/src/transport_layer/etcp.h index 35ab45ee..93de87fe 100644 --- a/src/transport_layer/etcp.h +++ b/src/transport_layer/etcp.h @@ -228,6 +228,8 @@ struct ETCP_CONN { uint8_t got_initial_pkt; // uint8_t initialized; // 0 - только созданный ETCP, 1 - хотя бы один линк проинициалзирован (обмен ключами произведен) uint32_t session_id; // случайный ID сессии (генерируется клиентом) для защиты от ложного reinit + uint64_t reset_id; // метка эпохи ресета (генерируется при локальном ресете, master авторитет) + uint64_t peer_reset_id; // последняя увиденная эпоха ресета пира uint8_t tx_state; // 0 - n/a, 1 - data_wait (queues empty), 2 - link_wait (link busy) uint8_t links_up; // 0 - канал не готов для передачи, 1 - канал готов для передачи (хотя бы один линк не down) uint8_t reset_done; // 0 - рукопожатие не завершено (реинит разрешён), 1 - соединение стабильно (реинит заблокирован) @@ -289,6 +291,12 @@ void etcp_conn_reset(struct ETCP_CONN* etcp); void etcp_links_reset(struct ETCP_CONN* etcp); void etcp_conn_reinit(struct ETCP_CONN* etcp, const char* reason); +// Реинит с конкретным reset_id (для adopt чужой эпохи) — без генерации нового id +void etcp_conn_reinit_id(struct ETCP_CONN* etcp, const char* reason, uint64_t reset_id); +// Фатальный ресет (normalizer desync / bad fragment size): генерит новую эпоху reset_id +void etcp_conn_fatal_reinit(struct ETCP_CONN* etcp, const char* reason); +// Применение чужого reset_id: при отличии — slave принимает эпоху пира и ресетится, master пересылает свою +void etcp_conn_apply_peer_reset_id(struct ETCP_CONN* conn, uint64_t peer_id, uint64_t peer_reset_id); void etcp_connection_ready(struct ETCP_CONN* etcp);// вызывается когда подключение инициализировано diff --git a/src/transport_layer/etcp_connections.c b/src/transport_layer/etcp_connections.c index 70f6c8b2..0e41fcbf 100644 --- a/src/transport_layer/etcp_connections.c +++ b/src/transport_layer/etcp_connections.c @@ -161,6 +161,7 @@ static void etcp_link_send_init(struct ETCP_LINK* link, uint8_t reset, uint8_t c req->code = reset ? ETCP_INIT_REQUEST : ETCP_INIT_REQUEST_NOINIT; *(uint64_t*)req->node_id = htobe64(link->etcp->instance->node_id); *(uint32_t*)req->session_id = htobe32(link->etcp->session_id); + *(uint64_t*)req->reset_id = htobe64(link->etcp->reset_id); *(uint16_t*)req->mtu = htobe16(link->mtu_local); *(uint16_t*)req->keepalive = htobe16(link->etcp->instance->keepalive_interval); *(uint16_t*)req->recovery = htobe16(link->recovery_interval / 100); @@ -1026,7 +1027,8 @@ void etcp_link_enter_ready_tcp(struct ETCP_LINK *link) { if (link->tcp_link) { uint8_t peer_gop = stcp_link_get_peer_got_initial_pkt(link->tcp_link); uint32_t peer_sid = stcp_link_get_peer_session_id(link->tcp_link); - if (etcp->got_initial_pkt == 1 && peer_gop == 0) { + uint64_t peer_rid = stcp_link_get_peer_reset_id(link->tcp_link); + if (etcp->got_initial_pkt == 1 && peer_gop == 0 && peer_rid != etcp->reset_id) { DEBUG_INFO(DEBUG_CATEGORY_ETCP, "[%s] TCP ready: I'm dirty (gop=1) but peer clean (gop=0) → reinit", etcp->log_name); etcp_conn_reinit(etcp, "tcp peer clean"); @@ -1037,6 +1039,7 @@ void etcp_link_enter_ready_tcp(struct ETCP_LINK *link) { etcp_conn_reinit(etcp, "tcp session changed"); } etcp->session_id = peer_sid; + etcp_conn_apply_peer_reset_id(etcp, etcp->peer_node_id, peer_rid); link->peer_device_type = stcp_link_get_peer_device_type(link->tcp_link); { uint16_t peer_ka = stcp_link_get_peer_keepalive_interval(link->tcp_link); @@ -1570,6 +1573,7 @@ static void send_init_response(struct ETCP_SOCKET* e_sock, struct ETCP_DGRAM* pk } *(uint64_t*)resp->node_id = htobe64(e_sock->instance->node_id); *(uint32_t*)resp->session_id = htobe32(conn->session_id); + *(uint64_t*)resp->reset_id = htobe64(conn->reset_id); resp->mtu[0]=link->mtu_local>>8; resp->mtu[1]=link->mtu_local; @@ -1704,7 +1708,8 @@ static int handle_init_response_client(struct ETCP_SOCKET* e_sock, struct ETCP_D struct ETCP_INIT_RESPONSE_PKT* resp = (struct ETCP_INIT_RESPONSE_PKT*)pkt->data; uint64_t server_node_id = be64toh(*(uint64_t*)resp->node_id); uint32_t resp_session_id = be32toh(*(uint32_t*)resp->session_id); - DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "INIT_RESPONSE session_id=%08x", resp_session_id); + uint64_t resp_reset_id = be64toh(*(uint64_t*)resp->reset_id); + DEBUG_TRACE(DEBUG_CATEGORY_CONNECTION, "INIT_RESPONSE session_id=%08x rid=%016llx", resp_session_id, (unsigned long long)resp_reset_id); // Check session_id: ignore response if it doesn't match our session if (resp_session_id != link->etcp->session_id) { DEBUG_WARN(DEBUG_CATEGORY_CONNECTION, "[%s] INIT_RESPONSE session_id mismatch: got %08x, expected %08x, ignoring", @@ -1712,6 +1717,7 @@ static int handle_init_response_client(struct ETCP_SOCKET* e_sock, struct ETCP_D memory_pool_free(e_sock->instance->pkt_pool, pkt); return 0; } + etcp_conn_apply_peer_reset_id(link->etcp, server_node_id, resp_reset_id); link->mtu_remote = be16toh(*(uint16_t*)resp->mtu); if (link->mtu_remote > PACKET_DATA_MAX_MTU) link->mtu_remote = PACKET_DATA_MAX_MTU; link->mtu = link->mtu_local < link->mtu_remote ? link->mtu_local : link->mtu_remote; @@ -2014,11 +2020,12 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) { struct ETCP_INIT_REQUEST_PKT* req = (struct ETCP_INIT_REQUEST_PKT*)pkt->data; peer_id = be64toh(*(uint64_t*)req->node_id); uint32_t session_id = be32toh(*(uint32_t*)req->session_id); + uint64_t reset_id = be64toh(*(uint64_t*)req->reset_id); uint16_t mtu = be16toh(*(uint16_t*)req->mtu); uint16_t src_port = be16toh(*(uint16_t*)req->src_port); - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "INIT received: peer=0x%016llx mtu=%u link=%u sock=%u session=%08x src=%s", + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "INIT received: peer=0x%016llx mtu=%u link=%u sock=%u session=%08x rid=%016llx src=%s", (unsigned long long)peer_id, mtu, req->link_id, req->socket_id, - session_id, sockaddr_storage_to_str(&addr).str); + session_id, (unsigned long long)reset_id, sockaddr_storage_to_str(&addr).str); struct ETCP_CONN* conn = NULL; { @@ -2073,6 +2080,12 @@ void etcp_connections_read_callback_socket(socket_t sock, void* arg) { }// коллизия - peer id совпал а ключи разные. } + // Синхронизация эпохи ресета (для нового conn — просто принять эпоху пира) + if (new_conn) { + conn->reset_id = reset_id; + conn->peer_reset_id = reset_id; + } + if (conn->fin_wait) { DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "[%s] received INIT during fin_wait, clearing", conn->log_name); conn->fin_wait = 0; @@ -2237,6 +2250,9 @@ create_new_link: memcpy(conn->peer_ed25519_pubkey, req->ed25519_pubkey, SC_PUBKEY_SIZE); + // Применяем эпоху ресета пира ПОСЛЕ всей reinit-логики (adopt должен быть финальным словом) + if (!new_conn) etcp_conn_apply_peer_reset_id(conn, peer_id, reset_id); + send_init_response(e_sock, pkt, link, conn, req, &addr, send_reset, req_src_ip, req_src_port, pkt_len); return; @@ -2308,7 +2324,9 @@ int etcp_packet_decrypted(struct ETCP_SOCKET* e_sock, struct ETCP_DGRAM* pkt, volatile int _halt = 1; while (_halt) {} } etcp_conn_input(pkt); - } else memory_pool_free(e_sock->instance->pkt_pool, pkt); + } else { + memory_pool_free(e_sock->instance->pkt_pool, pkt); + } return 0; } diff --git a/src/transport_layer/etcp_connections.h b/src/transport_layer/etcp_connections.h index 95d4c9bf..7e10a515 100644 --- a/src/transport_layer/etcp_connections.h +++ b/src/transport_layer/etcp_connections.h @@ -73,20 +73,21 @@ struct ETCP_INIT_REQUEST_PKT { uint8_t code; // 0: ETCP_INIT_REQUEST (0x02) или ETCP_INIT_REQUEST_NOINIT (0x04) uint8_t node_id[8]; // 1: sender node_id (big-endian) uint8_t session_id[4]; // 9: session (big-endian) - uint8_t mtu[2]; // 13: client MTU (big-endian) - uint8_t keepalive[2]; // 15: keepalive interval (big-endian) - uint8_t recovery[2]; // 17: recovery interval/100 (big-endian) - uint8_t link_id; // 19: client's local link id - uint8_t socket_id; // 20: client's socket id - uint8_t only_local; // 21: client only_local flag - uint8_t type; // 22: client socket type (CFG_SERVER_TYPE_*) + uint8_t reset_id[8]; // 13: reset epoch id (big-endian) + uint8_t mtu[2]; // 21: client MTU (big-endian) + uint8_t keepalive[2]; // 23: keepalive interval (big-endian) + uint8_t recovery[2]; // 25: recovery interval/100 (big-endian) + uint8_t link_id; // 27: client's local link id + uint8_t socket_id; // 28: client's socket id + uint8_t only_local; // 29: client only_local flag + uint8_t type; // 30: client socket type (CFG_SERVER_TYPE_*) // V2 fields (NAT_DIRECT detection): - uint8_t src_ipv4[4]; // 23: client interface_addr IPv4 (big-endian, 0 if N/A) - uint8_t src_port[2]; // 27: client interface_addr port (big-endian) - uint8_t collision; // 29: 1 = cross-connect, remote claims master - uint8_t ed25519_pubkey[SC_PUBKEY_SIZE]; // 30: client Ed25519 pubkey (32 bytes) + uint8_t src_ipv4[4]; // 31: client interface_addr IPv4 (big-endian, 0 if N/A) + uint8_t src_port[2]; // 35: client interface_addr port (big-endian) + uint8_t collision; // 37: 1 = cross-connect, remote claims master + uint8_t ed25519_pubkey[SC_PUBKEY_SIZE]; // 38: client Ed25519 pubkey (32 bytes) // V4 fields: - uint8_t device_type; // 62: CLIENT_TYPE_SERVER/DESKTOP/MOBILE + uint8_t device_type; // 70: CLIENT_TYPE_SERVER/DESKTOP/MOBILE } __attribute__((packed)); #define ETCP_INIT_REQ_SIZE sizeof(struct ETCP_INIT_REQUEST_PKT) @@ -95,19 +96,20 @@ struct ETCP_INIT_RESPONSE_PKT { uint8_t code; // 0: ETCP_INIT_RESPONSE (0x03) или ETCP_INIT_RESPONSE_NOINIT (0x05) uint8_t node_id[8]; // 1: server node_id (big-endian) uint8_t session_id[4]; // 9: session (big-endian) - uint8_t mtu[2]; // 13: server MTU (big-endian) - uint8_t link_id; // 15: server's local link id - uint8_t remote_socket_id; // 16: echo of client's socket_id - uint8_t only_local; // 17: server only_local flag - uint8_t type; // 18: server socket type (CFG_SERVER_TYPE_*) + uint8_t reset_id[8]; // 13: reset epoch id (big-endian) + uint8_t mtu[2]; // 21: server MTU (big-endian) + uint8_t link_id; // 23: server's local link id + uint8_t remote_socket_id; // 24: echo of client's socket_id + uint8_t only_local; // 25: server only_local flag + uint8_t type; // 26: server socket type (CFG_SERVER_TYPE_*) // V2 fields: - uint8_t peer_ipv4[4]; // 19: client's external NAT address (big-endian) - uint8_t peer_port[2]; // 23: client's external NAT port (big-endian) + uint8_t peer_ipv4[4]; // 27: client's external NAT address (big-endian) + uint8_t peer_port[2]; // 31: client's external NAT port (big-endian) // V3 fields: - uint8_t ed25519_pubkey[SC_PUBKEY_SIZE]; // 25: server Ed25519 pubkey (32 bytes) + uint8_t ed25519_pubkey[SC_PUBKEY_SIZE]; // 33: server Ed25519 pubkey (32 bytes) // V4 fields: - uint8_t device_type; // 57: CLIENT_TYPE_SERVER/DESKTOP/MOBILE - uint8_t keepalive[2]; // 58: server keepalive interval (big-endian) + uint8_t device_type; // 65: CLIENT_TYPE_SERVER/DESKTOP/MOBILE + uint8_t keepalive[2]; // 66: server keepalive interval (big-endian) } __attribute__((packed)); #define ETCP_INIT_RESP_SIZE sizeof(struct ETCP_INIT_RESPONSE_PKT) diff --git a/src/transport_layer/pkt_normalizer.c b/src/transport_layer/pkt_normalizer.c index d7cf2a29..e83ef5ee 100644 --- a/src/transport_layer/pkt_normalizer.c +++ b/src/transport_layer/pkt_normalizer.c @@ -369,7 +369,7 @@ static void pn_unpacker_cb(struct ll_queue* q, void* arg) { // Incomplete header, reset DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "reset state"); // pn_unpacker_reset_state(pn); - etcp_conn_reinit(pn->etcp, "normalizer desync"); + etcp_conn_fatal_reinit(pn->etcp, "normalizer desync"); break; } uint16_t part_size = payload[ptr] | (payload[ptr + 1] << 8); @@ -379,7 +379,7 @@ static void pn_unpacker_cb(struct ll_queue* q, void* arg) { if (part_size<1 || part_size>16384) { pn->logic_errors++; DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "PART_SIZE ERROR!!! %d", part_size); - etcp_conn_reinit(pn->etcp, "bad fragment size"); + etcp_conn_fatal_reinit(pn->etcp, "bad fragment size"); break; } diff --git a/src/transport_layer/stcp.h b/src/transport_layer/stcp.h index 61f15125..e1409233 100644 --- a/src/transport_layer/stcp.h +++ b/src/transport_layer/stcp.h @@ -26,10 +26,21 @@ typedef void (*stcp_ping_cb)(int success, uint16_t rtt, void *arg); #define STCP_HS_TIMEOUT 50000 // 5s in 0.1ms timebase units #define STCP_CONNECT_TIMEOUT 100000 // 10s in 0.1ms timebase units -#define STCP_HS_CLIENT_MIN (SC_PUBKEY_ENC_SIZE + STCP_HS_ENC_CLIENT) // 40+47=87 -#define STCP_HS_SERVER_MIN (SC_PUBKEY_ENC_SIZE + STCP_HS_ENC_SERVER) // 40+47=87 -#define STCP_HS_ENC_CLIENT 47 // ed25519_pubkey(32) + got_initial_pkt(1) + session_id(4) + padding_size(2) + device_type(1) + keepalive(2) + flags(1) + CRC32(4) -#define STCP_HS_ENC_SERVER 47 // ed25519_pubkey(32) + got_initial_pkt(1) + session_id(4) + padding_size(2) + device_type(1) + keepalive(2) + flags(1) + CRC32(4) +#define STCP_HS_CLIENT_MIN (SC_PUBKEY_ENC_SIZE + STCP_HS_ENC_CLIENT) // 40+55=95 +#define STCP_HS_SERVER_MIN (SC_PUBKEY_ENC_SIZE + STCP_HS_ENC_SERVER) // 40+55=95 +#define STCP_HS_ENC_CLIENT 55 // ed25519_pubkey(32) + got_initial_pkt(1) + session_id(4) + reset_id(8) + padding_size(2) + device_type(1) + keepalive(2) + flags(1) + CRC32(4) +#define STCP_HS_ENC_SERVER 55 // ed25519_pubkey(32) + got_initial_pkt(1) + session_id(4) + reset_id(8) + padding_size(2) + device_type(1) + keepalive(2) + flags(1) + CRC32(4) + +// handshake plain layout (before CRC): +// [0..31] ed25519_pubkey +// [32] got_initial_pkt +// [33..36] session_id +// [37..44] reset_id (epoch) +// [45..46] padding_size +// [47] device_type +// [48..49] keepalive_interval +// [50] flags +#define STCP_HS_PLAIN_SIZE 51 // без CRC32 #define STCP_HANDSHAKE_FLAG_PING 0x01 @@ -82,6 +93,8 @@ struct stcp_conn { uint8_t peer_got_initial_pkt; // remote, received from peer during handshake uint32_t session_id; // local, sent to peer during handshake uint32_t peer_session_id; // remote, received from peer during handshake + uint64_t reset_id; // local, sent to peer during handshake (reset epoch) + uint64_t peer_reset_id; // remote, received from peer during handshake uint8_t device_type; // CLIENT_TYPE_* sent to peer during handshake uint16_t keepalive_interval; // keepalive interval sent to peer during handshake uint8_t peer_device_type; // CLIENT_TYPE_* from peer handshake diff --git a/src/transport_layer/stcp_client.c b/src/transport_layer/stcp_client.c index fe4bb3a6..1220cf87 100644 --- a/src/transport_layer/stcp_client.c +++ b/src/transport_layer/stcp_client.c @@ -60,17 +60,19 @@ static void client_send_handshake(struct stcp_conn *c, const uint8_t *server_pub uint8_t gop = c->etcp_conn ? c->etcp_conn->got_initial_pkt : c->got_initial_pkt; uint32_t sid = c->etcp_conn ? c->etcp_conn->session_id : c->session_id; - uint8_t plain[43]; memcpy(plain, my_ed25519, 32); + uint64_t rid = c->etcp_conn ? c->etcp_conn->reset_id : c->reset_id; + uint8_t plain[STCP_HS_PLAIN_SIZE]; memcpy(plain, my_ed25519, 32); plain[32] = gop; memcpy(plain + 33, &sid, 4); - plain[37] = (uint8_t)padding; plain[38] = (uint8_t)(padding >> 8); - plain[39] = c->device_type; - *(uint16_t*)(plain + 40) = htobe16(c->keepalive_interval); - plain[42] = c->hs_flags; - uint32_t crc = crc32_calc(plain, 43); + { uint64_t rid_be = htobe64(rid); memcpy(plain + 37, &rid_be, 8); } + plain[45] = (uint8_t)padding; plain[46] = (uint8_t)(padding >> 8); + plain[47] = c->device_type; + *(uint16_t*)(plain + 48) = htobe16(c->keepalive_interval); + plain[50] = c->hs_flags; + uint32_t crc = crc32_calc(plain, STCP_HS_PLAIN_SIZE); uint8_t *enc_dst = hs + SC_PUBKEY_ENC_SIZE; - memcpy(enc_dst, plain, 43); - enc_dst[43] = (uint8_t)(crc >> 0); enc_dst[44] = (uint8_t)(crc >> 8); enc_dst[45] = (uint8_t)(crc >> 16); enc_dst[46] = (uint8_t)(crc >> 24); + memcpy(enc_dst, plain, STCP_HS_PLAIN_SIZE); + enc_dst[51] = (uint8_t)(crc >> 0); enc_dst[52] = (uint8_t)(crc >> 8); enc_dst[53] = (uint8_t)(crc >> 16); enc_dst[54] = (uint8_t)(crc >> 24); if (sc_stream_xor(&c->stream_send, enc_dst, STCP_HS_ENC_CLIENT) != SC_OK) { u_free(hs); stcp_conn_do_close(c, 2); return; } for (int i = 0; i < padding; i++) hs[SC_PUBKEY_ENC_SIZE + STCP_HS_ENC_CLIENT + i] = (uint8_t)(salt[0] ^ i); @@ -103,12 +105,13 @@ static void client_hs_cb(struct stcp_conn *c, uint8_t *data, size_t len) { memcpy(c->peer_ed25519_pubkey, enc_hs, SC_PUBKEY_SIZE); c->peer_ed25519_set = 1; c->peer_got_initial_pkt = enc_hs[32]; memcpy(&c->peer_session_id, enc_hs + 33, 4); - uint16_t padding_size = (uint16_t)enc_hs[37] | ((uint16_t)enc_hs[38] << 8); - c->peer_device_type = enc_hs[39]; - c->peer_keepalive_interval = ((uint16_t)enc_hs[40] << 8) | enc_hs[41]; - c->peer_flags = enc_hs[42]; - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "stcp_client: server response OK gop=%d sid=%08x padding=%u dev=%d ka=%u flags=%02x", - c->peer_got_initial_pkt, c->peer_session_id, padding_size, c->peer_device_type, c->peer_keepalive_interval, c->peer_flags); + { uint64_t rid_be; memcpy(&rid_be, enc_hs + 37, 8); c->peer_reset_id = be64toh(rid_be); } + uint16_t padding_size = (uint16_t)enc_hs[45] | ((uint16_t)enc_hs[46] << 8); + c->peer_device_type = enc_hs[47]; + c->peer_keepalive_interval = ((uint16_t)enc_hs[48] << 8) | enc_hs[49]; + c->peer_flags = enc_hs[50]; + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "stcp_client: server response OK gop=%d sid=%08x rid=%016llx padding=%u dev=%d ka=%u flags=%02x", + c->peer_got_initial_pkt, c->peer_session_id, (unsigned long long)c->peer_reset_id, padding_size, c->peer_device_type, c->peer_keepalive_interval, c->peer_flags); stcp_recv_set(c, padding_size, 0, client_hs_padding_cb); } diff --git a/src/transport_layer/stcp_link.c b/src/transport_layer/stcp_link.c index 7bca8b26..fe8eb74d 100644 --- a/src/transport_layer/stcp_link.c +++ b/src/transport_layer/stcp_link.c @@ -48,6 +48,7 @@ struct stcp_link { uint8_t peer_got_initial_pkt; // received from peer during handshake uint32_t peer_session_id; // received from peer during handshake + uint64_t peer_reset_id; // received from peer during handshake (reset epoch) uint8_t peer_device_type; // CLIENT_TYPE_* from peer handshake uint16_t peer_keepalive_interval; // keepalive interval from peer handshake }; @@ -114,6 +115,7 @@ static void server_accept_cb(struct stcp_conn *conn, void *arg) { link->conn = conn; link->peer_got_initial_pkt = conn->peer_got_initial_pkt; link->peer_session_id = conn->peer_session_id; + link->peer_reset_id = conn->peer_reset_id; link->peer_device_type = conn->peer_device_type; link->peer_keepalive_interval = conn->peer_keepalive_interval; @@ -141,6 +143,7 @@ static void client_ready_cb(struct stcp_conn *conn, void *arg) { link->conn = conn; link->peer_got_initial_pkt = conn->peer_got_initial_pkt; link->peer_session_id = conn->peer_session_id; + link->peer_reset_id = conn->peer_reset_id; link->peer_device_type = conn->peer_device_type; link->peer_keepalive_interval = conn->peer_keepalive_interval; @@ -370,6 +373,10 @@ uint32_t stcp_link_get_peer_session_id(struct stcp_link *link) { return link ? link->peer_session_id : 0; } +uint64_t stcp_link_get_peer_reset_id(struct stcp_link *link) { + return link ? link->peer_reset_id : 0; +} + uint8_t stcp_link_get_peer_device_type(struct stcp_link *link) { return link ? link->peer_device_type : 0; } diff --git a/src/transport_layer/stcp_link.h b/src/transport_layer/stcp_link.h index 4b57f3b8..27c0bbf4 100644 --- a/src/transport_layer/stcp_link.h +++ b/src/transport_layer/stcp_link.h @@ -62,6 +62,7 @@ const uint8_t *stcp_link_get_peer_pubkey(struct stcp_link *link); const uint8_t *stcp_link_get_peer_ed25519_pubkey(struct stcp_link *link); uint8_t stcp_link_get_peer_got_initial_pkt(struct stcp_link *link); uint32_t stcp_link_get_peer_session_id(struct stcp_link *link); +uint64_t stcp_link_get_peer_reset_id(struct stcp_link *link); uint8_t stcp_link_get_peer_device_type(struct stcp_link *link); uint16_t stcp_link_get_peer_keepalive_interval(struct stcp_link *link); diff --git a/src/transport_layer/stcp_server.c b/src/transport_layer/stcp_server.c index 7ce59744..9bfa0b9b 100644 --- a/src/transport_layer/stcp_server.c +++ b/src/transport_layer/stcp_server.c @@ -77,12 +77,13 @@ static void server_hs_phase1_cb(struct stcp_conn *c, uint8_t *data, size_t len) memcpy(c->peer_ed25519_pubkey, enc_hs, SC_PUBKEY_SIZE); c->peer_ed25519_set = 1; c->peer_got_initial_pkt = enc_hs[32]; memcpy(&c->peer_session_id, enc_hs + 33, 4); - uint16_t padding_size = (uint16_t)enc_hs[37] | ((uint16_t)enc_hs[38] << 8); - c->peer_device_type = enc_hs[39]; - c->peer_keepalive_interval = ((uint16_t)enc_hs[40] << 8) | enc_hs[41]; - c->peer_flags = enc_hs[42]; - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "stcp_server: client handshake OK gop=%d sid=%08x padding=%u dev=%d ka=%u flags=%02x", - c->peer_got_initial_pkt, c->peer_session_id, padding_size, c->peer_device_type, c->peer_keepalive_interval, c->peer_flags); + { uint64_t rid_be; memcpy(&rid_be, enc_hs + 37, 8); c->peer_reset_id = be64toh(rid_be); } + uint16_t padding_size = (uint16_t)enc_hs[45] | ((uint16_t)enc_hs[46] << 8); + c->peer_device_type = enc_hs[47]; + c->peer_keepalive_interval = ((uint16_t)enc_hs[48] << 8) | enc_hs[49]; + c->peer_flags = enc_hs[50]; + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "stcp_server: client handshake OK gop=%d sid=%08x rid=%016llx padding=%u dev=%d ka=%u flags=%02x", + c->peer_got_initial_pkt, c->peer_session_id, (unsigned long long)c->peer_reset_id, padding_size, c->peer_device_type, c->peer_keepalive_interval, c->peer_flags); stcp_recv_set(c, padding_size, 0, server_hs_phase2_cb); } @@ -93,14 +94,16 @@ static void server_hs_phase2_cb(struct stcp_conn *c, uint8_t *data, size_t len) uint8_t server_gop = 0; uint32_t server_sid = 0; + uint64_t server_rid = 0; if (c->inst && c->peer_pubkey_set) { uint64_t node_id = sc_derive_node_id_from_pubkey(c->peer_pubkey); struct ETCP_CONN *conn = instance_find_conn(c->inst, node_id); server_gop = conn ? conn->got_initial_pkt : 0; server_sid = conn ? conn->session_id : 0; - DEBUG_INFO(DEBUG_CATEGORY_ETCP, "stcp_server: node=0x%016llx conn=%p server_gop=%d server_sid=%08x client_gop=%d client_sid=%08x", - (unsigned long long)node_id, (void*)conn, server_gop, server_sid, - c->peer_got_initial_pkt, c->peer_session_id); + server_rid = conn ? conn->reset_id : 0; + DEBUG_INFO(DEBUG_CATEGORY_ETCP, "stcp_server: node=0x%016llx conn=%p server_gop=%d server_sid=%08x server_rid=%016llx client_gop=%d client_sid=%08x client_rid=%016llx", + (unsigned long long)node_id, (void*)conn, server_gop, server_sid, (unsigned long long)server_rid, + c->peer_got_initial_pkt, c->peer_session_id, (unsigned long long)c->peer_reset_id); } uint8_t salt2[SC_PUBKEY_ENC_SALT_SIZE]; @@ -115,17 +118,18 @@ static void server_hs_phase2_cb(struct stcp_conn *c, uint8_t *data, size_t len) memcpy(resp, salt2, SC_PUBKEY_ENC_SALT_SIZE); sc_obfuscate_pubkey(salt2, c->peer_pubkey, c->my_keys.public_key, resp + SC_PUBKEY_ENC_SALT_SIZE); - uint8_t plain_hs[43]; memcpy(plain_hs, c->my_ed25519_pubkey, SC_PUBKEY_SIZE); + uint8_t plain_hs[STCP_HS_PLAIN_SIZE]; memcpy(plain_hs, c->my_ed25519_pubkey, SC_PUBKEY_SIZE); plain_hs[32] = server_gop; memcpy(plain_hs + 33, &server_sid, 4); - plain_hs[37] = (uint8_t)padding; plain_hs[38] = (uint8_t)(padding >> 8); - plain_hs[39] = c->device_type; - *(uint16_t*)(plain_hs + 40) = htobe16(c->keepalive_interval); - plain_hs[42] = (c->peer_flags & STCP_HANDSHAKE_FLAG_PING); - uint32_t crc = crc32_calc(plain_hs, 43); + { uint64_t rid_be = htobe64(server_rid); memcpy(plain_hs + 37, &rid_be, 8); } + plain_hs[45] = (uint8_t)padding; plain_hs[46] = (uint8_t)(padding >> 8); + plain_hs[47] = c->device_type; + *(uint16_t*)(plain_hs + 48) = htobe16(c->keepalive_interval); + plain_hs[50] = (c->peer_flags & STCP_HANDSHAKE_FLAG_PING); + uint32_t crc = crc32_calc(plain_hs, STCP_HS_PLAIN_SIZE); uint8_t *enc_dst = resp + SC_PUBKEY_ENC_SIZE; - memcpy(enc_dst, plain_hs, 43); - enc_dst[43] = (uint8_t)(crc >> 0); enc_dst[44] = (uint8_t)(crc >> 8); enc_dst[45] = (uint8_t)(crc >> 16); enc_dst[46] = (uint8_t)(crc >> 24); + memcpy(enc_dst, plain_hs, STCP_HS_PLAIN_SIZE); + enc_dst[51] = (uint8_t)(crc >> 0); enc_dst[52] = (uint8_t)(crc >> 8); enc_dst[53] = (uint8_t)(crc >> 16); enc_dst[54] = (uint8_t)(crc >> 24); log_dump(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_CRYPTO, "stcp_server hs_resp BEFORE xor", enc_dst, STCP_HS_ENC_SERVER); if (sc_stream_xor(&c->stream_send, enc_dst, STCP_HS_ENC_SERVER) != SC_OK) { DEBUG_ERROR(DEBUG_CATEGORY_ETCP, "encrypt failed"); diff --git a/src/utun_instance.c b/src/utun_instance.c index 6ec9deb9..c185a1f9 100644 --- a/src/utun_instance.c +++ b/src/utun_instance.c @@ -230,6 +230,8 @@ static int instance_init_common(struct UTUN_INSTANCE* instance, struct UASYNC* u socket_monitor_init(instance); auto_socket_init(instance); + if (g_topo_group_enabled && instance->topo_groups) + topo_node_update_my_addresses(instance); // conn_mgr initialized inside topo_group_create (via topo_groups_init) diff --git a/tests/Makefile.am b/tests/Makefile.am index 1b34e108..6db0d0d7 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -28,6 +28,7 @@ check_PROGRAMS = \ test_etcp_simple_traffic \ test_etcp_100_packets \ test_etcp_reconnect \ + test_etcp_seq_collision \ test_pkt_normalizer_etcp \ test_lwip_tcp \ test_u_async_performance \ @@ -191,6 +192,10 @@ test_etcp_reconnect_SOURCES = test_etcp_reconnect.c test_etcp_reconnect_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_etcp_reconnect_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_etcp_seq_collision_SOURCES = test_etcp_seq_collision.c +test_etcp_seq_collision_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib +test_etcp_seq_collision_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) + test_pkt_normalizer_etcp_SOURCES = test_pkt_normalizer_etcp.c test_pkt_normalizer_etcp_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_pkt_normalizer_etcp_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) diff --git a/tests/test_etcp_seq_collision.c b/tests/test_etcp_seq_collision.c new file mode 100644 index 00000000..005834ac --- /dev/null +++ b/tests/test_etcp_seq_collision.c @@ -0,0 +1,243 @@ +#include +#include +#include +#include "../lib/platform_compat.h" +#include "test_utils.h" +#ifdef _WIN32 +#include +#include +#include +#define getpid _getpid +#else +#include +#endif +#include + +#include "etcp.h" +#include "etcp_connections.h" +#include "etcp_api.h" +#include "../src/config_parser.h" +#include "../src/utun_instance.h" +#include "routing.h" +#include "secure_channel.h" +#include "../lib/u_async.h" +#include "../lib/ll_queue.h" +#include "../lib/debug_config.h" + +#define TEST_TIMEOUT_MS 60000 +#define PACKET_SIZE 100 +#define PACKETS_FIRST 3 + +static struct UTUN_INSTANCE* server_instance = NULL; +static struct UTUN_INSTANCE* client_instance = NULL; +static struct UASYNC* ua = NULL; + +static char temp_dir[] = "/tmp/utun_col_XXXXXX"; +static char server_config_path[256]; +static char client_config_path[256]; +static int server_port = 0; +static int client_port = 0; + +static int test_completed = 0; +static void* global_timeout_id = NULL; + +static int phase = 0; +static int packets_sent = 0; +static int packets_received = 0; +static uint8_t packet_buffer[PACKET_SIZE]; + +// client = 0x1111 (master, smaller node_id), server = 0x2222 (slave) +static int create_temp_configs(void) { + if (test_mkdtemp(temp_dir) != 0) { fprintf(stderr, "Failed to create temp directory\n"); return -1; } + int base_port = 45000 + (getpid() % 10000); + server_port = base_port; + client_port = base_port + 1; + + snprintf(server_config_path, sizeof(server_config_path), "%s/server.conf", temp_dir); + snprintf(client_config_path, sizeof(client_config_path), "%s/client.conf", temp_dir); + + FILE* f = fopen(server_config_path, "w"); + if (!f) { fprintf(stderr, "Failed to create server config file\n"); return -1; } + fprintf(f, + "[global]\n" + "my_node_id=0x2222222222222222\n" + "my_private_key=38240cb82199e504686507f11f6eaa4f740fde6f0c425c495e49a523019a5d68\n" + "my_public_key=ce8871f07fa056c636d297115f231b08c29cdf94e0d440fce83a07c34416d36a\n" + "tun_ip=10.99.0.1/24\n" + "tun_ifname=tun99\n" + "\n" + "[server: test]\n" + "addr=127.0.0.1:%d\n" + "type=public\n" + "\n" + "[allowed_keys]\n" + "allow_all=1\n", + server_port); + fclose(f); + + f = fopen(client_config_path, "w"); + if (!f) { fprintf(stderr, "Failed to create client config file\n"); test_unlink(server_config_path); return -1; } + fprintf(f, + "[global]\n" + "my_node_id=0x1111111111111111\n" + "my_private_key=704f2e012c8fa8768130cb0f988a997dccb628372bc5ceccacc78dcbfec5916f\n" + "my_public_key=b3193173def895bd0fcea6f86af077c7d77216f10395275f627ac18242ec0f01\n" + "tun_ip=10.99.0.2/24\n" + "tun_ifname=tun98\n" + "\n" + "[server: test]\n" + "addr=127.0.0.1:%d\n" + "type=public\n" + "\n" + "[client: test_client]\n" + "keepalive=1\n" + "peer_public_key=ce8871f07fa056c636d297115f231b08c29cdf94e0d440fce83a07c34416d36a\n" + "link=test:127.0.0.1:%d\n", + client_port, server_port); + fclose(f); + return 0; +} + +static void cleanup_temp_configs(void) { + if (server_config_path[0]) test_unlink(server_config_path); + if (client_config_path[0]) test_unlink(client_config_path); + if (temp_dir[0]) test_rmdir(temp_dir); +} + +static struct ETCP_CONN* server_conn(void) { + if (!server_instance || !server_instance->connections || !server_instance->connections->head) return NULL; + return ((struct conn_queue_entry*)server_instance->connections->head->data)->conn; +} + +static struct ETCP_CONN* client_conn(void) { + if (!client_instance || !client_instance->connections || !client_instance->connections->head) return NULL; + return ((struct conn_queue_entry*)client_instance->connections->head->data)->conn; +} + +static int is_connected(void) { + struct ETCP_CONN* c = client_conn(); + return c && c->initialized && c->links_up; +} + +static void drain_received(int count_flag) { + struct ETCP_CONN* conn = server_conn(); + if (!conn || !conn->output_queue) return; + queue_set_callback(conn->output_queue, NULL, NULL); + struct ETCP_FRAGMENT* pkt; + while ((pkt = (struct ETCP_FRAGMENT*)queue_data_get(conn->output_queue)) != NULL) { + if (count_flag) packets_received++; + if (pkt->ll.dgram) memory_pool_free(conn->instance->data_pool, pkt->ll.dgram); + queue_entry_free((struct ll_entry*)pkt); + } +} + +static void send_packets(int n) { + struct ETCP_CONN* conn = client_conn(); + if (!conn || !conn->initialized) return; + for (int i = 0; i < n; i++) { + if (queue_entry_count(conn->input_queue) > 10) break; + for (int j = 0; j < PACKET_SIZE; j++) packet_buffer[j] = (uint8_t)((packets_sent + j) % 256); + if (etcp_int_send(conn, packet_buffer, PACKET_SIZE) != 0) break; + packets_sent++; + } +} + +static void monitor(void* arg) { + (void)arg; + if (test_completed) return; + + switch (phase) { + case 0: // wait connect + drain_received(0); + if (is_connected()) { + printf("=== Phase 0: connected ===\n"); + phase = 1; + } + break; + + case 1: // send PACKETS_FIRST + send_packets(PACKETS_FIRST - packets_sent); + drain_received(1); + if (packets_received >= PACKETS_FIRST) { + printf("=== Phase 1: %d packets delivered ===\n", packets_received); + printf(" forcing one-sided reinit on client (master 0x1111)\n"); + etcp_conn_reinit(client_conn(), "test seq collision"); + { struct ETCP_LINK* l = client_conn() ? client_conn()->links : NULL; + while (l) { etcp_link_enter_reinit(l); l = l->next; } } + phase = 2; + } + break; + + case 2: // wait client re-ready, then send 1 packet + drain_received(1); + if (is_connected()) { + printf(" client reinitialized, sending post-reinit packet\n"); + send_packets(1); + phase = 3; + } + break; + + case 3: // verify delivery + drain_received(1); + if (packets_received >= PACKETS_FIRST + 1) { + printf("\n[PASS] post-reinit packet delivered (recv=%d)\n", packets_received); + test_completed = 1; + return; + } + break; + } + + if (!test_completed) uasync_set_timeout(ua, 20, NULL, monitor, "test_monitor"); +} + +static void test_timeout(void* arg) { + (void)arg; + if (!test_completed) { + printf("\n[FAIL] timeout at phase %d: sent=%d recv=%d — post-reinit packet silently dropped\n", + phase, packets_sent, packets_received); + test_completed = 2; + } +} + +int main(void) { + if (create_temp_configs() != 0) return 1; + debug_config_init(); + utun_instance_set_tun_init_enabled(0); + + printf("=== ETCP Seq Collision Regression Test ===\n"); + + ua = uasync_create(); + server_instance = utun_instance_create(ua, server_config_path); + if (!server_instance || utun_instance_init(server_instance) < 0) { + fprintf(stderr, "Failed to create server\n"); + if (server_instance) utun_instance_destroy(server_instance); + uasync_destroy(ua, 0); cleanup_temp_configs(); return 1; + } + client_instance = utun_instance_create(ua, client_config_path); + if (!client_instance || utun_instance_init(client_instance) < 0) { + fprintf(stderr, "Failed to create client\n"); + utun_instance_destroy(server_instance); + if (client_instance) utun_instance_destroy(client_instance); + uasync_destroy(ua, 0); cleanup_temp_configs(); return 1; + } + + uasync_set_timeout(ua, 20, NULL, monitor, "test_monitor"); + global_timeout_id = uasync_set_timeout(ua, TEST_TIMEOUT_MS * 10, NULL, test_timeout, "test_timeout"); + + int elapsed = 0; + int poll_interval = 20; + while (!test_completed && elapsed < TEST_TIMEOUT_MS * 10 + 5000) { + uasync_poll(ua, poll_interval); + elapsed += poll_interval; + } + + if (global_timeout_id) uasync_cancel_timeout(ua, global_timeout_id); + if (server_instance) { server_instance->running = 0; utun_instance_destroy(server_instance); server_instance = NULL; } + if (client_instance) { client_instance->running = 0; utun_instance_destroy(client_instance); client_instance = NULL; } + if (ua) { uasync_destroy(ua, 0); ua = NULL; } + cleanup_temp_configs(); + + if (test_completed == 1) { printf("=== TEST PASSED ===\n"); return 0; } + printf("=== TEST FAILED ===\n"); + return 1; +} diff --git a/tools/chatgui-android/.gitignore b/tools/chatgui-android/.gitignore index 9a0fe6e9..36f2d824 100644 --- a/tools/chatgui-android/.gitignore +++ b/tools/chatgui-android/.gitignore @@ -1,4 +1,5 @@ build/ +build_headless/ .gradle/ *.apk *.iml diff --git a/tools/chatgui/src/accountlist.cpp b/tools/chatgui/src/accountlist.cpp index a3d0b831..1767b10a 100644 --- a/tools/chatgui/src/accountlist.cpp +++ b/tools/chatgui/src/accountlist.cpp @@ -200,14 +200,13 @@ void AccountList::showMemberDetail(quint64 nodeId) { sqlite3* d = m_db->m_db; sqlite3_stmt* st = nullptr; if (sqlite3_prepare_v2(d, - "SELECT name, online, last_seen_at, created_at FROM nodes WHERE node_id=?", + "SELECT name, last_seen_at, created_at FROM nodes WHERE node_id=?", -1, &st, nullptr) == SQLITE_OK) { sqlite3_bind_int64(st, 1, (sqlite3_int64)nodeId); if (sqlite3_step(st) == SQLITE_ROW) { m_detailName = DbManager::colText(st, 0); - m_detailOnline = sqlite3_column_int(st, 1) != 0; - m_detailLastSeen = sqlite3_column_int64(st, 2); - m_detailCreated = sqlite3_column_int64(st, 3); + m_detailLastSeen = sqlite3_column_int64(st, 1); + m_detailCreated = sqlite3_column_int64(st, 2); } sqlite3_finalize(st); } @@ -283,6 +282,7 @@ void AccountList::onMemberDetailData(uint64_t nodeId, const uint8_t* data, int l m_snapConnPresent = (flags & 1) != 0; m_snapConnUp = (flags & 2) != 0; m_snapConnInit = (flags & 4) != 0; + m_detailOnline = (flags & 8) != 0; uint16_t ownTcpActive; memcpy(&ownTcpActive, p, 2); p += 2; m_snapOwnTcpActive = ownTcpActive; diff --git a/tools/chatgui/src/mainwindow.cpp b/tools/chatgui/src/mainwindow.cpp index b957f72c..8d1b5981 100644 --- a/tools/chatgui/src/mainwindow.cpp +++ b/tools/chatgui/src/mainwindow.cpp @@ -82,7 +82,7 @@ static void onMemberListCallback(const uint8_t* data, int len) { } static void onMemberUpdatedCallback(const uint8_t* data, int len) { - if (len < 1 + 64 + 1 + 79) return; + if (len < 1) return; uint8_t cl = data[0]; if (len < 1 + (int)cl + 1 + 79) return; if (s_mainWindow) @@ -90,7 +90,7 @@ static void onMemberUpdatedCallback(const uint8_t* data, int len) { } static void onMemberRemovedCallback(const uint8_t* data, int len) { - if (len < 1 + 64 + 8) return; + if (len < 1) return; uint8_t cl = data[0]; if (len < 1 + (int)cl + 8) return; uint64_t nid; memcpy(&nid, data + 1 + cl, 8); diff --git a/tools/chatgui/src/nodespage.cpp b/tools/chatgui/src/nodespage.cpp index 06a6d2ee..68a0d559 100644 --- a/tools/chatgui/src/nodespage.cpp +++ b/tools/chatgui/src/nodespage.cpp @@ -62,8 +62,8 @@ NodesPage::NodesPage(DbManager* db, QWidget* parent) title->setStyleSheet("font-weight: bold; font-size: 13px;"); layout->addWidget(title); - m_table = new QTableWidget(0, 6, this); - m_table->setHorizontalHeaderLabels({"Node ID", "Name", "Online", "x25519", "ed25519", "Addrs"}); + m_table = new QTableWidget(0, 5, this); + m_table->setHorizontalHeaderLabels({"Node ID", "Name", "x25519", "ed25519", "Addrs"}); m_table->horizontalHeader()->setStretchLastSection(true); m_table->setSelectionBehavior(QAbstractItemView::SelectRows); m_table->setSelectionMode(QAbstractItemView::SingleSelection); @@ -121,7 +121,7 @@ void NodesPage::refreshNodes() { sqlite3_stmt* st = nullptr; if (sqlite3_prepare_v2(d, - "SELECT node_id, name, online, x25519_pubkey, ed25519_pubkey," + "SELECT node_id, name, x25519_pubkey, ed25519_pubkey," " last_seen_at, created_at FROM nodes ORDER BY node_id", -1, &st, nullptr) != SQLITE_OK) return; @@ -136,23 +136,22 @@ void NodesPage::refreshNodes() { setCell(0, hex16(nid)); setCell(1, DbManager::colText(st, 1)); - setCell(2, sqlite3_column_int(st, 2) ? "yes" : "no"); - QByteArray x25 = DbManager::colBlob(st, 3); - QByteArray ed = DbManager::colBlob(st, 4); - setCell(3, x25.isEmpty() ? "-" : shortenKey(x25)); - setCell(4, ed.isEmpty() ? "-" : shortenKey(ed)); + QByteArray x25 = DbManager::colBlob(st, 2); + QByteArray ed = DbManager::colBlob(st, 3); + setCell(2, x25.isEmpty() ? "-" : shortenKey(x25)); + setCell(3, ed.isEmpty() ? "-" : shortenKey(ed)); sqlite3_stmt* ac = nullptr; sqlite3_prepare_v2(d, "SELECT COUNT(*) FROM node_addresses WHERE node_id=?", -1, &ac, nullptr); if (ac) { sqlite3_bind_int64(ac, 1, (sqlite3_int64)nid); - setCell(5, sqlite3_step(ac) == SQLITE_ROW + setCell(4, sqlite3_step(ac) == SQLITE_ROW ? QString::number(sqlite3_column_int(ac, 0)) : "0"); sqlite3_finalize(ac); } else { - setCell(5, "0"); + setCell(4, "0"); } row++; } @@ -164,7 +163,7 @@ void NodesPage::refreshNodes() { st = nullptr; sqlite3_prepare_v2(d, - "SELECT node_id, name, online, x25519_pubkey, ed25519_pubkey," + "SELECT node_id, name, x25519_pubkey, ed25519_pubkey," " last_seen_at, created_at FROM nodes ORDER BY node_id", -1, &st, nullptr); if (!st) { m_dumpText->setPlainText("Query error"); return; } @@ -176,16 +175,14 @@ void NodesPage::refreshNodes() { quint64 nid = (quint64)sqlite3_column_int64(st, 0); QString name = DbManager::colText(st, 1); - int online = sqlite3_column_int(st, 2); - QByteArray x25 = DbManager::colBlob(st, 3); - QByteArray ed = DbManager::colBlob(st, 4); - qint64 last_seen = sqlite3_column_int64(st, 5); - qint64 created = sqlite3_column_int64(st, 6); + QByteArray x25 = DbManager::colBlob(st, 2); + QByteArray ed = DbManager::colBlob(st, 3); + qint64 last_seen = sqlite3_column_int64(st, 4); + qint64 created = sqlite3_column_int64(st, 5); dump += QString("Node: %1").arg(hex16(nid)); if (!name.isEmpty()) dump += QString(" (%1)").arg(name); dump += QString(" id=0x%1\n").arg(nid, 16, 16, QChar('0')); - dump += QString(" online: %1\n").arg(online ? "yes" : "no"); dump += QString(" x25519: %1\n").arg(x25.isEmpty() ? "-" : shortenKey(x25)); dump += QString(" ed25519: %1\n").arg(ed.isEmpty() ? "-" : shortenKey(ed)); if (last_seen > 0)