Browse Source

conn_mgr: handle-based API (conn_mgr_open/close), фазы DIRECT/REVERSE/INDIRECT, тесты phases+handles

proxy
evgeny 2 weeks ago
parent
commit
1de1409699
  1. 1
      src/routing_layer/conn_mgr.h
  2. 139
      src/routing_layer/conn_mgr_core.c
  3. 221
      src/routing_layer/conn_mgr_doc.md
  4. 47
      src/routing_layer/conn_mgr_indirect.c
  5. 16
      src/routing_layer/conn_mgr_monitor.c
  6. 15
      src/routing_layer/conn_mgr_priv.h
  7. 10
      src/routing_layer/route_connectivity.c
  8. 4
      src/routing_layer/route_connectivity.h
  9. 10
      tests/Makefile.am
  10. 82
      tests/test_conn_mgr_handles.c
  11. 322
      tests/test_conn_mgr_phases.c
  12. 34
      tools/chatgui-android/build.sh

1
src/routing_layer/conn_mgr.h

@ -43,7 +43,6 @@ struct TOPO_GROUP;
*
* bg_ping: раз в ~100ms один узел группы за тик. Поддерживает свежие RTT.
* candidate_ping: раз в ~2с, удаляет протухших кандидатов.
* idle: раз в ~1с, рвёт соединение при неактивности.
*
* === Refcounting и закрытие ===
*

139
src/routing_layer/conn_mgr_core.c

