Browse Source

feat: auto-connect Phase 3 (local/netif addresses) + restart fix

Phase 1: connected peers (connected=1)
Phase 2: direct/EIM (addr_type=1,2), exclude self
Phase 3: local/strict NAT (addr_type=0,3), exclude self
Timeout: 5000→2000ms per attempt
Restart: DOWN triggers reconnect when phase=DONE & active=0

Debug logs: init, Phase1 per-peer, UP phase info, restart branches
topo_upd
Evgeny 2 months ago
parent
commit
8facce46ac
  1. 2
      src/routing_layer/conn_mgr.h
  2. 75
      src/routing_layer/topo_group_connect.c
  3. 35
      src/routing_layer/topo_node_sqlite.c
  4. 3
      src/routing_layer/topo_node_sqlite.h

2
src/routing_layer/conn_mgr.h

@ -69,7 +69,7 @@ struct TOPO_GROUP;
#define CONN_MGR_IDLE_CHECK_INTERVAL_TB 10000 ///< Интервал проверки idle-таймаута (tb, ~1 сек) #define CONN_MGR_IDLE_CHECK_INTERVAL_TB 10000 ///< Интервал проверки idle-таймаута (tb, ~1 сек)
#define CONN_MGR_LOCAL_SCAN_ATTEMPTS 3 ///< Число попыток локального сканирования #define CONN_MGR_LOCAL_SCAN_ATTEMPTS 3 ///< Число попыток локального сканирования
#define CONN_MGR_LOCAL_SCAN_TIMEOUT_MS 100 ///< Таймаут локального сканирования (мс) #define CONN_MGR_LOCAL_SCAN_TIMEOUT_MS 100 ///< Таймаут локального сканирования (мс)
#define CONN_MGR_CONNECT_DIRECT_TIMEOUT_MS 5000 ///< Таймаут фазы DIRECT (мс) #define CONN_MGR_CONNECT_DIRECT_TIMEOUT_MS 2000 ///< Таймаут фазы DIRECT (мс)
#define CONN_MGR_CONNECT_REVERSE_TIMEOUT_MS 15000 ///< Таймаут фазы REVERSE (мс) #define CONN_MGR_CONNECT_REVERSE_TIMEOUT_MS 15000 ///< Таймаут фазы REVERSE (мс)
#define CONN_MGR_INTERM_EXCHANGE_TIMEOUT_MS 15000 ///< Таймаут фазы INDIRECT (мс) #define CONN_MGR_INTERM_EXCHANGE_TIMEOUT_MS 15000 ///< Таймаут фазы INDIRECT (мс)

75
src/routing_layer/topo_group_connect.c

