diff --git a/src/routing_layer/conn_mgr.h b/src/routing_layer/conn_mgr.h index 29250d1d..58808f2c 100644 --- a/src/routing_layer/conn_mgr.h +++ b/src/routing_layer/conn_mgr.h @@ -43,7 +43,6 @@ struct TOPO_GROUP; * * bg_ping: раз в ~100ms один узел группы за тик. Поддерживает свежие RTT. * candidate_ping: раз в ~2с, удаляет протухших кандидатов. - * idle: раз в ~1с, рвёт соединение при неактивности. * * === Refcounting и закрытие === * diff --git a/src/routing_layer/conn_mgr_core.c b/src/routing_layer/conn_mgr_core.c index b1f30c86..7eae361c 100644 --- a/src/routing_layer/conn_mgr_core.c +++ b/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;iaddr_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); } diff --git a/src/routing_layer/conn_mgr_doc.md b/src/routing_layer/conn_mgr_doc.md index a69e3d35..10d1f05c 100644 --- a/src/routing_layer/conn_mgr_doc.md +++ b/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()` | diff --git a/src/routing_layer/conn_mgr_indirect.c b/src/routing_layer/conn_mgr_indirect.c index f0bee25a..e4c40711 100644 --- a/src/routing_layer/conn_mgr_indirect.c +++ b/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;iyour_count&&cnt<8;i++) { - uint16_t our=0; for(uint8_t j=0;jmgr->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;iyour_candidates[i].node_id; + if (nid==self||nid==target) continue; /* сам/цель не может быть посредником */ + uint16_t our=0; for(uint8_t j=0;jmgr->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;imy_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;imy_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;iCONN_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;icached.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(lenmy_count>4||resp->your_count>4)return; + if(lenyour_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,lencached)?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;icached.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(lenentries->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(lenstate==CONN_MGR_STATE_CONNECTED)return; /* симметричный случай — UP уже доставлен инициаторским путём */ + if(entry->main.phase!=3)return; /* не INDIRECT-фаза */ for(uint8_t i=0;iintermediaries[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;icandidate_count&&i<4;i++){ struct TOPO_GROUP_NODE* nq=topo_node_find_by_id(mgr->group,req->candidates[i].node_id); diff --git a/src/routing_layer/conn_mgr_monitor.c b/src/routing_layer/conn_mgr_monitor.c index c0577827..d554d14e 100644 --- a/src/routing_layer/conn_mgr_monitor.c +++ b/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; } diff --git a/src/routing_layer/conn_mgr_priv.h b/src/routing_layer/conn_mgr_priv.h index 0150f3d5..bafe263b 100644 --- a/src/routing_layer/conn_mgr_priv.h +++ b/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); diff --git a/src/routing_layer/route_connectivity.c b/src/routing_layer/route_connectivity.c index 72622dbf..fb6da9eb 100644 --- a/src/routing_layer/route_connectivity.c +++ b/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; diff --git a/src/routing_layer/route_connectivity.h b/src/routing_layer/route_connectivity.h index 9e47b2ba..8a002924 100644 --- a/src/routing_layer/route_connectivity.h +++ b/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 } diff --git a/tests/Makefile.am b/tests/Makefile.am index 86afeee7..2ff70aee 100644 --- a/tests/Makefile.am +++ b/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) diff --git a/tests/test_conn_mgr_handles.c b/tests/test_conn_mgr_handles.c index 1c71d9fe..492423e0 100644 --- a/tests/test_conn_mgr_handles.c +++ b/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 #include @@ -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; } diff --git a/tests/test_conn_mgr_phases.c b/tests/test_conn_mgr_phases.c new file mode 100644 index 00000000..7d52a222 --- /dev/null +++ b/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 +#include +#include +#include +#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; +} diff --git a/tools/chatgui-android/build.sh b/tools/chatgui-android/build.sh index 632cf1b5..c80b86d8 100755 --- a/tools/chatgui-android/build.sh +++ b/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"