@ -180,7 +180,6 @@ struct CONN_MGR_ENTRY* cm_ensure_entry(struct CONN_MGR* mgr, uint64_t node_id) {
void cm_entry_cleanup(struct CONN_MGR_ENTRY* entry) {
if (!entry || entry->state == CONN_MGR_STATE_DISCONNECTED) return;
if (entry->idle_timer) { uasync_cancel_timeout(entry->mgr->instance->ua, entry->idle_timer); entry->idle_timer = NULL; }
if (entry->main.timer) { uasync_cancel_timeout(entry->mgr->instance->ua, entry->main.timer); entry->main.timer = NULL; }
{ struct cm_exchange_pending** pp = &entry->mgr->exchange_pending;
while (*pp) {
@ -204,6 +203,7 @@ void cm_clear_nodeinfo(struct CONN_MGR* mgr, uint64_t node_id) {
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(mgr->group, node_id);
if (!nq) return;
nq->conn_mgr_type = CONN_TYPE_NONE; nq->conn_mgr_intermediariy_count = 0;
nq->conn_presence &= ~NCONN_INDIRECT; nq->conn_up &= ~NCONN_INDIRECT;
}
void cm_update_nodeinfo(struct CONN_MGR_ENTRY* entry) {
@ -240,7 +240,7 @@ void cm_cleanup_db_node(struct CONN_MGR_ENTRY* entry) {
* TIMEOUT — NCD таймаут истёк, переходим к REVERSE или INDIRECT, либо фейлим db_node.
* DOWN — соединение упало, доставляем DOWN. */
void cm_ncd_callback(struct NODE_CONN_DIRECT* ncd_h, enum ncd_event ncd_ev, void* arg) {
struct CONN_MGR_ENTRY* entry = (struct CONN_MGR_ENTRY*)arg; (void)ncd_h;
struct CONN_MGR_ENTRY* entry = (struct CONN_MGR_ENTRY*)arg;
if (!entry || !entry->mgr || entry->state == CONN_MGR_STATE_DISCONNECTED) return;
switch (ncd_ev) {
case NCD_EVENT_UP:
@ -248,9 +248,16 @@ void cm_ncd_callback(struct NODE_CONN_DIRECT* ncd_h, enum ncd_event ncd_ev, void
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "link: ETCP handshake OK with 0x%016llx — connection ESTABLISHED",
(unsigned long long)entry->node_id);
if (entry->main.timer) { uasync_cancel_timeout(entry->mgr->instance->ua, entry->main.timer); entry->main.timer = NULL; }
entry->main_connect_state = CM_TRY_OK; entry->conn_type = CONN_TYPE_DIRECT;
entry->main_connect_state = CM_TRY_OK;
entry->conn_type = (entry->main.phase == 2) ? CONN_TYPE_REVERSE : CONN_TYPE_DIRECT;
entry->state = CONN_MGR_STATE_CONNECTED;
cm_update_nodeinfo(entry); cm_deliver_event(entry, CONN_EVENT_UP);
{ struct cm_reverse_pending** pp = &entry->mgr->reverse_pending;
while (*pp) { if ((*pp)->entry == entry) { struct cm_reverse_pending* rp = *pp; *pp = rp->next; u_free(rp); break; } pp = &(*pp)->next; }
}
cm_update_nodeinfo(entry);
{ struct ETCP_CONN* c = node_conn_direct_get_conn(ncd_h);
if (c && entry->mgr && entry->mgr->group) topo_group_new_conn(entry->mgr->group, c); }
cm_deliver_event(entry, CONN_EVENT_UP);
break;
case NCD_EVENT_TIMEOUT:
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "link: ETCP connect TIMEOUT to 0x%016llx — entry=%p mgr=%p db_loaded=%d",
@ -267,20 +274,21 @@ void cm_ncd_callback(struct NODE_CONN_DIRECT* ncd_h, enum ncd_event ncd_ev, void
return;
}
entry->main_connect_state = CM_TRY_FAILED;
/* Entry без ожидающих handle (цель REVERSE) — эскалация бессмысленна, просто чистим */
if (!entry->handles) {
{ struct NODE_CONN_DIRECT* ncd = entry->ncd_handle; entry->ncd_handle = NULL;
if (ncd) node_conn_direct_force_close(ncd); }
cm_entry_cleanup(entry);
return;
}
{ struct TOPO_GROUP* g = entry->mgr->group;
struct TOPO_GROUP_NODE* t = topo_node_find_by_id(g, entry->node_id);
if (entry->local_scan_state != CM_TRY_OK) {
{ struct NODE_CONN_DIRECT* ncd = entry->ncd_handle; entry->ncd_handle = NULL;
if (ncd) node_conn_direct_force_close(ncd); }
if (cm_has_direct_ip(entry->mgr, g->local_node) && !(t && cm_has_direct_ip(entry->mgr, t)))
{ cm_start_phase_reverse(entry); return; }
cm_start_phase_indirect(entry); return;
}
if (cm_has_direct_ip(entry->mgr, g->local_node) && !(t && cm_has_direct_ip(entry->mgr, t)))
{ cm_start_phase_reverse(entry); return; } /* keep ncd_handle: target will reuse this conn for incoming INIT */
{ struct NODE_CONN_DIRECT* ncd = entry->ncd_handle; entry->ncd_handle = NULL;
if (ncd) node_conn_direct_force_close(ncd); }
cm_start_phase_indirect(entry); return;
}
if (entry->local_scan_state == CM_TRY_FAILED) cm_deliver_event(entry, CONN_EVENT_TIMEOUT);
{ struct NODE_CONN_DIRECT* ncd = entry->ncd_handle; entry->ncd_handle = NULL;
if (ncd) node_conn_direct_force_close(ncd); }
break;
case NCD_EVENT_DOWN:
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "link: ETCP link DROPPED to 0x%016llx",
(unsigned long long)entry->node_id);
@ -341,7 +349,7 @@ void conn_mgr_destroy(struct CONN_MGR* mgr) {
/* ═══════ публичное API ═══════ */
int cm_open(struct CONN_MGR* mgr, uint64_t node_id, uint32_t idle_timeout_ms,
int cm_open(struct CONN_MGR* mgr, uint64_t node_id,
conn_mgr_cb_t cb, void* cb_arg, struct CONN_MGR_HANDLE** out_handle) {
if (out_handle) *out_handle = NULL;
if (!mgr) { if (cb) cb(NULL, node_id, 0, CONN_EVENT_TIMEOUT, cb_arg); return -1; }
@ -395,7 +403,7 @@ int cm_open(struct CONN_MGR* mgr, uint64_t node_id, uint32_t idle_timeout_ms,
entry = cm_ensure_entry(mgr, node_id);
if (!entry) return -1;
entry->state = CONN_MGR_STATE_CONNECTED; entry->conn_type = CONN_TYPE_DIRECT;
entry->idle_timeout_ms = idle_timeout_ms; entry->last_traffic_tb = get_time_tb();
entry->last_traffic_tb = get_time_tb();
{ struct TOPO_NODE* ni = topo_node_registry_find(mgr->instance->topo_groups, node_id);
struct TOPO_NODE ni_e; memset(&ni_e, 0, sizeof(ni_e)); ni_e.node_id = node_id;
if (ni) memcpy(ni_e.public_key, ni->public_key, SC_PUBKEY_SIZE);
@ -411,7 +419,7 @@ int cm_open(struct CONN_MGR* mgr, uint64_t node_id, uint32_t idle_timeout_ms,
entry->state = CONN_MGR_STATE_CONNECTING;
struct CONN_MGR_HANDLE* h = cm_handle_new(entry, gid, cb, cb_arg);
if (!h) { cm_entry_cleanup(entry); return -1; } if (out_handle) *out_handle = h;
entry->idle_timeout_ms = idle_timeout_ms; entry->last_traffic_tb = get_time_tb();
entry->last_traffic_tb = get_time_tb();
entry->local_scan_state = CM_TRY_NONE; entry->main_connect_state = CM_TRY_NONE;
entry->db_loaded = (uint8_t)loaded_from_db;
if (cm_has_local_addr(mgr, group->local_node) && cm_has_local_addr(mgr, target))
@ -441,7 +449,7 @@ void conn_mgr_close(struct CONN_MGR_HANDLE* h) {
struct ll_entry* qe = queue_entry_new(0);
if (qe) { qe->dgram = u_malloc(1+sizeof(pkt)); qe->dgram[0] = ETCP_RT_ID_CONN_MGR; memcpy(qe->dgram+1, &pkt, sizeof(pkt)); qe->len = (uint16_t)(1+sizeof(pkt));
etcp_route_send(entry->mgr->instance, entry->mgr->group->group_id, entry->node_id, qe, 1, 0); }
/* NCD CLOSE — транспортный уровень (отправит CLOSE/KEEP_ALIVE через node_conn_direct) */
/* Транспортный уровень рвётся ниже в cm_entry_cleanup (node_conn_direct_force_close) */
}
cm_entry_cleanup(entry); u_free(h);
}
@ -457,7 +465,7 @@ int conn_mgr_send(struct CONN_MGR_HANDLE* h, struct ll_entry* e) {
if (!entry || entry->state != CONN_MGR_STATE_CONNECTED) { queue_entry_free(e); return -1; }
entry->last_traffic_tb = get_time_tb();
struct ETCP_ROUTER_CONN* r = etcp_router_conn_get(entry->mgr->instance, entry->mgr->group->group_id, entry->node_id, ETCP_RT_ID_CONN_MGR);
if (!r) { queue_entry_free(e); return -1; }
if (!r) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "conn_mgr: no route to 0x%016llx, dropping %u bytes", (unsigned long long)entry->node_id, (unsigned)e->len); queue_entry_free(e); return -1; }
int ret = etcp_router_conn_send(r, e->dgram, e->len, 0); queue_entry_free(e); return ret;
}
@ -466,7 +474,7 @@ int conn_mgr_send(struct CONN_MGR_HANDLE* h, struct ll_entry* e) {
static void cm_ping_cb(int success, uint16_t rtt, void* arg, uint64_t nonce, const uint8_t* resp_data, size_t resp_data_len);
static void cm_ping_cb_impl(int success, uint16_t rtt, void* arg, uint64_t nonce, const uint8_t* resp_data, size_t resp_data_len);
struct cm_ping_ctx { struct CONN_MGR_ENTRY* entry; struct sockaddr_storage addr; struct ETCP_SOCKET* sock; uint8_t phase, attempt; };
struct cm_ping_ctx { struct CONN_MGR* mgr; uint64_t node_id; struct sockaddr_storage addr; struct ETCP_SOCKET* sock; uint8_t phase, attempt; };
/* LAN broadcast ping для обнаружения узла в локальной сети. Если ответит —
* соединение готово быстрее чем через интернет, без NAT. Запускается
@ -486,7 +494,7 @@ void cm_start_local_scan(struct CONN_MGR_ENTRY* entry) {
if ((s->type == CFG_SERVER_TYPE_PRIVATE || s->type == CFG_SERVER_TYPE_LOCAL) && s->local_addr.ss_family == AF_INET) {
struct cm_ping_ctx* ctx = u_calloc(1, sizeof(struct cm_ping_ctx));
if (!ctx) continue;
ctx->entry = entry; ctx->sock = s; ctx->phase = 0; ctx->attempt = 0;
ctx->mgr = entry->mgr; ctx->node_id = entry->node_id; ctx->sock = s; ctx->phase = 0; ctx->attempt = 0;
struct sockaddr_in sin; memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET;
memcpy(&sin.sin_addr.s_addr, a->addr, 4); sin.sin_port = htons(a->port);
memcpy(&ctx->addr, &sin, sizeof(sin));
@ -504,14 +512,15 @@ static void cm_ping_cb(int success, uint16_t rtt, void* arg, uint64_t nonce, con
{ cm_ping_cb_impl(success, rtt, arg, nonce, resp_data, resp_data_len); }
static void cm_ping_cb_impl(int success, uint16_t rtt, void* arg, uint64_t nonce, const uint8_t* resp_data, size_t resp_data_len) {
struct cm_ping_ctx* ctx = (struct cm_ping_ctx*)arg;
struct CONN_MGR_ENTRY* entry = ctx->entry; (void)nonce; (void)resp_data; (void)resp_data_len;
struct cm_ping_ctx* ctx = (struct cm_ping_ctx*)arg; (void)nonce; (void)resp_data; (void)resp_data_len;
struct CONN_MGR_ENTRY* entry = cm_find_entry(ctx->mgr, ctx->node_id);
if (!entry || entry->state == CONN_MGR_STATE_DISCONNECTED) { u_free(ctx); return; }
if (entry->main_connect_state == CM_TRY_OK) { u_free(ctx); return; }
if (success && ctx->phase == 0) {
entry->local_scan_state = CM_TRY_OK; entry->conn_type = CONN_TYPE_DIRECT; entry->state = CONN_MGR_STATE_CONNECTED;
if (entry->main.timer) { uasync_cancel_timeout(entry->mgr->instance->ua, entry->main.timer); entry->main.timer = NULL; }
entry->local_scan_state = CM_TRY_OK;
topo_node_ping_update_rtt(entry->mgr->instance->topo_groups, entry->node_id, rtt);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: local scan OK for 0x%016llx rtt=%u", (unsigned long long)entry->node_id, rtt);
cm_update_nodeinfo(entry); cm_deliver_event(entry, CONN_EVENT_UP); u_free(ctx); return;
u_free(ctx); return;
}
if (++ctx->attempt < CONN_MGR_LOCAL_SCAN_ATTEMPTS && ctx->phase == 0) {
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(entry->mgr->group, entry->node_id);
@ -617,7 +626,7 @@ static void cm_direct_add_links(struct ETCP_CONN* conn, struct CONN_MGR_ENTRY* e
* NCD управляет таймаутом (2с). При успехе — cm_ncd_callback доставит UP,
* при таймауте — переключится на REVERSE или INDIRECT. */
void cm_start_phase_direct(struct CONN_MGR_ENTRY* entry) {
if (entry->main_connect_state == CM_TRY_OK || entry->local_scan_state == CM_TRY_OK) return;
if (entry->main_connect_state == CM_TRY_OK) return;
struct TOPO_GROUP* group = entry->mgr->group;
struct TOPO_GROUP_NODE* target = topo_node_find_by_id(group, entry->node_id);
struct TOPO_NODE* ni = target ? topo_node_registry_find(entry->mgr->instance->topo_groups, target->node_id) : NULL;
@ -634,26 +643,13 @@ void cm_start_phase_direct(struct CONN_MGR_ENTRY* entry) {
/* ═══════ REVERSE фаза ═══════ */
void cm_reverse_init_cb(struct ETCP_CONN* conn, int event, void* arg) { (void)event;
struct cm_reverse_pending* rp = (struct cm_reverse_pending*)arg;
if (!rp) return; if (!conn) { u_free(rp); return; }
struct CONN_MGR_ENTRY* entry = rp->entry;
if (!entry) { DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: reverse ready for 0x%016llx (no entry, freeing)", (unsigned long long)conn->peer_node_id); u_free(rp); return; }
if (conn->peer_node_id != entry->node_id) { u_free(rp); return; }
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: reverse ready for 0x%016llx", (unsigned long long)entry->node_id);
if (entry->main.timer) { uasync_cancel_timeout(entry->mgr->instance->ua, entry->main.timer); entry->main.timer = NULL; }
{ struct cm_reverse_pending** pp = &entry->mgr->reverse_pending;
while (*pp) { if (*pp == rp) { *pp = rp->next; break; } pp = &(*pp)->next; }
}
entry->main_connect_state = CM_TRY_OK; entry->conn_type = CONN_TYPE_REVERSE; entry->state = CONN_MGR_STATE_CONNECTED;
cm_update_nodeinfo(entry); cm_deliver_event(entry, CONN_EVENT_UP); u_free(rp);
}
void cm_reverse_timeout_cb(void* arg) {
struct CONN_MGR_ENTRY* entry = (struct CONN_MGR_ENTRY*)arg; entry->main.timer = NULL;
{ struct cm_reverse_pending** pp = &entry->mgr->reverse_pending;
while (*pp) { if ((*pp)->entry == entry) { struct cm_reverse_pending* rp = *pp; *pp = rp->next; u_free(rp); break; } pp = &(*pp)->next; }
}
{ struct NODE_CONN_DIRECT* ncd = entry->ncd_handle; entry->ncd_handle = NULL;
if (ncd) node_conn_direct_force_close(ncd); }
cm_start_phase_indirect(entry);
}
@ -697,31 +693,52 @@ void cm_start_phase_reverse(struct CONN_MGR_ENTRY* entry) {
/* ═══════ обработка DIRECT_REQ (REVERSE входящие) ═══════ */
/* Принимающая сторона REVERSE: получили DIRECT_REQ от инициатора — открываем
* NCD-соединение к нему, добавляем линки по адресам из запроса. При INIT
* вызывается cm_reverse_init_cb. */
* NCD-соединение к нему, добавляем линки по адресам из запроса с NAT-фильтрацией
* (sock_match). Целевая сторона может быть за NAT — поэтому линк создаём от любого
* совместимого сокета (не только public), чтобы её исходящий INIT достиг инициатора. */
void cm_handle_direct_req(uint64_t src, struct CONN_MGR* mgr, const uint8_t* data, size_t len) {
if (!mgr) return;
struct CM_DIRECT_REQ* req = (struct CM_DIRECT_REQ*)data;
if (len < CM_DIRECT_REQ_H_SIZE + (size_t)req->addr_count * 8) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: DIRECT_REQ too short"); return; }
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: got DIRECT_REQ from 0x%016llx with %u addrs",
(unsigned long long)src, req->addr_count);
struct TOPO_NODE* src_ni = topo_node_registry_find(mgr->instance->topo_groups, src); if (!src_ni) return;
struct TOPO_NODE* src_ni = topo_node_registry_find(mgr->instance->topo_groups, src);
if (!src_ni) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: DIRECT_REQ from unknown node 0x%016llx", (unsigned long long)src); return; }
/* Целевая сторона REVERSE. Если уже есть entry к инициатору (своё open идёт
* или связь уже есть) — ничего не делаем: обратный INIT сматчится с нашим conn. */
struct CONN_MGR_ENTRY* entry = cm_find_entry(mgr, src);
if (entry && entry->state == CONN_MGR_STATE_CONNECTED) return;
if (entry && entry->state == CONN_MGR_STATE_CONNECTING) return;
if (!entry) entry = cm_ensure_entry(mgr, src);
if (!entry) return;
entry->state = CONN_MGR_STATE_CONNECTING; entry->main.phase = 2;
entry->main_connect_state = CM_TRY_PENDING; entry->db_loaded = 0; entry->conn_type = CONN_TYPE_NONE;
struct TOPO_NODE nbuf; memset(&nbuf,0,sizeof(nbuf)); nbuf.node_id=src; memcpy(nbuf.public_key,src_ni->public_key,SC_PUBKEY_SIZE);
struct NODE_CONN_DIRECT* ncd_h = NULL;
if (node_conn_direct_open_node(mgr->instance, src, NULL, NULL, &ncd_h, &nbuf, NULL) == NCD_ERR) return;
struct ETCP_CONN* newc = node_conn_direct_get_conn(ncd_h); if (!newc) { node_conn_direct_close(ncd_h); return; }
struct cm_reverse_pending* rp = u_calloc(1,sizeof(*rp));
if (!rp) { node_conn_direct_close(ncd_h); return; }
rp->request_id=req->request_id; rp->next=mgr->reverse_pending; mgr->reverse_pending=rp;
if (node_conn_direct_open_node(mgr->instance, src, cm_ncd_callback, entry, &entry->ncd_handle, &nbuf, NULL) == NCD_ERR) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: DIRECT_REQ open_node failed for 0x%016llx", (unsigned long long)src);
cm_entry_cleanup(entry); return;
}
struct ETCP_CONN* newc = node_conn_direct_get_conn(entry->ncd_handle); if (!newc) { node_conn_direct_close(entry->ncd_handle); entry->ncd_handle = NULL; return; }
struct sock_view views[CM_MAX_SOCKET_VIEWS];
int nviews = sock_collect_views(mgr->instance, views, CM_MAX_SOCKET_VIEWS);
uint8_t* p = (uint8_t*)data + CM_DIRECT_REQ_H_SIZE;
for (uint8_t i=0;i<req->addr_count;i++) {
uint8_t type=*p++; (void)type; uint8_t ip[4]; memcpy(ip,p,4);p+=4; uint16_t port; memcpy(&port,p,2);p+=2; uint8_t sid=*p++; (void)sid;
struct ETCP_SOCKET* s = mgr->instance->etcp_sockets;
while (s) { if ((s->type==CFG_SERVER_TYPE_PUBLIC||s->type==CFG_SERVER_TYPE_UNKNOWN)&&s->local_addr.ss_family==AF_INET)
{ cm_add_v4_link(newc, ip, port, s); break; } s=s->next; }
uint8_t type=*p++; uint8_t ip[4]; memcpy(ip,p,4);p+=4; uint16_t port; memcpy(&port,p,2);p+=2; uint8_t sid=*p++;
uint32_t tip; memcpy(&tip, ip, 4);
uint8_t tcfg = CFG_SERVER_TYPE_UNKNOWN, tnat = NAT_TYPE_UNKNOWN;
for (const struct TOPO_SOCKMETA4* m = src_ni->v4_sock_meta; m; m = m->next)
if (m->id == sid) { tcfg = m->config_type; tnat = m->nat_type; break; }
for (int vi = 0; vi < nviews; vi++) {
const struct sock_view* v = &views[vi];
if (v->is_tcp || !sock_match(v, type, tcfg, tnat, tip)) continue;
cm_add_v4_link(newc, ip, port, v->sock);
break;
}
}
etcp_conn_add_cbk(newc, cm_reverse_init_cb, rp, ETCP_CBK_EVENT_INIT);
}
/* ═══════ disconnect / recv ═══════ */
@ -738,7 +755,9 @@ void conn_mgr_router_recv_handler(struct ETCP_CONN* conn, struct ll_entry* entry
DEBUG_ERROR(DEBUG_CATEGORY_ROUTING, "conn_mgr_router_recv: conn=%p", (void*)conn);
return;
}
uint8_t* d = entry->dgram + ROUTER_SVC_PAYLOAD_OFF; size_t len = entry->len - ROUTER_SVC_PAYLOAD_OFF; uint8_t sub = d[1];
uint8_t* d = entry->dgram + ROUTER_SVC_PAYLOAD_OFF; size_t len = entry->len - ROUTER_SVC_PAYLOAD_OFF;
if (d[0] != ETCP_RT_ID_CONN_MGR) return;
uint8_t sub = d[1];
/* conn_mgr — per-group: резолвим группу из заголовка доставки, а не instance->conn_mgr */
uint64_t group_id; memcpy(&group_id, entry->dgram + ROUTER_SVC_GROUP_OFF, 8);
uint64_t src; memcpy(&src, entry->dgram + ROUTER_SVC_SRC_OFF, 8);
@ -750,7 +769,7 @@ void conn_mgr_router_recv_handler(struct ETCP_CONN* conn, struct ll_entry* entry
case CM_SUBCMD_DIRECT_RESP: break;
case CM_SUBCMD_INTERM_EXCHANGE_REQ: if (len >= CM_EXCHANGE_REQ_SIZE) cm_handle_interm_exchange_req(src, mgr, (struct CM_EXCHANGE_REQ*)d); break;
case CM_SUBCMD_INTERM_EXCHANGE_RESP: if (mgr) cm_handle_interm_exchange_resp(mgr, d, len); break;
case CM_SUBCMD_INTERM_SELECTED: if (mgr) cm_handle_interm_selected(mgr, d, len); break;
case CM_SUBCMD_INTERM_SELECTED: if (mgr) cm_handle_interm_selected(mgr, src, d, len); break;
case CM_SUBCMD_DISCONNECT: if (len >= CM_DISCONNECT_SIZE && mgr) cm_handle_disconnect(mgr, ((struct CM_DISCONNECT*)d)->node_id); break;
}
}
@ -772,5 +791,5 @@ int conn_mgr_open(struct UTUN_INSTANCE* inst,
struct CONN_MGR* mgr = group->conn_mgr;
if (!mgr) return -1;
return cm_open(mgr, node_id, 0, cb, cb_arg, out_handle);
return cm_open(mgr, node_id, cb, cb_arg, out_handle);
}

221
src/routing_layer/conn_mgr_doc.md

@ -3,134 +3,213 @@
## 1. Назначение
Трёхфазный асинхронный менеджер подключения к удалённым нодам. Отвечает за установку ETCP-соединения с выбором оптимального способа:
1. **DIRECT** — прямое INIT-рукопожатие со всеми известными IPv4-адресами цели (с проверкой NAT-совместимости)
2. **REVERSE** — если у нас прямой IP, а цель за NAT: отправляем наши адреса через BGP, цель подключается к нам
1. **DIRECT** — прямое INIT-рукопожатие со всеми известными адресами цели (IPv4/IPv6, TCP/UDP, с NAT-фильтрацией через `sock_match`)
2. **REVERSE** — если у нас прямой IP, а цель за NAT: отправляем наши адреса через BGP (`DIRECT_REQ`), цель подключается к нам
3. **INDIRECT** — обмениваемся через BGP списками кандидатов-посредников, выбираем общих по минимальной сумме RTT, трафик идёт через `etcp_router`
Также выполняет фоновые задачи: периодический ping всех BGP-нод, поддержание списка лучших посредников, idle-таймаут для неактивных соединений, локальное сканирование сети.
Также выполняет фоновые задачи: периодический ping всех BGP-нод, поддержание списка лучших посредников, idle-таймаут для неактивных соединений, локальное сканирование сети (LAN broadcast ping).
Типы соединений (итоговый результат):
- `CONN_TYPE_DIRECT` (1) — прямое соединение
- `CONN_TYPE_REVERSE` (2) — цель подключилась к нам
- `CONN_TYPE_INDIRECT` (3) — через посредника
Один `CONN_MGR` обслуживает одну `TOPO_GROUP` (создаётся в `conn_mgr_init(group)`). Внутри — кеш `entries` (по одной записи на целевой узел) и список `handle`'ов (по одному на каждый вызов `conn_mgr_open`).
## 2. Как пользоваться
### Инициализация
```c
struct CONN_MGR* mgr = conn_mgr_init(instance);
// Регистрирует обработчик в etcp_router для ETCP_RT_ID_CONN_MGR (0x11)
struct CONN_MGR* mgr = conn_mgr_init(group);
// Регистрируется в etcp_router через conn_mgr_router_recv_handler (ETCP_RT_ID_CONN_MGR)
// Запускает фоновые таймеры: bg_ping, candidate_ping
```
### Подключение
### Подключение (handle-based)
```c
static void my_connect_cb(int result, uint64_t node_id, void* arg) {
if (result == CONN_MGR_OK) { /* подключено */ }
else if (result == CONN_MGR_ERR_TIMEOUT) { /* таймаут */ }
else if (result == CONN_MGR_ERR_UNREACHABLE) { /* недостижима */ }
static void my_cb(struct CONN_MGR_HANDLE* h, uint64_t node_id, uint64_t group_id,
enum conn_mgr_event event, void* arg) {
if (event == CONN_EVENT_UP) { /* связь есть */ }
else if (event == CONN_EVENT_DOWN) { /* временный обрыв, сам восстановится */ }
else if (event == CONN_EVENT_TIMEOUT) { /* не удалось, закройте handle */ }
}
conn_mgr_connect_node(mgr, target_node_id, idle_timeout_ms, my_connect_cb, arg);
struct CONN_MGR_HANDLE* h = NULL;
conn_mgr_open(inst, group_id, node_id, my_cb, arg, &h);
```
### Отправка данных
```c
struct ll_entry* entry = queue_entry_new(0);
entry->dgram = data; entry->len = len;
conn_mgr_send(mgr, node_id, entry); // entry освобождается внутри
struct ll_entry* e = queue_entry_new(0);
e->dgram = data; e->len = len;
conn_mgr_send(h, e); // забирает владение e, освобождает внутри
struct ETCP_CONN* conn = conn_mgr_get_conn(h); // NULL для INDIRECT
```
### Закрытие
```c
conn_mgr_close(h); // убирает ОДИН handle; последний рвёт соединение
```
### Кандидаты-посредники
Внешний код (обычно `route_connectivity`) вызывает `conn_mgr_update_best_candidates()` для заполнения списка лучших посредников (топ-3 по RTT). Список используется фазой INDIRECT.
`cm_candidate_ping_timer_cb` (раз в ~2с) пополняет `best_candidates` (топ-3 по RTT) из узлов группы с активным соединением и пингует/probe'ит их. Список используется фазой INDIRECT. Внешний код (`route_connectivity`) лишь заполняет `nq->connectivity`, откуда берутся свежие RTT.
### Нюансы
- Повторный `connect_node` для уже подключённой ноды вернёт `CONN_MGR_ERR_ALREADY_CONNECTED`
- Для ноды в состоянии CONNECTING новый вызов добавит ещё один callback в список ожидания
- Для INDIRECT-соединений `conn_mgr_send` использует round-robin по выбранным посредникам
- Пока существует активный ETCP-линк к ноде, немедленно возвращает CONNECTED
- Несколько `conn_mgr_open` на один узел → несколько handle на одном entry; все получают события
- Повторный `open` для уже подключённой ноды вернёт handle сразу, UP доставляется отложенно (`uasync_call_soon`)
- Для `open` в состоянии CONNECTING новый handle добавляется в список ожидания
- Для INDIRECT-соединений `conn_mgr_get_conn` возвращает `NULL`, `conn_mgr_send` идёт через `etcp_router` (первый посредник с живым conn)
- Пока есть активный ETCP-линк к ноде — соединение считается DIRECT
## 3. API
### События (enum conn_mgr_event)
| Событие | Значение | Смысл |
|---------|----------|-------|
| `CONN_EVENT_UP` | 0 | Связь есть, можно отправлять данные |
| `CONN_EVENT_DOWN` | 1 | Временный обрыв, handle жив, само восстановится |
| `CONN_EVENT_TIMEOUT` | 2 | Подключение не удалось, handle жив, закройте сами |
### Коллбэк
```c
typedef void (*conn_mgr_cb_t)(struct CONN_MGR_HANDLE* h,
uint64_t node_id, uint64_t group_id,
enum conn_mgr_event event, void* arg);
```
`h` может быть `NULL` (ошибки до создания handle). `node_id`/`group_id` всегда валидны.
### Корневые структуры
| Структура | Описание |
|-----------|----------|
| `CONN_MGR` | Корневая структура менеджера. Хранит массив `entries[]`, топ-3 `best_candidates`, таймеры (`bg_ping_timer`, `candidate_ping_timer`), списки `reverse_pending`/`exchange_pending`, `next_request_id` |
| `CONN_MGR_ENTRY` | Запись о подключении к одной ноде: `state`, `conn_type`, `idle_timeout_ms`, список коллбэков `cb_list`, массив `intermediaries[]`, состояния `local_scan_state`/`main_connect_state`, текущая `main.phase` |
| `CONN_MGR` | Корневая структура менеджера. `instance`, `group`, очередь `entries`, топ-3 `best_candidates`, таймеры (`bg_ping_timer`, `candidate_ping_timer`), `bg_ping_cursor`, списки `reverse_pending`/`exchange_pending`, `next_request_id` |
| `CONN_MGR_ENTRY` | Запись о подключении к одной ноде: `state`, `conn_type`, `last_traffic_tb`, список `handles`, `ncd_handle`, `intermediaries[3]`/`rr_idx`, `local_scan_state`/`main_connect_state`, `db_loaded`, `main.phase` |
| `CONN_MGR_HANDLE` | Handle для одного вызова `conn_mgr_open`: `node_id`, `group_id`, `cb`, `cb_arg`, ссылка на `entry`, `next`, `deliver_up_token` (отложенный UP) |
| `CONN_MGR_CANDIDATE` | Пара `{node_id, rtt}` — запись о кандидате-посреднике (packed) |
| `cm_cb_node` | Узел связного списка коллбэков для асинхронного результата подключения |
| `cm_reverse_pending` | Запись об ожидании REVERSE-подключения (связный список) |
| `cm_exchange_pending` | Запись об ожидании INDIRECT-обмена, включает кешированный ответ и таймер probe-повторов |
| `cm_exchange_pending` | Запись об ожидании INDIRECT-обмена: кешированный `EXCHANGE_RESP`, флаги `received`/`probed`, таймер probe-повторов |
### Жизненный цикл
| Функция | Описание |
|---------|----------|
| `conn_mgr_init(instance)` | Выделяет CONN_MGR, регистрирует `ETCP_RT_ID_CONN_MGR` в маршрутизаторе, запускает bg_ping и candidate_ping таймеры |
| `conn_mgr_init(group)` | Выделяет CONN_MGR, создаёт очередь entries, запускает bg_ping и candidate_ping таймеры |
| `conn_mgr_destroy(mgr)` | Отменяет таймеры, освобождает все entries, exchange_pending, reverse_pending |
| `conn_mgr_router_recv_handler(conn, entry)` | Обработчик входящих пакетов `ETCP_RT_ID_CONN_MGR` из etcp_router; диспетчеризация по subcmd |
### Подключение / отключение
| Функция | Описание |
|---------|----------|
| `conn_mgr_connect_node(mgr, id, timeout, cb, arg)` | Запускает 3-фазное подключение. Если уже подключена — сразу callback с OK. Если в процессе — добавляет callback в список |
| `conn_mgr_disconnect_node(mgr, id)` | Шлёт DISCONNECT пакет, сбрасывает `conn_mgr_type` в `NODEINFO_Q`, удаляет entry |
| `conn_mgr_send(mgr, id, entry)` | Отправляет данные подключённой ноде через `etcp_router`; для INDIRECT — round-robin по посредникам; обновляет `last_traffic_tb` |
| `conn_mgr_set_idle_timeout(mgr, id, ms)` | Устанавливает/снимает idle-таймер для ноды. При неактивности дольше `ms` — автоотключение |
| `conn_mgr_get_status(mgr, id, &state, &type)` | Возвращает текущее состояние и тип соединения |
| `conn_mgr_open(inst, group_id, node_id, cb, cb_arg, &h)` | Единая точка входа. Резолвит группу, запускает 3-фазное подключение. Если уже подключена — сразу отложенный UP. Если в процессе — добавляет handle в список |
| `conn_mgr_close(h)` | Убирает один handle. Последний handle при DIRECT/REVERSE шлёт DISCONNECT и чистит entry |
| `conn_mgr_send(h, entry)` | Отправляет данные через `etcp_router`; для INDIRECT — через первого посредника с живым conn; обновляет `last_traffic_tb` |
| `conn_mgr_get_conn(h)` | ETCP_CONN для DIRECT/REVERSE, NULL для INDIRECT |
### Кандидаты-посредники
| Функция | Описание |
|---------|----------|
| `conn_mgr_update_best_candidates(mgr, id, rtt)` | Вставляет/обновляет кандидата в отсортированном топ-3 по RTT. Худшие кандидаты вытесняются |
### Протокольные подкоманды (ETCP_RT_ID_CONN_MGR = 0x11)
| Константа | Описание |
|-----------|----------|
| `CONN_MGR_SUBCMD_DIRECT_REQ` (0x01) | Фаза 2: "я за NAT, подключись ко мне" — наши адреса |
| `CONN_MGR_SUBCMD_DIRECT_RESP` (0x02) | Ответ на DIRECT_REQ (зарезервирован) |
| `CONN_MGR_SUBCMD_INTERM_EXCHANGE_REQ` (0x03) | Фаза 3 инициатор: наш топ-4 кандидатов |
| `CONN_MGR_SUBCMD_INTERM_EXCHANGE_RESP` (0x04) | Фаза 3 ответ: свои + чужие кандидаты с RTT с обеих сторон |
| `CONN_MGR_SUBCMD_INTERM_SELECTED` (0x05) | Фаза 3 финал: выбранные посредники |
| `CONN_MGR_SUBCMD_DISCONNECT` (0x06) | Уведомление о разрыве |
| `conn_mgr_update_best_candidates(mgr, node_id, rtt)` | Вставляет/обновляет кандидата в отсортированном топ-3 по RTT. Худшие кандидаты вытесняются |
## 4. Алгоритм подключения
> Нумерация фаз в коде ведётся через `entry->main.phase`: DIRECT = 0, REVERSE = 2, INDIRECT = 3 (значение 1 не используется).
### `conn_mgr_open` / `cm_open` (core)
1. Отклоняет `node_id == 0`.
2. Ищет целевой `TOPO_GROUP_NODE`; если нет — подгружает из SQLite (`db_loaded = 1`).
3. Уже `CONNECTED` → новый handle + отложенный UP.
4. Уже `CONNECTING` → добавляет handle в список ожидания.
5. Есть живой ETCP-линк → помечает `CONNECTED` (DIRECT), открывает NCD, отложенный UP.
6. Иначе: новый entry → `CONNECTING`. При локальных адресах у обеих сторон — запускает локальный скан параллельно, затем фазу DIRECT.
### Фаза 1 — DIRECT (`cm_start_phase_direct`)
- `node_conn_direct_open_node` открывает NCD-соединение (пустой `ni`, коллбэк `cm_ncd_callback`).
- `cm_direct_add_links` вручную добавляет линки: перебор адресов в приоритете `NAT > INTERFACE`, v4+v6, TCP+UDP, сопоставление сокетов через `sock_match`.
- Хоть один линк → ждём INIT handshake. Ноль линков → закрыть NCD, перейти к REVERSE/INDIRECT (для `db_loaded` — cleanup + TIMEOUT).
- NCD держит таймаут 2с. Результаты через `cm_ncd_callback`:
- `NCD_EVENT_UP` → `conn_type = (main.phase==2) ? REVERSE : DIRECT`, `CONNECTED`, регистрация conn в группе (`topo_group_new_conn`) и доставка UP.
- `NCD_EVENT_TIMEOUT` → выбор следующей фазы.
- `NCD_EVENT_DOWN` → доставка DOWN (handle жив).
### Локальный скан (`cm_start_local_scan`)
Параллельно DIRECT. Если у обеих сторон приватные/локальные адреса — LAN ping (`etcp_send_ping_to_socket`, до 3 попыток по 100мс). Успех только обновляет RTT (`topo_node_ping_update_rtt`) и помечает узел достижимым локально (`local_scan_state`); соединение при этом устанавливает DIRECT-фаза обычным путём.
### Фаза 2 — REVERSE (`cm_start_phase_reverse`)
Запускается при `NCD_EVENT_TIMEOUT`, если у нас прямой IP, а у цели нет (`cm_has_direct_ip(local) && !cm_has_direct_ip(target)`):
1. Собираем наши публичные/EIM/direct адреса.
2. Шлём `DIRECT_REQ` (наши адреса) цели через BGP-маршрут `etcp_route_send`, таймер 15с. При таймауте NCD-соединение фазы DIRECT **не закрывается** — оно остаётся в `instance->connections`, чтобы обратный INIT цели переиспользовал его. (В пути «ноль линков» из `cm_direct_add_links` NCD закрывается — линков для INIT-matching всё равно нет.)
3. Цель (`cm_handle_direct_req`) открывает NCD к нам и добавляет линки к нашим адресам с NAT-фильтрацией (`sock_match`) — линк создаётся от любого совместимого сокета (цель может быть за NAT).
4. Цель дозванивается → её INIT матчится с нашим conn по `peer_node_id` → conn поднимается → `cm_ncd_callback(UP)` с `main.phase==2` → `conn_type=REVERSE`, `CONNECTED`, UP.
5. Таймаут → `cm_reverse_timeout_cb` (закрывает NCD) → INDIRECT.
### Фаза 3 — INDIRECT (indirect.c)
1. `cm_start_phase_indirect`: пустой `best_candidates` → TIMEOUT. Иначе шлёт `INTERM_EXCHANGE_REQ` со своим топ-4 кандидатов, таймер 15с.
2. Цель (`cm_handle_interm_exchange_req`) пробегает RTT до кандидатов инициатора (probe при протухших), шлёт `EXCHANGE_RESP` со своими кандидатами + замерами RTT.
3. Цель (`cm_handle_interm_exchange_resp`) кеширует ответ, проверяет свежесть RTT до своих кандидатов; протухшие — probe + ретрай 200мс; свежие → `cm_compute_intermediaries`.
4. `cm_compute_intermediaries` суммирует RTT по общим узлам, сортирует, берёт топ-3, шлёт `INTERM_SELECTED`, помечает `conn_type=INDIRECT`, `CONNECTED`, UP.
5. Цель (`cm_handle_interm_selected`) находит свой entry по `src`, сохраняет посредников, помечает `CONNECTED`, `cm_update_nodeinfo` + UP (если ещё не было).
Трафик идёт через `etcp_router` (первый посредник с живым conn; `rr_idx` зарезервирован, не используется).
## 5. Протокол (ETCP_RT_ID_CONN_MGR)
### Подкоманды
| Константа | Значение | Описание |
|-----------|----------|----------|
| `CM_SUBCMD_DIRECT_REQ` | 0x01 | Фаза 2: "я за NAT, подключись ко мне" — наши адреса |
| `CM_SUBCMD_DIRECT_RESP` | 0x02 | Ответ на DIRECT_REQ (зарезервирован) |
| `CM_SUBCMD_INTERM_EXCHANGE_REQ` | 0x03 | Фаза 3 инициатор: наш топ-4 кандидатов |
| `CM_SUBCMD_INTERM_EXCHANGE_RESP` | 0x04 | Фаза 3 ответ: свои + чужие кандидаты с RTT обеих сторон |
| `CM_SUBCMD_INTERM_SELECTED` | 0x05 | Фаза 3 финал: выбранные посредники |
| `CM_SUBCMD_DISCONNECT` | 0x06 | Уведомление о разрыве |
### Пакеты протокола (все packed)
| Структура | Поля |
|-----------|------|
| `CONN_MGR_DIRECT_REQ` | `cmd, subcmd, request_id, addr_count` + массив адресов `{type, ip[4], port, sock_id}` |
| `CONN_MGR_DIRECT_RESP` | `cmd, subcmd, request_id, accepted` |
| `CONN_MGR_INTERM_EXCHANGE_REQ` | `cmd, subcmd, request_id, candidate_count, candidates[4]` |
| `CONN_MGR_INTERM_EXCHANGE_RESP` | `cmd, subcmd, request_id, my_count, your_count, my_candidates[4], your_candidates[4]` |
| `CONN_MGR_INTERM_SELECTED` | `cmd, subcmd, request_id, count, selected[3]` |
| `CONN_MGR_DISCONNECT` | `cmd, subcmd, node_id` |
### Коды возврата
| Код | Значение |
|-----|----------|
| `CONN_MGR_OK` (0) | Успех |
| `CONN_MGR_ERR_NOT_FOUND` (-1) | Нода не найдена в BGP |
| `CONN_MGR_ERR_NO_ADDRESSES` (-2) | Нет адресов для подключения |
| `CONN_MGR_ERR_TIMEOUT` (-3) | Таймаут |
| `CONN_MGR_ERR_UNREACHABLE` (-4) | Недостижима (все фазы) |
| `CONN_MGR_ERR_ALREADY_CONNECTED` (-6) | Уже подключены |
| `CONN_MGR_ERR_INTERNAL` (-7) | Внутренняя ошибка |
### Ключевые константы
| `CM_DIRECT_REQ` | `cmd, subcmd, request_id, addr_count` + массив адресов `{type, ip[4], port, sock_id}` |
| `CM_DIRECT_RESP` | `cmd, subcmd, request_id, accepted` |
| `CM_EXCHANGE_REQ` | `cmd, subcmd, request_id, candidate_count, candidates[4]` |
| `CM_EXCHANGE_RESP` | `cmd, subcmd, request_id, my_count, your_count, my_candidates[4], your_candidates[4]` |
| `CM_INTERM_SEL` | `cmd, subcmd, request_id, count, selected[3]` |
| `CM_DISCONNECT` | `cmd, subcmd, node_id` |
## 6. Фоновые процессы (monitor.c)
| Таймер | Интервал | Описание |
|--------|----------|----------|
| `bg_ping` | ~100мс (один узел за тик, цикл ≥10с) | Пинг всех BGP-нод без активного соединения, поддерживает свежие RTT для best_candidates |
| `candidate_ping` | ~2с | Удаляет протухших кандидатов (>30с), обновляет `best_candidates`, probe'ит неактивных |
## 7. Ключевые константы
| Константа | Значение | Смысл |
|-----------|----------|-------|
| `CONN_MGR_MAX_CANDIDATES` | 3 | Макс. число лучших посредников в списке |
| `CONN_MGR_MAX_INTERMEDIARIES` | 3 | Макс. число выбранных посредников для одного соединения |
| `CONN_MGR_CONNECT_DIRECT_TIMEOUT_MS` | 5000 | Таймаут фазы DIRECT |
| `CONN_MGR_CONNECT_DIRECT_TIMEOUT_MS` | 2000 | Таймаут фазы DIRECT |
| `CONN_MGR_CONNECT_REVERSE_TIMEOUT_MS` | 15000 | Таймаут фазы REVERSE |
| `CONN_MGR_INTERM_EXCHANGE_TIMEOUT_MS` | 15000 | Таймаут фазы INDIRECT |
| `CONN_MGR_CANDIDATE_STALE_TB` | 300000 | Порог устаревания кандидата (~30 сек) |
| `CONN_MGR_CANDIDATE_PING_TB` | 20000 | Интервал ping кандидатов (~2 сек) |
| `CONN_MGR_BG_PING_INTERVAL_TB` | 1000 | Интервал фонового ping нод (~100 мс) |
| `CONN_MGR_BG_PING_CYCLE_MIN_TB` | 100000 | Мин. длительность цикла обхода всех нод (~10 сек) |
| `CONN_MGR_IDLE_CHECK_INTERVAL_TB` | 10000 | Интервал проверки idle (~1 сек) |
| `CONN_MGR_CANDIDATE_STALE_TB` | 300000 | Порог устаревания кандидата (~30с) |
| `CONN_MGR_CANDIDATE_PING_TB` | 20000 | Интервал ping кандидатов (~2с) |
| `CONN_MGR_BG_PING_INTERVAL_TB` | 1000 | Интервал фонового ping нод (~100мс) |
| `CONN_MGR_BG_PING_CYCLE_MIN_TB` | 100000 | Мин. длительность цикла обхода всех нод (~10с) |
| `CONN_MGR_LOCAL_SCAN_ATTEMPTS` | 3 | Попыток локального сканирования |
| `CONN_MGR_LOCAL_SCAN_TIMEOUT_MS` | 100 | Таймаут одной попытки локального скана |
| `CACHE_MS` | 20000 | Срок свежести RTT-замера (~20с) |
## 8. Типы соединений (итоговый результат)
| Тип | Значение | Описание |
|-----|----------|----------|
| `CONN_TYPE_NONE` | 0 | Нет соединения |
| `CONN_TYPE_DIRECT` | 1 | Прямое соединение |
| `CONN_TYPE_REVERSE` | 2 | Цель подключилась к нам |
| `CONN_TYPE_INDIRECT` | 3 | Через посредника |
## 9. Тесты
| Тест | Покрытие |
|------|----------|
| `tests/test_conn_mgr.c` | DIRECT (существующее соединение + узел из БД), TIMEOUT |
| `tests/test_conn_mgr_already_connected.c` | DIRECT к уже подключённому узлу |
| `tests/test_conn_mgr_phases.c` | Все три фазы на loopback-топологии C(relay)↔A/B: DIRECT, REVERSE (A public + B nat), INDIRECT (A/B nat через C). Адреса цели «портятся» на неверный порт, чтобы DIRECT честно провалился; таймаут probe уменьшен до 10мс через `route_connectivity_set_probe_timeout_ms()` |

47
src/routing_layer/conn_mgr_indirect.c

@ -49,14 +49,21 @@ void cm_exchange_timeout_cb(void* arg) {
* Отправляет INTERM_SELECTED инициатору и доставляет CONN_EVENT_UP. */
void cm_compute_intermediaries(struct CONN_MGR_ENTRY* entry, struct CM_EXCHANGE_RESP* resp) {
struct { uint64_t nid; uint64_t rtt; } all[8]; uint8_t cnt=0;
for (uint8_t i=0;i<resp->your_count&&cnt<8;i++) {
uint16_t our=0; for(uint8_t j=0;j<CONN_MGR_MAX_CANDIDATES;j++) if(entry->mgr->best_candidates[j].node_id==resp->your_candidates[i].node_id){our=entry->mgr->best_candidates[j].rtt;break;}
all[cnt]=(typeof(all[0])){resp->your_candidates[i].node_id,(uint64_t)our+resp->your_candidates[i].rtt}; cnt++;
uint8_t your_n = resp->your_count > 4 ? 4 : resp->your_count;
uint8_t my_n = resp->my_count > 4 ? 4 : resp->my_count;
uint64_t self = entry->mgr->instance->node_id, target = entry->node_id;
for (uint8_t i=0;i<your_n&&cnt<8;i++) {
uint64_t nid=resp->your_candidates[i].node_id;
if (nid==self||nid==target) continue; /* сам/цель не может быть посредником */
uint16_t our=0; for(uint8_t j=0;j<CONN_MGR_MAX_CANDIDATES;j++) if(entry->mgr->best_candidates[j].node_id==nid){our=entry->mgr->best_candidates[j].rtt;break;}
all[cnt]=(typeof(all[0])){nid,(uint64_t)our+resp->your_candidates[i].rtt}; cnt++;
}
for (uint8_t i=0;i<resp->my_count&&cnt<8;i++) {
struct TOPO_GROUP_NODE* nq=topo_node_find_by_id(entry->mgr->group,resp->my_candidates[i].node_id);
for (uint8_t i=0;i<my_n&&cnt<8;i++) {
uint64_t nid=resp->my_candidates[i].node_id;
if (nid==self||nid==target) continue; /* сам/цель не может быть посредником */
struct TOPO_GROUP_NODE* nq=topo_node_find_by_id(entry->mgr->group,nid);
uint16_t our=nq?topo_get_chain_rtt(nq):0; if(our==0xFFFF)our=0;
all[cnt]=(typeof(all[0])){resp->my_candidates[i].node_id,(uint64_t)our+resp->my_candidates[i].rtt}; cnt++;
all[cnt]=(typeof(all[0])){nid,(uint64_t)our+resp->my_candidates[i].rtt}; cnt++;
}
for (uint8_t i=0;i<cnt;i++) for(uint8_t j=i+1;j<cnt;j++) if(all[j].rtt<all[i].rtt){typeof(all[0])t=all[i];all[i]=all[j];all[j]=t;}
uint8_t sel=cnt>CONN_MGR_MAX_INTERMEDIARIES?CONN_MGR_MAX_INTERMEDIARIES:cnt;
@ -78,7 +85,8 @@ void cm_exchange_probe_retry_cb(void* arg) {
uint8_t need=0; uint64_t now=get_time_tb();
for(uint8_t i=0;i<ep->cached.my_count&&i<4;i++){
struct TOPO_GROUP_NODE* nq=topo_node_find_by_id(ep->entry->mgr->group,ep->cached.my_candidates[i].node_id);
if(!nq||!cm_is_rtt_fresh(nq,now)){need++;break;}
if(!nq) continue; /* self / не в группе — не посредник, пропускаем */
if(!cm_is_rtt_fresh(nq,now)){need++;break;}
}
if(!need){cm_compute_intermediaries(ep->entry,&ep->cached);struct cm_exchange_pending**pp=&ep->entry->mgr->exchange_pending;while(*pp){if(*pp==ep){*pp=ep->next;break;}pp=&(*pp)->next;}u_free(ep);}
else ep->timer=uasync_set_timeout(ep->entry->mgr->instance->ua,2000,ep,cm_exchange_probe_retry_cb,"cm_exch_probe");
@ -86,19 +94,24 @@ void cm_exchange_probe_retry_cb(void* arg) {
/* Инициатор получил EXCHANGE_RESP: кеширует ответ, проверяет свежесть RTT
* до СВОИХ кандидатов (тех, что в my_candidates ответа). Если протухли —
* probe + retry-таймер 2с. Если свежие — cm_compute_intermediaries. */
* probe + retry-таймер 200мс. Если свежие — cm_compute_intermediaries. */
void cm_handle_interm_exchange_resp(struct CONN_MGR* mgr, const uint8_t* data, size_t len) {
const struct CM_EXCHANGE_RESP* resp=(const struct CM_EXCHANGE_RESP*)data;
if(len<offsetof(struct CM_EXCHANGE_RESP,your_candidates))return;
if(resp->my_count>4||resp->your_count>4)return;
if(len<offsetof(struct CM_EXCHANGE_RESP,your_candidates)+(size_t)resp->your_count*sizeof(struct CONN_MGR_CANDIDATE))return;
struct cm_exchange_pending* ep=mgr->exchange_pending;
while(ep){if(ep->request_id==resp->request_id)break;ep=ep->next;}
if(!ep)return;
if(!ep){DEBUG_INFO(DEBUG_CATEGORY_GENERAL,"conn_mgr: EXCHANGE_RESP req=%u no pending",resp->request_id);return;}
DEBUG_INFO(DEBUG_CATEGORY_GENERAL,"conn_mgr: EXCHANGE_RESP req=%u my=%u your=%u",
resp->request_id, resp->my_count, resp->your_count);
if(ep->entry->main.timer){uasync_cancel_timeout(mgr->instance->ua,ep->entry->main.timer);ep->entry->main.timer=NULL;}
memcpy(&ep->cached,resp,len<sizeof(ep->cached)?len:sizeof(ep->cached)); ep->received=1; ep->probed=0;
uint8_t need=0; uint64_t now=get_time_tb();
for(uint8_t i=0;i<ep->cached.my_count&&i<4;i++){
struct TOPO_GROUP_NODE* nq=topo_node_find_by_id(mgr->group,ep->cached.my_candidates[i].node_id);
if(!nq||!cm_is_rtt_fresh(nq,now)){need++;route_connectivity_probe_node(mgr->instance,mgr->group,nq);}
if(!nq) continue; /* self / не в группе — не посредник, пропускаем */
if(!cm_is_rtt_fresh(nq,now)){need++;route_connectivity_probe_node(mgr->instance,mgr->group,nq);}
}
if(!need){ep->probed=1;cm_compute_intermediaries(ep->entry,&ep->cached);struct cm_exchange_pending**pp=&mgr->exchange_pending;while(*pp){if(*pp==ep){*pp=ep->next;break;}pp=&(*pp)->next;}u_free(ep);return;}
ep->timer=uasync_set_timeout(mgr->instance->ua,2000,ep,cm_exchange_probe_retry_cb,"cm_exch_probe");
@ -106,16 +119,19 @@ void cm_handle_interm_exchange_resp(struct CONN_MGR* mgr, const uint8_t* data, s
/* Финальный выбор посредников от инициатора (приходит второй стороне INDIRECT-
* обмена). Сохраняет выбранные intermediaries в entry, помечает CONNECTED. */
void cm_handle_interm_selected(struct CONN_MGR* mgr, const uint8_t* data, size_t len) {
void cm_handle_interm_selected(struct CONN_MGR* mgr, uint64_t src, const uint8_t* data, size_t len) {
const struct CM_INTERM_SEL* sel=(const struct CM_INTERM_SEL*)data;
if(len<offsetof(struct CM_INTERM_SEL,selected))return;
struct CONN_MGR_ENTRY* entry = NULL; struct ll_entry* ee = mgr->entries->head;
while (ee) { struct CONN_MGR_ENTRY* e = (struct CONN_MGR_ENTRY*)ee; if (e->state == CONN_MGR_STATE_CONNECTING) { entry = e; break; } ee = ee->next; }
if(!entry)return;
uint8_t c=sel->count>CONN_MGR_MAX_INTERMEDIARIES?CONN_MGR_MAX_INTERMEDIARIES:sel->count;
if(len<offsetof(struct CM_INTERM_SEL,selected)+(size_t)c*sizeof(struct CONN_MGR_CANDIDATE))return;
struct CONN_MGR_ENTRY* entry = cm_find_entry(mgr, src);
if(!entry)return;
if(entry->state==CONN_MGR_STATE_CONNECTED)return; /* симметричный случай — UP уже доставлен инициаторским путём */
if(entry->main.phase!=3)return; /* не INDIRECT-фаза */
for(uint8_t i=0;i<c;i++)entry->intermediaries[i]=sel->selected[i].node_id;
entry->intermediariy_count=c; entry->rr_idx=0; entry->conn_type=CONN_TYPE_INDIRECT; entry->state=CONN_MGR_STATE_CONNECTED;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL,"conn_mgr: got INTERM_SELECTED count=%u",c);
cm_update_nodeinfo(entry); cm_deliver_event(entry,CONN_EVENT_UP);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL,"conn_mgr: got INTERM_SELECTED from 0x%016llx count=%u",(unsigned long long)src,c);
}
/* Цель получила EXCHANGE_REQ: измеряет RTT до кандидатов инициатора (если
@ -123,6 +139,7 @@ void cm_handle_interm_selected(struct CONN_MGR* mgr, const uint8_t* data, size_t
* кандидатами + замерами RTT до кандидатов инициатора. */
void cm_handle_interm_exchange_req(uint64_t src, struct CONN_MGR* mgr, struct CM_EXCHANGE_REQ* req) {
if (!mgr) return;
if (req->candidate_count > 4) return;
uint64_t now=get_time_tb();
for(uint8_t i=0;i<req->candidate_count&&i<4;i++){
struct TOPO_GROUP_NODE* nq=topo_node_find_by_id(mgr->group,req->candidates[i].node_id);

16
src/routing_layer/conn_mgr_monitor.c

@ -55,22 +55,6 @@ void conn_mgr_update_best_candidates(struct CONN_MGR* mgr, uint64_t node_id, uin
{ struct CONN_MGR_CANDIDATE t = mgr->best_candidates[i]; mgr->best_candidates[i]=mgr->best_candidates[i-1]; mgr->best_candidates[i-1]=t; }
}
/* ═══════ idle ═══════ */
/* Проверка неактивности: если трафика не было > idle_timeout_ms — доставляет
* CONN_EVENT_TIMEOUT всем handle'ам и чистит entry. Таймер самоперезапускается. */
void cm_idle_timer_cb(void* arg) {
struct CONN_MGR_ENTRY* entry = (struct CONN_MGR_ENTRY*)arg;
if (entry->state != CONN_MGR_STATE_CONNECTED || !entry->idle_timeout_ms) return;
uint64_t idle = get_time_tb() - entry->last_traffic_tb;
if (idle > (uint64_t)entry->idle_timeout_ms * 10) {
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: node 0x%016llx idle %llums, disconnecting",
(unsigned long long)entry->node_id, (unsigned long long)(idle / 10));
cm_deliver_event(entry, CONN_EVENT_TIMEOUT); cm_entry_cleanup(entry); return;
}
entry->idle_timer = uasync_set_timeout(entry->mgr->instance->ua, CONN_MGR_IDLE_CHECK_INTERVAL_TB, entry, cm_idle_timer_cb, "conn_mgr_idle");
}
/* ═══════ bg_ping ═══════ */
static void cm_bg_ping_noop_cb(int s, uint16_t r, void* a, uint64_t n, const uint8_t* d, size_t l) { (void)s;(void)r;(void)a;(void)n;(void)d;(void)l; }

15
src/routing_layer/conn_mgr_priv.h

@ -25,13 +25,12 @@ struct stcp_link;
#define CONN_MGR_CANDIDATE_PING_TB 20000
#define CONN_MGR_BG_PING_INTERVAL_TB 1000
#define CONN_MGR_BG_PING_CYCLE_MIN_TB 100000
#define CONN_MGR_IDLE_CHECK_INTERVAL_TB 10000
#define CONN_MGR_LOCAL_SCAN_ATTEMPTS 3
#define CONN_MGR_LOCAL_SCAN_TIMEOUT_MS 100
#define CONN_MGR_CONNECT_DIRECT_TIMEOUT_MS 2000
#define CONN_MGR_CONNECT_REVERSE_TIMEOUT_MS 15000
#define CONN_MGR_INTERM_EXCHANGE_TIMEOUT_MS 15000
#define CACHE_MS 20000 /* время жизни RTT-замера (через cm_is_rtt_fresh) */
#define CACHE_MS 20000 /* время жизни RTT-замера ~20с (через cm_is_rtt_fresh, мс) */
#define CONN_MGR_MAX_INTERMEDIARIES 3
/* ═══════ состояния ═══════ */
@ -58,10 +57,6 @@ enum {
CM_SUBCMD_INTERM_EXCHANGE_RESP= 0x04,
CM_SUBCMD_INTERM_SELECTED = 0x05,
CM_SUBCMD_DISCONNECT = 0x06,
CM_SUBCMD_BRIDGE_REQUEST = 0x10,
CM_SUBCMD_BRIDGE_RESPONSE = 0x11,
CM_SUBCMD_BRIDGE_CLOSE = 0x12,
CM_SUBCMD_BRIDGE_CLOSED = 0x13,
};
struct CONN_MGR_CANDIDATE { uint64_t node_id; uint16_t rtt; } __attribute__((packed));
@ -102,7 +97,7 @@ struct CONN_MGR_ENTRY {
struct ll_entry ll; /* ll_queue entry — первый элемент */
uint64_t node_id; /* хеш-ключ 8B, offset от data = 0 */
uint8_t state, conn_type;
uint32_t idle_timeout_ms; uint64_t last_traffic_tb; void* idle_timer;
uint64_t last_traffic_tb;
struct CONN_MGR_HANDLE* handles; struct NODE_CONN_DIRECT* ncd_handle;
uint64_t intermediaries[3]; uint8_t intermediariy_count, rr_idx;
uint8_t local_scan_state, main_connect_state, db_loaded;
@ -126,7 +121,7 @@ struct CONN_MGR {
/* ═══════ декларации внутренних функций ═══════ */
/* core — публичные (в priv.h, не в conn_mgr.h) */
int cm_open(struct CONN_MGR* mgr, uint64_t node_id, uint32_t idle_timeout_ms,
int cm_open(struct CONN_MGR* mgr, uint64_t node_id,
conn_mgr_cb_t cb, void* cb_arg, struct CONN_MGR_HANDLE** out_handle);
/* core */
@ -144,7 +139,6 @@ void cm_start_phase_indirect(struct CONN_MGR_ENTRY* entry);
void cm_start_local_scan(struct CONN_MGR_ENTRY* entry);
void cm_handle_direct_req(uint64_t src, struct CONN_MGR* mgr, const uint8_t* data, size_t len);
void cm_handle_disconnect(struct CONN_MGR* mgr, uint64_t node_id);
void cm_reverse_init_cb(struct ETCP_CONN* conn, int event, void* arg);
void cm_reverse_timeout_cb(void* arg);
void cm_add_v4_link(struct ETCP_CONN* conn, const uint8_t* addr, uint16_t port, struct ETCP_SOCKET* s);
void cm_add_v6_link(struct ETCP_CONN* conn, const uint8_t addr[16], uint16_t port, struct ETCP_SOCKET* s);
@ -165,10 +159,9 @@ void cm_compute_intermediaries(struct CONN_MGR_ENTRY* entry, struct CM_EXCHANGE_
void cm_exchange_probe_retry_cb(void* arg);
void cm_handle_interm_exchange_req(uint64_t src, struct CONN_MGR* mgr, struct CM_EXCHANGE_REQ* req);
void cm_handle_interm_exchange_resp(struct CONN_MGR* mgr, const uint8_t* data, size_t len);
void cm_handle_interm_selected(struct CONN_MGR* mgr, const uint8_t* data, size_t len);
void cm_handle_interm_selected(struct CONN_MGR* mgr, uint64_t src, const uint8_t* data, size_t len);
/* monitor */
void cm_idle_timer_cb(void* arg);
void cm_bg_ping_timer_cb(void* arg);
void cm_candidate_ping_timer_cb(void* arg);
uint64_t cm_get_node_max_probe_time(struct TOPO_GROUP_NODE* nq);

10
src/routing_layer/route_connectivity.c

@ -24,6 +24,12 @@
#define CONN_MAX_SOCKET_CANDIDATES 8
static uint16_t g_probe_timeout_ms = CONN_PROBE_TIMEOUT_MS;
void route_connectivity_set_probe_timeout_ms(uint16_t ms) {
if (ms) g_probe_timeout_ms = ms;
}
struct conn_probe_ctx {
struct conn_probe_ctx* next; /* linked list in nq->connectivity.probe_list */
struct UTUN_INSTANCE* instance;
@ -317,7 +323,7 @@ void route_connectivity_probe_node(struct UTUN_INSTANCE* instance, struct TOPO_G
struct conn_probe_ctx* ctx = u_calloc(1, sizeof(struct conn_probe_ctx)); if (!ctx) continue;
ctx->instance = instance; ctx->group = group; ctx->nq = nq; ctx->addr_type = a->type; ctx->target_addr = target;
memcpy(ctx->peer_pubkey, ni->public_key, SC_PUBKEY_SIZE);
ctx->timeout_ms = CONN_PROBE_TIMEOUT_MS; ctx->best_across_sockets = 65535;
ctx->timeout_ms = g_probe_timeout_ms; ctx->best_across_sockets = 65535;
if (a->protocol & TOPO_PROTO_TCP) {
ctx->is_tcp = 1;
@ -355,7 +361,7 @@ void route_connectivity_probe_node(struct UTUN_INSTANCE* instance, struct TOPO_G
struct conn_probe_ctx* ctx = u_calloc(1, sizeof(struct conn_probe_ctx)); if (!ctx) continue;
ctx->instance = instance; ctx->nq = nq; ctx->addr_type = a6->type; ctx->target_addr = target;
memcpy(ctx->peer_pubkey, ni->public_key, SC_PUBKEY_SIZE);
ctx->timeout_ms = CONN_PROBE_TIMEOUT_MS; ctx->best_across_sockets = 65535;
ctx->timeout_ms = g_probe_timeout_ms; ctx->best_across_sockets = 65535;
if (a6->protocol & TOPO_PROTO_TCP) {
ctx->is_tcp = 1;

4
src/routing_layer/route_connectivity.h

@ -23,6 +23,10 @@ void route_connectivity_cancel_node(struct UTUN_INSTANCE* instance,
// Отменяет все pending пробы для всех узлов (при destroy)
void route_connectivity_cancel_all(struct UTUN_INSTANCE* instance);
// Устанавливает таймаут одного probe-пинга (по умолчанию CONN_PROBE_TIMEOUT_MS).
// Используется тестами для ускорения probe-циклов.
void route_connectivity_set_probe_timeout_ms(uint16_t ms);
#ifdef __cplusplus
}

10
tests/Makefile.am

@ -52,6 +52,8 @@ check_PROGRAMS = \
test_bgp_triangle \
test_broadcast \
test_conn_mgr \
test_conn_mgr_phases \
test_conn_mgr_handles \
test_conn_mgr_already_connected \
test_sock_match \
test_invite_group_create \
@ -345,6 +347,14 @@ test_conn_mgr_SOURCES = test_conn_mgr.c
test_conn_mgr_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib
test_conn_mgr_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS)
test_conn_mgr_phases_SOURCES = test_conn_mgr_phases.c
test_conn_mgr_phases_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib
test_conn_mgr_phases_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS)
test_conn_mgr_handles_SOURCES = test_conn_mgr_handles.c
test_conn_mgr_handles_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib
test_conn_mgr_handles_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS)
test_conn_mgr_already_connected_SOURCES = test_conn_mgr_already_connected.c
test_conn_mgr_already_connected_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib
test_conn_mgr_already_connected_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS)

82
tests/test_conn_mgr_handles.c

@ -1,10 +1,10 @@
/**
* @file test_conn_mgr_handles.c
* @brief conn_mgr handle lifecycle: multiple handles sharing one connection
* @brief conn_mgr: жизненный цикл handle (refcount нескольких handle на один entry)
*
* Test 1: Two conn_mgr_connect_node → both OK, same entry, one conn
* Test 2: Release first handle → conn stays (second handle still active)
* Test 3: Release last handle → conn torn down
* Test 1: два conn_mgr_open на один узел → оба получают CONN_EVENT_UP, entry CONNECTED/DIRECT
* Test 2: закрыть первый handle → entry остаётся живым (второй handle держит)
* Test 3: закрыть последний handle → entry очищен (conn_mgr_type == CONN_TYPE_NONE)
*/
#include <stdio.h>
#include <stdlib.h>
@ -51,15 +51,23 @@ static void fail(const char* msg) { fprintf(stderr, "FAIL: %s\n", msg); fflush(s
static void done(void) { if (result == 0) result = 1; }
static void to_cb(void* arg) { (void)arg; fail("timeout"); done(); }
static void ccb1(int r, uint64_t id, void* arg) {
(void)arg; cb1_ok = (r == CONN_MGR_OK);
fprintf(stderr, "ccb1: result=%d node=0x%llx ok=%d\n", r, (unsigned long long)id, cb1_ok); fflush(stderr);
static void ccb1(struct CONN_MGR_HANDLE* h, uint64_t id, uint64_t gid, enum conn_mgr_event ev, void* arg) {
(void)h; (void)arg; (void)gid;
cb1_ok = (ev == CONN_EVENT_UP);
fprintf(stderr, "ccb1: event=%d node=0x%llx ok=%d\n", ev, (unsigned long long)id, cb1_ok); fflush(stderr);
}
static void ccb2(int r, uint64_t id, void* arg) {
(void)arg; cb2_ok = (r == CONN_MGR_OK);
fprintf(stderr, "ccb2: result=%d node=0x%llx ok=%d\n", r, (unsigned long long)id, cb2_ok); fflush(stderr);
static void ccb2(struct CONN_MGR_HANDLE* h, uint64_t id, uint64_t gid, enum conn_mgr_event ev, void* arg) {
(void)h; (void)arg; (void)gid;
cb2_ok = (ev == CONN_EVENT_UP);
fprintf(stderr, "ccb2: event=%d node=0x%llx ok=%d\n", ev, (unsigned long long)id, cb2_ok); fflush(stderr);
}
static uint8_t entry_type(void) {
struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(topo_groups_get_default(g_a->topo_groups), g_nid_b);
return nq ? nq->conn_mgr_type : CONN_TYPE_NONE;
}
static void test1(void* arg);
static void test1_check(void* arg);
static void test2_check(void* arg);
static void test2_verify(void* arg);
@ -68,54 +76,49 @@ static void test3_verify(void* arg);
static void test1(void* arg) {
(void)arg; if (result) return;
if (!topo_node_find_by_id(topo_groups_get_default(g_a->topo_groups), g_nid_b)) { uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test1, "t1w"); return; }
struct ETCP_CONN* c = topo_group_find_conn_for_node(g_a->conn_mgr->group, g_nid_b);
if (!c) c = instance_find_conn(g_a, g_nid_b);
if (!c || !c->links) { uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test1, "t1w2"); return; }
struct ETCP_CONN* c = instance_find_conn(g_a, g_nid_b);
if (!c || !c->links) { uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test1, "t1w"); return; }
int has = 0; struct ETCP_LINK* l = c->links;
while (l) { if (l->link_state == 3 && l->initialized) { has = 1; break; } l = l->next; }
if (!has) { uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test1, "t1w3"); return; }
fprintf(stderr, "Test 1: two conn_mgr_connect_node to same node\n"); fflush(stderr);
if (!has) { uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test1, "t1w2"); return; }
fprintf(stderr, "Test 1: two conn_mgr_open to same node\n"); fflush(stderr);
cb1_ok = 0; cb2_ok = 0;
gh1 = conn_mgr_connect_node(g_a->conn_mgr, g_nid_b, 0, ccb1, NULL);
gh2 = conn_mgr_connect_node(g_a->conn_mgr, g_nid_b, 0, ccb2, NULL);
if (!gh1 || !gh2) { fail("test1: handle returned NULL"); done(); return; }
if (conn_mgr_open(g_a, TOPO_GROUP_UTUN, g_nid_b, ccb1, NULL, &gh1) != 0 || !gh1) { fail("test1: open1 failed"); done(); return; }
if (conn_mgr_open(g_a, TOPO_GROUP_UTUN, g_nid_b, ccb2, NULL, &gh2) != 0 || !gh2) { fail("test1: open2 failed"); done(); return; }
uasync_set_timeout(ua, 200, NULL, (timeout_callback_t)test1_check, "t1c");
}
static void test1_check(void* arg) {
(void)arg;
if (!cb1_ok || !cb2_ok) { fail("test1: both callbacks should be OK"); done(); return; }
uint8_t st, ty; conn_mgr_get_status(g_a->conn_mgr, g_nid_b, &st, &ty);
if (st != CONN_MGR_STATE_CONNECTED || ty != CONN_TYPE_DIRECT) { fail("test1: status mismatch"); done(); return; }
fprintf(stderr, " OK: both handles returned OK, state=%d type=%d\n", st, ty); fflush(stderr);
uasync_call_soon(ua, NULL, (timeout_callback_t)test2_check);
if (!cb1_ok || !cb2_ok) { fail("test1: both callbacks should be UP"); done(); return; }
uint8_t ty = entry_type();
if (ty != CONN_TYPE_DIRECT) { fail("test1: conn_type != DIRECT"); done(); return; }
fprintf(stderr, " OK: both handles UP, conn_type=DIRECT\n"); fflush(stderr);
uasync_call_soon(ua, NULL, test2_check);
}
static void test2_check(void* arg) {
(void)arg;
fprintf(stderr, "Test 2: release first handle — conn stays alive\n"); fflush(stderr);
conn_mgr_release_handle(gh1); gh1 = NULL;
fprintf(stderr, "Test 2: close first handle — entry stays alive\n"); fflush(stderr);
conn_mgr_close(gh1); gh1 = NULL;
uasync_set_timeout(ua, 100, NULL, (timeout_callback_t)test2_verify, "t2v");
}
static void test2_verify(void* arg) {
(void)arg;
uint8_t st, ty; conn_mgr_get_status(g_a->conn_mgr, g_nid_b, &st, &ty);
if (st != CONN_MGR_STATE_CONNECTED) { fail("test2: conn died after first handle close"); done(); return; }
fprintf(stderr, " OK: conn alive after h1 close, state=%d\n", st); fflush(stderr);
uasync_call_soon(ua, NULL, (timeout_callback_t)test3_check);
if (entry_type() != CONN_TYPE_DIRECT) { fail("test2: entry died after first handle close"); done(); return; }
fprintf(stderr, " OK: entry alive after h1 close\n"); fflush(stderr);
uasync_call_soon(ua, NULL, test3_check);
}
static void test3_check(void* arg) {
(void)arg;
fprintf(stderr, "Test 3: release last handle — conn torn down\n"); fflush(stderr);
conn_mgr_release_handle(gh2); gh2 = NULL;
fprintf(stderr, "Test 3: close last handle — entry cleaned up\n"); fflush(stderr);
conn_mgr_close(gh2); gh2 = NULL;
uasync_set_timeout(ua, 200, NULL, (timeout_callback_t)test3_verify, "t3v");
}
static void test3_verify(void* arg) {
(void)arg;
uint8_t st, ty; conn_mgr_get_status(g_a->conn_mgr, g_nid_b, &st, &ty);
if (st != CONN_MGR_STATE_DISCONNECTED) { fail("test3: conn should be disconnected"); done(); return; }
fprintf(stderr, " OK: conn DISCONNECTED after last handle close\n"); fflush(stderr);
if (entry_type() != CONN_TYPE_NONE) { fail("test3: entry should be cleaned up"); done(); return; }
fprintf(stderr, " OK: entry cleaned up (conn_type=NONE)\n"); fflush(stderr);
fprintf(stderr, "=== ALL DONE ===\n"); fflush(stderr);
done();
}
@ -146,9 +149,12 @@ int main(void) {
{ int el = 0; while (!result && el < TIMEOUT_TB/10 + 500) { uasync_poll(ua, POLL_MS); el += POLL_MS; } }
done:
if (ttimer) { uasync_cancel_timeout(ua, ttimer); }
g_a->running = 0; if (g_b) g_b->running = 0;
if (ua) uasync_destroy(ua, 1); // force cleanup
if (ttimer) { uasync_cancel_timeout(ua, ttimer); ttimer = NULL; }
if (gh1) { conn_mgr_close(gh1); gh1 = NULL; }
if (gh2) { conn_mgr_close(gh2); gh2 = NULL; }
if (g_a) { g_a->running = 0; utun_instance_destroy(g_a); g_a = NULL; }
if (g_b) { g_b->running = 0; utun_instance_destroy(g_b); g_b = NULL; }
if (ua) { uasync_destroy(ua, 0); ua = NULL; }
test_unlink(ca); test_unlink(cb); test_rmdir(tdir);
return (result == 1) ? 0 : 1;
}

322
tests/test_conn_mgr_phases.c

@ -0,0 +1,322 @@
/**
* @file test_conn_mgr_phases.c
* @brief conn_mgr: все три варианта подключения (DIRECT / REVERSE / INDIRECT)
*
* Топология каждого сценария — три узла на loopback: C (публичный релей),
* A и B подключаются к C как клиенты, через BGP узнают друг друга.
*
* Сценарий 1 DIRECT: A=public, B=public → A подключается к B напрямую.
* Сценарий 2 REVERSE: A=public, B=nat → A не пробивает B (адрес испорчен),
* шлёт DIRECT_REQ, B подключается к A (обратное подключение).
* Сценарий 3 INDIRECT: A=nat, B=nat → оба без прямого адреса, выбирают
* посредника C через обмен кандидатами.
*
* Проверка честная: фазы выбираются реальным кодом (адреса цели испорчены на
* «неверный» порт, чтобы DIRECT провалился), RTT кандидата зреет настоящими
* probe-циклами (таймаут probe уменьшен до 10мс для скорости).
*/
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <stdarg.h>
#include "test_utils.h"
#include "etcp.h"
#include "etcp_connections.h"
#include "../src/config_parser.h"
#include "../src/config_updater.h"
#include "../src/utun_instance.h"
#include "topo_group.h"
#include "topo_node.h"
#include "conn_mgr.h"
#include "conn_mgr_priv.h"
#include "route_connectivity.h"
#include "../lib/u_async.h"
#include "../lib/debug_config.h"
#include "../lib/mem.h"
#define POLL_MS 5
#define LINK_TIMEOUT_MS 8000
#define BGP_TIMEOUT_MS 8000
#define CONN_TIMEOUT_MS 15000
static struct UASYNC* ua = NULL;
static struct UTUN_INSTANCE *inst_a = NULL, *inst_b = NULL, *inst_c = NULL;
static uint64_t nid_a = 0, nid_b = 0, nid_c = 0;
static volatile int ev_a = -1, ev_b = -1;
static struct CONN_MGR_HANDLE *h_a = NULL, *h_b = NULL;
static int g_failed = 0;
static char temp_dir[] = "/tmp/utun_cm_phases_XXXXXX";
#define TEMP_DIR_TEMPLATE "/tmp/utun_cm_phases_XXXXXX"
static char cfg_a[256], cfg_b[256], cfg_c[256];
static int base_port = 0;
/* ─── утилиты ─── */
static int write_file(const char* path, const char* fmt, ...) {
va_list ap; FILE* f = fopen(path, "w");
if (!f) return -1;
va_start(ap, fmt); vfprintf(f, fmt, ap); va_end(ap); fclose(f);
return 0;
}
static char* get_pubkey(const char* path) {
struct utun_config* cfg = parse_config(path);
if (!cfg) return NULL;
char* pub = u_strdup(cfg->global.my_public_key_hex);
free_config(cfg);
return pub;
}
static void fail(const char* msg) { fprintf(stderr, "FAIL: %s\n", msg); fflush(stderr); g_failed = 1; }
static void conn_cb_a(struct CONN_MGR_HANDLE* h, uint64_t node, uint64_t grp, enum conn_mgr_event ev, void* arg) {
(void)h; (void)node; (void)grp; (void)arg; ev_a = (int)ev;
}
static void conn_cb_b(struct CONN_MGR_HANDLE* h, uint64_t node, uint64_t grp, enum conn_mgr_event ev, void* arg) {
(void)h; (void)node; (void)grp; (void)arg; ev_b = (int)ev;
}
static int poll_until(int (*cond)(void), uint32_t timeout_ms) {
uint64_t start = get_time_tb();
while (!cond() && (get_time_tb() - start) < (uint64_t)timeout_ms * 10) uasync_poll(ua, POLL_MS);
return cond();
}
/* forward-декларации условий сценариев */
static int ev_a_up(void);
static int ev_ab_up(void);
static int candidates_a_and_b(void);
static int has_initialized_link(struct UTUN_INSTANCE* inst) {
if (!inst || !inst->connections) return 0;
struct ll_entry* e = inst->connections->head;
while (e) {
struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data;
struct ETCP_LINK* l = ce->conn->links;
while (l) { if (l->initialized) return 1; l = l->next; }
e = e->next;
}
return 0;
}
static struct TOPO_GROUP* defgrp(struct UTUN_INSTANCE* inst) { return topo_groups_get_default(inst->topo_groups); }
static int learned(struct UTUN_INSTANCE* inst, uint64_t node_id) {
return defgrp(inst) && topo_node_find_by_id(defgrp(inst), node_id) != NULL;
}
/* Меняем порт адресов цели на «неверный» (закрытый) — DIRECT создаст линк,
* но handshake не получит ответа и фаза провалится по таймауту NCD. */
static void mangle_addrs(struct UTUN_INSTANCE* inst, uint64_t node_id, uint16_t wrong_port) {
struct TOPO_NODE* ni = topo_node_registry_find(inst->topo_groups, node_id);
if (!ni) return;
for (struct TOPO_ADDR4* a = ni->v4_addrs; a; a = a->next) a->port = wrong_port;
for (struct TOPO_ADDR6* a6 = ni->v6_addrs; a6; a6 = a6->next) a6->port = wrong_port;
}
static int best_candidates_ready(struct UTUN_INSTANCE* inst) {
struct TOPO_GROUP* g = defgrp(inst);
return g && g->conn_mgr && g->conn_mgr->best_candidate_count > 0;
}
static uint8_t conn_type_of(struct UTUN_INSTANCE* inst, uint64_t node_id) {
struct TOPO_GROUP_NODE* nq = defgrp(inst) ? topo_node_find_by_id(defgrp(inst), node_id) : NULL;
return nq ? nq->conn_mgr_type : CONN_TYPE_NONE;
}
/* ─── конфиги ─── */
static int write_node_config(const char* path, const char* tun_ip, const char* tun_if,
int srv_port, const char* type, const char* peer_pub, int peer_port) {
if (peer_pub) {
return write_file(path,
"[global]\n"
"tun_ip=%s/24\n"
"tun_ifname=%s\n"
"[server: srv]\n"
"addr=127.0.0.1:%d\n"
"type=%s\n"
"[client: to_relay]\n"
"keepalive=1\n"
"peer_public_key=%s\n"
"link=srv:127.0.0.1:%d\n"
"[allowed_keys]\n"
"allow_all=1\n",
tun_ip, tun_if, srv_port, type, peer_pub, peer_port);
}
return write_file(path,
"[global]\n"
"tun_ip=%s/24\n"
"tun_ifname=%s\n"
"[server: srv]\n"
"addr=127.0.0.1:%d\n"
"type=%s\n"
"[allowed_keys]\n"
"allow_all=1\n",
tun_ip, tun_if, srv_port, type);
}
/* Создаёт конфиги C (public relay), A и B с заданными типами сокетов. */
static int create_configs(const char* type_a, const char* type_b) {
strcpy(temp_dir, TEMP_DIR_TEMPLATE);
if (test_mkdtemp(temp_dir) != 0) { fprintf(stderr, "mkdtemp failed\n"); return -1; }
base_port = 43000 + (getpid() % 12000);
int pc = base_port, pa = base_port + 1, pb = base_port + 2;
snprintf(cfg_c, sizeof(cfg_c), "%s/c.conf", temp_dir);
snprintf(cfg_a, sizeof(cfg_a), "%s/a.conf", temp_dir);
snprintf(cfg_b, sizeof(cfg_b), "%s/b.conf", temp_dir);
if (write_node_config(cfg_c, "10.90.0.3", "tun_c", pc, "public", NULL, 0) != 0) return -1;
if (config_ensure_keys_and_node_id(cfg_c) != 0) return -1;
char* pub_c = get_pubkey(cfg_c);
if (!pub_c) return -1;
if (write_node_config(cfg_a, "10.90.0.1", "tun_a", pa, type_a, pub_c, pc) != 0) { u_free(pub_c); return -1; }
if (write_node_config(cfg_b, "10.90.0.2", "tun_b", pb, type_b, pub_c, pc) != 0) { u_free(pub_c); return -1; }
u_free(pub_c);
if (config_ensure_keys_and_node_id(cfg_a) != 0) return -1;
if (config_ensure_keys_and_node_id(cfg_b) != 0) return -1;
return 0;
}
static void cleanup_configs(void) {
test_unlink(cfg_a); test_unlink(cfg_b); test_unlink(cfg_c);
test_rmdir(temp_dir);
}
/* ─── общий каркас сценария ─── */
static int has_initialized_link_all(void) {
return has_initialized_link(inst_a) && has_initialized_link(inst_b) && has_initialized_link(inst_c);
}
static int learned_ab(void) { return learned(inst_a, nid_b) && learned(inst_b, nid_a); }
/* Инициализирует инстансы и ждёт линков + BGP (A↔B узнают друг друга). */
static int scenario_setup(void) {
inst_c = utun_instance_create(ua, cfg_c);
inst_a = utun_instance_create(ua, cfg_a);
inst_b = utun_instance_create(ua, cfg_b);
if (!inst_c || !inst_a || !inst_b) { fail("instance create"); return -1; }
if (init_connections(inst_c) != 0 || init_connections(inst_a) != 0 || init_connections(inst_b) != 0) {
fail("init_connections"); return -1;
}
nid_a = inst_a->node_id; nid_b = inst_b->node_id; nid_c = inst_c->node_id;
if (!poll_until(has_initialized_link_all, LINK_TIMEOUT_MS)) { fail("links not initialized"); return -1; }
if (!poll_until(learned_ab, BGP_TIMEOUT_MS)) { fail("BGP: A/B did not learn each other"); return -1; }
return 0;
}
static void scenario_teardown(void) {
if (h_a) { conn_mgr_close(h_a); h_a = NULL; }
if (h_b) { conn_mgr_close(h_b); h_b = NULL; }
if (inst_a) { inst_a->running = 0; utun_instance_destroy(inst_a); inst_a = NULL; }
if (inst_b) { inst_b->running = 0; utun_instance_destroy(inst_b); inst_b = NULL; }
if (inst_c) { inst_c->running = 0; utun_instance_destroy(inst_c); inst_c = NULL; }
cleanup_configs();
}
/* ─── Сценарий 1: DIRECT ─── */
static int run_direct(void) {
fprintf(stderr, "\n=== Scenario 1: DIRECT ===\n"); fflush(stderr);
if (create_configs("public", "public") != 0) { fail("configs"); return -1; }
if (scenario_setup() != 0) return -1;
ev_a = -1;
conn_mgr_open(inst_a, TOPO_GROUP_UTUN, nid_b, conn_cb_a, NULL, &h_a);
if (!poll_until(ev_a_up, CONN_TIMEOUT_MS)) { fail("DIRECT: no CONN_EVENT_UP"); return -1; }
if (conn_type_of(inst_a, nid_b) != CONN_TYPE_DIRECT) { fail("DIRECT: conn_type != DIRECT"); return -1; }
fprintf(stderr, " OK: DIRECT\n"); fflush(stderr);
return 0;
}
static int ev_a_up(void) { return ev_a == CONN_EVENT_UP; }
/* ─── Сценарий 2: REVERSE ─── */
static int run_reverse(void) {
fprintf(stderr, "\n=== Scenario 2: REVERSE ===\n"); fflush(stderr);
if (create_configs("public", "nat") != 0) { fail("configs"); return -1; }
if (scenario_setup() != 0) return -1;
/* Испортим адреса B в реестре A — DIRECT провалится, A уйдёт в REVERSE. */
mangle_addrs(inst_a, nid_b, (uint16_t)(base_port + 1000));
inst_a->etcp_connect_timeout_tb = 3000; /* 300мс — быстрый провал DIRECT */
ev_a = -1;
conn_mgr_open(inst_a, TOPO_GROUP_UTUN, nid_b, conn_cb_a, NULL, &h_a);
if (!poll_until(ev_a_up, CONN_TIMEOUT_MS)) { fail("REVERSE: no CONN_EVENT_UP"); return -1; }
if (conn_type_of(inst_a, nid_b) != CONN_TYPE_REVERSE) { fail("REVERSE: conn_type != REVERSE"); return -1; }
fprintf(stderr, " OK: REVERSE\n"); fflush(stderr);
return 0;
}
/* ─── Сценарий 3: INDIRECT ─── */
static int run_indirect(void) {
fprintf(stderr, "\n=== Scenario 3: INDIRECT ===\n"); fflush(stderr);
if (create_configs("nat", "nat") != 0) { fail("configs"); return -1; }
if (scenario_setup() != 0) return -1;
/* Ждём, пока candidate_ping реально заполнит список посредников (C). */
if (!poll_until(candidates_a_and_b, 8000)) { fail("INDIRECT: no candidates (C)"); return -1; }
/* Испортим адреса цели с обеих сторон — DIRECT провалится у обоих. */
mangle_addrs(inst_a, nid_b, (uint16_t)(base_port + 1000));
mangle_addrs(inst_b, nid_a, (uint16_t)(base_port + 1000));
inst_a->etcp_connect_timeout_tb = 3000;
inst_b->etcp_connect_timeout_tb = 3000;
ev_a = -1; ev_b = -1;
conn_mgr_open(inst_a, TOPO_GROUP_UTUN, nid_b, conn_cb_a, NULL, &h_a);
conn_mgr_open(inst_b, TOPO_GROUP_UTUN, nid_a, conn_cb_b, NULL, &h_b);
if (!poll_until(ev_ab_up, CONN_TIMEOUT_MS)) { fail("INDIRECT: no CONN_EVENT_UP on both sides"); return -1; }
if (conn_type_of(inst_a, nid_b) != CONN_TYPE_INDIRECT) { fail("INDIRECT: A conn_type != INDIRECT"); return -1; }
if (conn_type_of(inst_b, nid_a) != CONN_TYPE_INDIRECT) { fail("INDIRECT: B conn_type != INDIRECT"); return -1; }
/* Посредником должен быть C. */
struct TOPO_GROUP_NODE* nq_b = topo_node_find_by_id(defgrp(inst_a), nid_b);
if (!nq_b || nq_b->conn_mgr_intermediariy_count < 1 || nq_b->conn_mgr_intermediaries[0] != nid_c) {
fail("INDIRECT: intermediary != C"); return -1;
}
fprintf(stderr, " OK: INDIRECT via C\n"); fflush(stderr);
return 0;
}
static int candidates_a_and_b(void) { return best_candidates_ready(inst_a) && best_candidates_ready(inst_b); }
static int ev_ab_up(void) { return ev_a == CONN_EVENT_UP && ev_b == CONN_EVENT_UP; }
/* ─── main ─── */
int main(void) {
debug_config_init();
debug_set_level(DEBUG_LEVEL_ERROR);
utun_instance_set_tun_init_enabled(0);
route_connectivity_set_probe_timeout_ms(10);
ua = uasync_create();
if (!ua) { fprintf(stderr, "uasync_create failed\n"); return 1; }
if (run_direct() != 0) g_failed = 1;
scenario_teardown();
ev_a = -1; ev_b = -1; h_a = NULL; h_b = NULL;
if (!g_failed && run_reverse() != 0) g_failed = 1;
scenario_teardown();
ev_a = -1; ev_b = -1; h_a = NULL; h_b = NULL;
if (!g_failed && run_indirect() != 0) g_failed = 1;
scenario_teardown();
if (ua) { uasync_destroy(ua, 0); ua = NULL; }
fprintf(stderr, "%s\n", g_failed ? "RESULT: FAIL" : "RESULT: PASS"); fflush(stderr);
return g_failed ? 1 : 0;
}

34
tools/chatgui-android/build.sh

@ -2,6 +2,8 @@
set -eo pipefail
cd "$(dirname "$0")"
export ANDROID_HOME="${ANDROID_HOME:-/home/user/android}"
APK=app/build/outputs/apk/debug/app-debug.apk
echo "Building..."
TMP=$(mktemp)
./gradlew clean > /dev/null 2>&1
@ -12,9 +14,31 @@ if ! ./gradlew assembleDebug > "$TMP" 2>&1; then
fi
rm -f "$TMP"
echo "Build OK"
if [ "$1" != "noinstall" ]; then
DEV=${1:-$(adb devices 2>/dev/null | awk 'NR==2{print $1}')}
ADB=adb; [ -n "$DEV" ] && ADB="adb -s $DEV"
$ADB install -r app/build/outputs/apk/debug/app-debug.apk
echo "Install OK"
[ "$1" = "noinstall" ] && exit 0
if [ -n "$1" ]; then
DEVICES=("$1")
else
mapfile -t DEVICES < <(adb devices | awk 'NR>1 && $2=="device"{print $1}')
fi
if [ ${#DEVICES[@]} -eq 0 ]; then
echo "No devices connected" >&2
exit 1
fi
echo "Devices (${#DEVICES[@]}):"
i=1
for d in "${DEVICES[@]}"; do
name=$(adb -s "$d" shell getprop ro.product.model 2>/dev/null | tr -d '\r')
[ -z "$name" ] && name=$(adb -s "$d" shell getprop ro.product.marketname 2>/dev/null | tr -d '\r')
echo " $i) $d $name"
i=$((i+1))
done
for d in "${DEVICES[@]}"; do
echo "Installing on $d..."
adb -s "$d" install -r "$APK"
done
echo "Install OK"

Loading…
Cancel
Save