@ -5,7 +5,9 @@
* Таймаут = direct_timeout_ms. По таймауту отменяем незавершённые и помечаем * Таймаут = direct_timeout_ms. По таймауту отменяем незавершённые и помечаем
* connected=0 в БД. Если хоть один подключился → done. * connected=0 в БД. Если хоть один подключился → done.
* *
* Phase 2: последовательный перебор пиров с публичными адресами до первого успеха. * Phase 2: последовательный перебор пиров с прямыми/EIM адресами (addr_type=1,2).
*
* Phase 3: последовательный перебор пиров с локальными/strict NAT (addr_type=0,3).
* *
* Счётчик active_conn_count — число уникальных peer-узлов с реальными соединениями. * Счётчик active_conn_count — число уникальных peer-узлов с реальными соединениями.
* UP: инкремент только для первого соединения к peer_node_id. * UP: инкремент только для первого соединения к peer_node_id.
@ -29,7 +31,8 @@
#define TGC_ID "topo_group_connect" #define TGC_ID "topo_group_connect"
#define TGC_PHASE_ONE 0 #define TGC_PHASE_ONE 0
#define TGC_PHASE_TWO 1 #define TGC_PHASE_TWO 1
#define TGC_PHASE_DONE 2 #define TGC_PHASE_THREE 2
#define TGC_PHASE_DONE 3
struct TOPO_GROUP_CONNECT { struct TOPO_GROUP_CONNECT {
struct TOPO_GROUP* group; struct TOPO_GROUP* group;
@ -46,6 +49,8 @@ struct TOPO_GROUP_CONNECT {
static void tgc_phase2_try_next(struct TOPO_GROUP_CONNECT* gc); static void tgc_phase2_try_next(struct TOPO_GROUP_CONNECT* gc);
static void tgc_phase2_result(int result, uint64_t node_id, void* arg); static void tgc_phase2_result(int result, uint64_t node_id, void* arg);
static void tgc_phase3_try_next(struct TOPO_GROUP_CONNECT* gc);
static void tgc_phase3_result(int result, uint64_t node_id, void* arg);
static void tgc_phase1_timeout(void* arg); static void tgc_phase1_timeout(void* arg);
static void tgc_phase1_result(int result, uint64_t node_id, void* arg); static void tgc_phase1_result(int result, uint64_t node_id, void* arg);
@ -63,6 +68,7 @@ int topo_group_connect_init(struct TOPO_GROUP* group) {
uint64_t* ids = NULL; int count = 0; uint64_t* ids = NULL; int count = 0;
sqlite3* db = group->instance->topo_sqlite_db; sqlite3* db = group->instance->topo_sqlite_db;
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: init ch=%s db=%p", TGC_ID, group->channel_id, (void*)db);
if (!db || topo_node_sqlite_get_connected_peers(db, group->channel_id, &ids, &count) != 0 || count == 0) { if (!db || topo_node_sqlite_get_connected_peers(db, group->channel_id, &ids, &count) != 0 || count == 0) {
if (ids) { u_free(ids); ids = NULL; } if (ids) { u_free(ids); ids = NULL; }
gc->phase = TGC_PHASE_TWO; gc->phase = TGC_PHASE_TWO;
@ -73,8 +79,10 @@ int topo_group_connect_init(struct TOPO_GROUP* group) {
gc->phase = TGC_PHASE_ONE; gc->pending = count; gc->phase = TGC_PHASE_ONE; gc->pending = count;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: ch=%s Phase 1 launching %d connects", TGC_ID, group->channel_id, count); DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: ch=%s Phase 1 launching %d connects", TGC_ID, group->channel_id, count);
for (int i = 0; i < count; i++) for (int i = 0; i < count; i++) {
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: Phase1 connect to 0x%016llx (%d/%d)", TGC_ID, (unsigned long long)ids[i], i + 1, count);
conn_mgr_connect_node(group->conn_mgr, ids[i], 0, tgc_phase1_result, gc); conn_mgr_connect_node(group->conn_mgr, ids[i], 0, tgc_phase1_result, gc);
}
u_free(ids); u_free(ids);
gc->phase_timer = uasync_set_timeout(group->instance->ua, gc->phase_timer = uasync_set_timeout(group->instance->ua,
@ -93,9 +101,9 @@ void topo_group_connect_destroy(struct TOPO_GROUP* group) {
void topo_group_connect_restart(struct TOPO_GROUP* group) { void topo_group_connect_restart(struct TOPO_GROUP* group) {
struct TOPO_GROUP_CONNECT* gc = group->connect; struct TOPO_GROUP_CONNECT* gc = group->connect;
if (!gc) { topo_group_connect_init(group); return; } if (!gc) { DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: restart ch=%s no gc → init", TGC_ID, group->channel_id); topo_group_connect_init(group); return; }
if (gc->phase != TGC_PHASE_DONE) return; if (gc->phase != TGC_PHASE_DONE) { DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: restart ch=%s skip phase=%d", TGC_ID, group->channel_id, gc->phase); return; }
if (gc->active_conn_count > 0) return; if (gc->active_conn_count > 0) { DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: restart ch=%s skip active=%d", TGC_ID, group->channel_id, gc->active_conn_count); return; }
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: restart ch=%s", TGC_ID, group->channel_id); DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: restart ch=%s", TGC_ID, group->channel_id);
topo_group_connect_destroy(group); topo_group_connect_destroy(group);
topo_group_connect_init(group); topo_group_connect_init(group);
@ -133,6 +141,8 @@ void topo_group_connect_on_up(struct TOPO_GROUP* group, struct ETCP_CONN* conn)
if (db) topo_node_sqlite_set_connected(db, group->channel_id, peer, 1); if (db) topo_node_sqlite_set_connected(db, group->channel_id, peer, 1);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: UP peer=0x%016llx ch=%s active=%d", DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: UP peer=0x%016llx ch=%s active=%d",
TGC_ID, (unsigned long long)peer, group->channel_id, gc->active_conn_count); TGC_ID, (unsigned long long)peer, group->channel_id, gc->active_conn_count);
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: UP phase=%d candidate_count=%d",
TGC_ID, gc->phase, gc->candidate_count);
} }
void topo_group_connect_on_down(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { void topo_group_connect_on_down(struct TOPO_GROUP* group, struct ETCP_CONN* conn) {
@ -173,6 +183,10 @@ void topo_group_connect_on_down(struct TOPO_GROUP* group, struct ETCP_CONN* conn
if (db) topo_node_sqlite_set_connected(db, group->channel_id, peer, 0); if (db) topo_node_sqlite_set_connected(db, group->channel_id, peer, 0);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: DOWN peer=0x%016llx ch=%s active=%d", DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: DOWN peer=0x%016llx ch=%s active=%d",
TGC_ID, (unsigned long long)peer, group->channel_id, gc->active_conn_count); TGC_ID, (unsigned long long)peer, group->channel_id, gc->active_conn_count);
if (gc->phase == TGC_PHASE_DONE && gc->active_conn_count == 0) {
DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: DOWN restart eligible ch=%s → restarting", TGC_ID, group->channel_id);
topo_group_connect_restart(group);
}
} }
/* ═══════════════════════════════════════════════════════════════════════ /* ═══════════════════════════════════════════════════════════════════════
@ -225,7 +239,8 @@ static void tgc_phase2_try_next(struct TOPO_GROUP_CONNECT* gc) {
if (gc->phase != TGC_PHASE_TWO) return; if (gc->phase != TGC_PHASE_TWO) return;
if (gc->candidate_count == 0) { if (gc->candidate_count == 0) {
sqlite3* db = gc->group->instance->topo_sqlite_db; sqlite3* db = gc->group->instance->topo_sqlite_db;
topo_node_sqlite_get_public_peers(db, gc->group->channel_id, &gc->candidate_ids, &gc->candidate_count); topo_node_sqlite_get_public_peers(db, gc->group->channel_id, &gc->candidate_ids, &gc->candidate_count,
gc->group->instance->node_id);
gc->cursor = 0; gc->cursor = 0;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 loaded %d public peers for ch=%s", DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 loaded %d public peers for ch=%s",
TGC_ID, gc->candidate_count, gc->group->channel_id); TGC_ID, gc->candidate_count, gc->group->channel_id);
@ -244,7 +259,8 @@ static void tgc_phase2_try_next(struct TOPO_GROUP_CONNECT* gc) {
return; return;
} }
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 exhausted — no connections, ch=%s", TGC_ID, gc->group->channel_id); DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 exhausted — no connections, ch=%s", TGC_ID, gc->group->channel_id);
gc->phase = TGC_PHASE_DONE; gc->phase = TGC_PHASE_THREE; gc->candidate_count = 0;
tgc_phase3_try_next(gc);
} }
static void tgc_phase2_result(int result, uint64_t node_id, void* arg) { static void tgc_phase2_result(int result, uint64_t node_id, void* arg) {
@ -258,3 +274,46 @@ static void tgc_phase2_result(int result, uint64_t node_id, void* arg) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 0x%016llx FAIL — try next", TGC_ID, (unsigned long long)node_id); DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 0x%016llx FAIL — try next", TGC_ID, (unsigned long long)node_id);
tgc_phase2_try_next(gc); tgc_phase2_try_next(gc);
} }
/* ═══════════════════════════════════════════════════════════════════════
* Phase 3 — локальные/strict NAT адреса (addr_type=0,3)
* ══════════════════════════════════════════════════════════════════════ */
static void tgc_phase3_try_next(struct TOPO_GROUP_CONNECT* gc) {
if (gc->phase != TGC_PHASE_THREE) return;
if (gc->candidate_count == 0) {
sqlite3* db = gc->group->instance->topo_sqlite_db;
topo_node_sqlite_get_local_peers(db, gc->group->channel_id, &gc->candidate_ids, &gc->candidate_count,
gc->group->instance->node_id);
gc->cursor = 0;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase3 loaded %d local peers for ch=%s",
TGC_ID, gc->candidate_count, gc->group->channel_id);
}
while (gc->cursor < gc->candidate_count) {
uint64_t nid = gc->candidate_ids[gc->cursor++];
uint8_t status; uint8_t conn_type;
conn_mgr_get_status(gc->group->conn_mgr, nid, &status, &conn_type);
if (status == CONN_MGR_STATE_CONNECTED) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase3 skip 0x%016llx already connected", TGC_ID, (unsigned long long)nid);
continue;
}
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase3 trying 0x%016llx (%d/%d)", TGC_ID,
(unsigned long long)nid, gc->cursor, gc->candidate_count);
conn_mgr_connect_node(gc->group->conn_mgr, nid, 0, tgc_phase3_result, gc);
return;
}
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase3 exhausted — no connections, ch=%s", TGC_ID, gc->group->channel_id);
gc->phase = TGC_PHASE_DONE;
}
static void tgc_phase3_result(int result, uint64_t node_id, void* arg) {
struct TOPO_GROUP_CONNECT* gc = (struct TOPO_GROUP_CONNECT*)arg;
if (gc->phase != TGC_PHASE_THREE) return;
if (result == CONN_MGR_OK) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase3 connected to 0x%016llx — done", TGC_ID, (unsigned long long)node_id);
gc->phase = TGC_PHASE_DONE;
return;
}
DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase3 0x%016llx FAIL — try next", TGC_ID, (unsigned long long)node_id);
tgc_phase3_try_next(gc);
}

35
src/routing_layer/topo_node_sqlite.c

@ -762,18 +762,45 @@ int topo_node_sqlite_get_connected_peers(sqlite3* db, const char* channel_id,
} }
int topo_node_sqlite_get_public_peers(sqlite3* db, const char* channel_id, int topo_node_sqlite_get_public_peers(sqlite3* db, const char* channel_id,
uint64_t** out_ids, int* out_count) { uint64_t** out_ids, int* out_count, uint64_t self_node_id) {
if (!db || !channel_id || !out_ids || !out_count) return -1;
*out_ids = NULL; *out_count = 0;
char peers_tbl[80]; peers_table_name(channel_id, peers_tbl, sizeof(peers_tbl));
char sql[300];
snprintf(sql, sizeof(sql),
"SELECT p.node_id FROM \"%s\" p"
" WHERE p.node_id!=? AND EXISTS (SELECT 1 FROM node_addresses a"
" WHERE a.node_id=p.node_id AND a.family=4 AND a.addr_type IN (1,2))"
" ORDER BY p.node_RTT", peers_tbl);
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) return -1;
sqlite3_bind_int64(st, 1, (sqlite3_int64)self_node_id);
int cnt = 0;
while (sqlite3_step(st) == SQLITE_ROW) cnt++;
if (cnt == 0) { sqlite3_finalize(st); return 0; }
uint64_t* ids = u_malloc((size_t)cnt * sizeof(uint64_t));
if (!ids) { sqlite3_finalize(st); return -1; }
sqlite3_reset(st); int i = 0;
while (sqlite3_step(st) == SQLITE_ROW) ids[i++] = (uint64_t)sqlite3_column_int64(st, 0);
sqlite3_finalize(st);
*out_ids = ids; *out_count = cnt;
return 0;
}
int topo_node_sqlite_get_local_peers(sqlite3* db, const char* channel_id,
uint64_t** out_ids, int* out_count, uint64_t self_node_id) {
if (!db || !channel_id || !out_ids || !out_count) return -1; if (!db || !channel_id || !out_ids || !out_count) return -1;
*out_ids = NULL; *out_count = 0; *out_ids = NULL; *out_count = 0;
char peers_tbl[80]; peers_table_name(channel_id, peers_tbl, sizeof(peers_tbl)); char peers_tbl[80]; peers_table_name(channel_id, peers_tbl, sizeof(peers_tbl));
char sql[280]; char sql[300];
snprintf(sql, sizeof(sql), snprintf(sql, sizeof(sql),
"SELECT p.node_id FROM \"%s\" p" "SELECT p.node_id FROM \"%s\" p"
" WHERE EXISTS (SELECT 1 FROM node_addresses a" " WHERE p.node_id!=? AND EXISTS (SELECT 1 FROM node_addresses a"
" WHERE a.node_id=p.node_id AND a.family=4 AND a.addr_type!=0)" " WHERE a.node_id=p.node_id AND a.family=4 AND a.addr_type IN (0,3))"
" ORDER BY p.node_RTT", peers_tbl); " ORDER BY p.node_RTT", peers_tbl);
sqlite3_stmt* st = NULL; sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) return -1; if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) return -1;
sqlite3_bind_int64(st, 1, (sqlite3_int64)self_node_id);
int cnt = 0; int cnt = 0;
while (sqlite3_step(st) == SQLITE_ROW) cnt++; while (sqlite3_step(st) == SQLITE_ROW) cnt++;
if (cnt == 0) { sqlite3_finalize(st); return 0; } if (cnt == 0) { sqlite3_finalize(st); return 0; }

3
src/routing_layer/topo_node_sqlite.h

@ -69,6 +69,7 @@ struct TOPO_NODE* topo_node_sqlite_node_load(sqlite3* db, struct TOPO_GROUPS* gr
int topo_node_sqlite_set_connected(sqlite3* db, const char* channel_id, uint64_t node_id, int connected); int topo_node_sqlite_set_connected(sqlite3* db, const char* channel_id, uint64_t node_id, int connected);
int topo_node_sqlite_get_connected_peers(sqlite3* db, const char* channel_id, uint64_t** out_ids, int* out_count); int topo_node_sqlite_get_connected_peers(sqlite3* db, const char* channel_id, uint64_t** out_ids, int* out_count);
int topo_node_sqlite_get_public_peers(sqlite3* db, const char* channel_id, uint64_t** out_ids, int* out_count); int topo_node_sqlite_get_public_peers(sqlite3* db, const char* channel_id, uint64_t** out_ids, int* out_count, uint64_t self_node_id);
int topo_node_sqlite_get_local_peers(sqlite3* db, const char* channel_id, uint64_t** out_ids, int* out_count, uint64_t self_node_id);
#endif #endif

Loading…
Cancel
Save