Browse Source

etcp: reset_id эпоха ресета — фикс потери NODEINFO при флапе

Добавлен 64-битный reset_id (эпоха ресета) в INIT-фреймы и STCP-handshake:
- генерится один раз при создании conn, меняется только при фатальном ресете
  (normalizer desync / bad fragment size);
- при приёме чужого id slave принимает эпоху пира и ресетится, master держит свою;
- первый id принимается без ресета (свежие при init);
- «tcp peer clean» реинит гейтится по смене id — нет повторного сброса seq.

Это устраняет коллизию seq=1 (RX dup) из-за двойного реинита нескольких TCP-линков,
из-за которой NODEINFO пира терялся и member не возвращался online.

Также: chat member online теперь из topo_group (без БД/merkle), тест test_etcp_seq_collision.
v2
evgeny 4 weeks ago
parent
commit
61c572cae4
  1. 89
      BUG_bgp_reconnect.md
  2. 24
      src/chat/chat_core.c
  3. 70
      src/chat/chat_member.c
  4. 4
      src/chat/chat_member.h
  5. 8
      src/chat/chat_profile.c
  6. 10
      src/chat/chat_status.c
  7. 44
      src/chat/chat_sync.c
  8. 33
      src/chat/member_sync.c
  9. 25
      src/chat/member_sync.h
  10. 1
      src/chat/merkle_sync.h
  11. 194
      src/routing_layer/topo_node.c
  12. 23
      src/routing_layer/topo_node_sqlite.c
  13. 45
      src/transport_layer/etcp.c
  14. 8
      src/transport_layer/etcp.h
  15. 28
      src/transport_layer/etcp_connections.c
  16. 46
      src/transport_layer/etcp_connections.h
  17. 4
      src/transport_layer/pkt_normalizer.c
  18. 21
      src/transport_layer/stcp.h
  19. 31
      src/transport_layer/stcp_client.c
  20. 7
      src/transport_layer/stcp_link.c
  21. 1
      src/transport_layer/stcp_link.h
  22. 38
      src/transport_layer/stcp_server.c
  23. 2
      src/utun_instance.c
  24. 5
      tests/Makefile.am
  25. 243
      tests/test_etcp_seq_collision.c
  26. 1
      tools/chatgui-android/.gitignore
  27. 8
      tools/chatgui/src/accountlist.cpp
  28. 4
      tools/chatgui/src/mainwindow.cpp
  29. 31
      tools/chatgui/src/nodespage.cpp

89
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"`.

24
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';

70
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 <stdlib.h>
/* ─── 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)");
}

4
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;

8
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();
}

10
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) {

44
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;

33
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);

25
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), в котором произошло изменение.

1
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);

194
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);
}

23
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;

45
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;

8
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);// вызывается когда подключение инициализировано

28
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;
}

46
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)

4
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;
}

21
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

31
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);
}

7
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;
}

1
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);

38
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");

2
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)

5
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)

243
tests/test_etcp_seq_collision.c

@ -0,0 +1,243 @@
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include "../lib/platform_compat.h"
#include "test_utils.h"
#ifdef _WIN32
#include <windows.h>
#include <direct.h>
#include <process.h>
#define getpid _getpid
#else
#include <unistd.h>
#endif
#include <time.h>
#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;
}

1
tools/chatgui-android/.gitignore vendored

@ -1,4 +1,5 @@
build/
build_headless/
.gradle/
*.apk
*.iml

8
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;

4
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);

31
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)

Loading…
Cancel
Save