From 5838e63cde18c89114a0174bdac6e4d57e4871be Mon Sep 17 00:00:00 2001 From: evgeny Date: Mon, 28 Sep 2026 11:07:25 +0300 Subject: [PATCH] Use shared group requests for config and chat connection planners --- src/chat/chat_sync.c | 2 +- src/chat/member_sync.c | 3 +- src/routing_layer/topo_group.c | 57 +- src/routing_layer/topo_group.h | 7 +- src/routing_layer/topo_group_connect.c | 824 +++++-------------- src/routing_layer/topo_group_connect.h | 36 +- src/routing_layer/topo_group_invite.c | 49 +- src/transport_layer/etcp_connections.c | 39 +- src/utun_instance.c | 2 +- src/utun_instance.h | 2 +- tests/bbr_integration/test_bbr_integration.c | 16 +- tests/test_group_recovery.c | 38 + tests/test_ncd_config.c | 7 +- 13 files changed, 345 insertions(+), 737 deletions(-) diff --git a/src/chat/chat_sync.c b/src/chat/chat_sync.c index 6f729fd6..ee9778a4 100644 --- a/src/chat/chat_sync.c +++ b/src/chat/chat_sync.c @@ -1440,7 +1440,7 @@ static void cs_handle_join_ready(struct chat_sync* cs, uint64_t peer, struct TOPO_GROUP* g = topo_groups_find(cs->inst->topo_groups, group_id); if (g && g->group_type == TOPO_GROUP_TYPE_CHAT) { struct ETCP_CONN* c = cs_find_conn_for_node(cs->inst, peer); - if (c) topo_group_new_conn(g, c); + if (c) topo_group_connect_node_once(g, c->peer_node_id); } } chat_core_ensure_channel_ready(cs->inst, ch_id); diff --git a/src/chat/member_sync.c b/src/chat/member_sync.c index 3e0f23de..d99f34d0 100644 --- a/src/chat/member_sync.c +++ b/src/chat/member_sync.c @@ -1,3 +1,4 @@ +#include "../routing_layer/topo_group_connect.h" #include "member_sync.h" #include "../routing_layer/topo_node_sqlite.h" @@ -706,7 +707,7 @@ static void member_committed(struct UTUN_INSTANCE* inst, const char* ns, uint64_ merkle_sync_changed(inst, ns); struct TOPO_GROUP* g = inst->topo_groups ? topo_groups_find(inst->topo_groups, strtoull(ns, NULL, 10)) : NULL; struct ETCP_CONN* conn = instance_find_conn(inst, node_id); - if (g && g->group_type == TOPO_GROUP_TYPE_CHAT && conn) topo_group_new_conn(g, conn); + if (g && g->group_type == TOPO_GROUP_TYPE_CHAT && conn) topo_group_connect_node_once(g, node_id); } int member_sync_apply_record(struct UTUN_INSTANCE* inst, const char* ch_id, diff --git a/src/routing_layer/topo_group.c b/src/routing_layer/topo_group.c index 12fa4bd4..28c62f23 100644 --- a/src/routing_layer/topo_group.c +++ b/src/routing_layer/topo_group.c @@ -71,6 +71,7 @@ static void topo_group_peer_free(struct TOPO_GROUP_CONN_ITEM* peer) { } if (peer->handle) node_conn_direct_close(peer->handle); topo_recovery_changed(peer->group); + topo_group_connect_changed(peer->group); } static void topo_group_pending_remove(struct TOPO_GROUP_CONN_ITEM* peer) { @@ -93,7 +94,10 @@ static void topo_group_peer_event(struct NODE_CONN_DIRECT* handle, enum ncd_even DEBUG_WARN(DEBUG_CATEGORY_BGP, "group peer transport timeout peer=%016llx group=%016llx", (unsigned long long)peer->node_id, (unsigned long long)peer->group->group_id); if (peer->conn) topo_group_remove_conn(peer->group, peer->conn, TOPO_REMOVE_TRANSPORT_DOWN); - else topo_group_pending_remove(peer); + else if ((event == NCD_EVENT_DOWN || event == NCD_EVENT_TIMEOUT) && peer->requests && node_conn_direct_get_conn(handle)) { + peer->transport_down = 1; peer->retained = 0; topo_recovery_changed(peer->group); + topo_group_connect_changed(peer->group); + } else topo_group_pending_remove(peer); } int topo_group_peer_open(struct TOPO_GROUP* group, uint64_t node_id, struct TOPO_PEER_REQUEST** out) { @@ -125,11 +129,22 @@ int topo_group_peer_open(struct TOPO_GROUP* group, uint64_t node_id, struct TOPO enum topo_peer_phase topo_group_peer_phase(const struct TOPO_PEER_REQUEST* request) { struct TOPO_GROUP_CONN_ITEM* peer = request ? request->peer : NULL; - if (!peer || peer->tx_failed) return TOPO_PEER_FAILED; + if (!peer || peer->tx_failed || peer->transport_down) return TOPO_PEER_FAILED; if (!peer->conn || !peer->conn->links_up) return TOPO_PEER_CONNECTING; return topo_group_peer_ready(peer->group, peer->node_id) ? TOPO_PEER_READY : TOPO_PEER_SYNCING; } +struct ETCP_CONN* topo_group_peer_conn(const struct TOPO_PEER_REQUEST* request) { + return request && request->peer ? node_conn_direct_get_conn(request->peer->handle) : NULL; +} + +void topo_group_peer_leave(struct TOPO_GROUP* group, uint64_t node_id) { + struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, node_id); + if (!peer) return; + if (peer->conn) topo_group_leave_peer(group, peer->handle); + else topo_group_pending_remove(peer); +} + uint64_t topo_group_peer_progress(const struct TOPO_PEER_REQUEST* request) { return request && request->peer ? request->peer->progress : 0; } @@ -153,10 +168,13 @@ static void topo_group_peer_progressed(struct TOPO_GROUP* group, uint64_t node_i struct TOPO_GROUP_CONN_ITEM* peer = topo_group_peer(group, node_id); if (peer) peer->progress = get_time_tb(); topo_recovery_changed(group); + topo_group_connect_changed(group); } static void topo_group_log_exchange(struct TOPO_GROUP* group, struct TOPO_GROUP_CONN_ITEM* peer) { - if (topo_group_peer_ready(group, peer->node_id)) peer->retained = 1; + if (topo_group_peer_ready(group, peer->node_id)) { + peer->retained = 1; topo_group_connect_on_up(group, peer->conn); + } topo_group_peer_progressed(group, peer->node_id); DEBUG_INFO(DEBUG_CATEGORY_BGP, "group exchange: group=%016llx peer=%016llx exchange=%016llx sent=%u received=%u ready=%d", (unsigned long long)group->group_id, (unsigned long long)peer->conn->peer_node_id, @@ -192,7 +210,7 @@ static int topo_send_control(struct ETCP_CONN* conn, const void* data, size_t si } static int topo_group_reset_exchange(struct TOPO_GROUP_CONN_ITEM* peer) { - topo_tx_clear(peer); peer->tx_failed = 0; + topo_tx_clear(peer); peer->tx_failed = 0; peer->transport_down = 0; peer->accepted = 0; peer->table_received = 0; peer->table_sent = 0; peer->peer_epoch = 0; uint64_t epoch; if (random_bytes((uint8_t*)&epoch, sizeof(epoch)) != 0 || !epoch) { @@ -329,10 +347,10 @@ static void topo_group_broadcast_withdraw(struct TOPO_GROUP* group, uint64_t nod // Приём пакетов // ============================================================================ -static struct NODE_CONN_DIRECT* topo_config_handle(struct UTUN_INSTANCE* instance, uint64_t node_id) { +static int topo_configured(struct UTUN_INSTANCE* instance, uint64_t node_id) { for (struct CONFIG_CONN_HANDLE* owner = instance->config_conn_handles; owner; owner = owner->next) - if (owner->node_id == node_id) return owner->handle; - return NULL; + if (owner->node_id == node_id) return 1; + return 0; } static int topo_group_has_sender(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { @@ -473,7 +491,7 @@ static void topo_group_conn_status(struct ETCP_CONN* conn, int status, void* arg if (status == ETCP_CONN_STATUS_UP && conn->peer_node_id) { /* UTUN присоединяет пир через UP принадлежащего ему NCD handle. * Сам по себе транспортный UP не означает участия в UTUN-группе. */ - if (!topo_config_handle(conn->instance, conn->peer_node_id)) topo_group_send_resync(conn); + if (!topo_configured(conn->instance, conn->peer_node_id)) topo_group_send_resync(conn); /* CHAT: только члены каналов (peers_* в БД). */ sqlite3* db = conn->instance ? conn->instance->topo_sqlite_db : NULL; uint64_t* chs = NULL; int chn = 0; @@ -780,8 +798,7 @@ struct TOPO_GROUP* topo_groups_create_group(struct TOPO_GROUPS* g, uint64_t grou } u_free(chs); } else if (group_type == TOPO_GROUP_TYPE_UTUN && cqe->conn->links_up) { - struct NODE_CONN_DIRECT* handle = topo_config_handle(g->instance, cqe->peer_node_id); - if (handle) topo_group_join_peer(group, handle); + if (topo_configured(g->instance, cqe->peer_node_id)) topo_group_new_conn(group, cqe->conn); } } entry = entry->hash_next; @@ -909,7 +926,11 @@ void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn, en if (pending && !pending->conn && node_conn_direct_get_conn(pending->handle) == conn) { DEBUG_INFO(DEBUG_CATEGORY_BGP, "group pending peer removed: peer=%016llx group=%016llx reason=%d", (unsigned long long)pending->node_id, (unsigned long long)group->group_id, reason); - topo_group_pending_remove(pending); return; + if (reason == TOPO_REMOVE_TRANSPORT_DOWN && pending->requests) { + pending->transport_down = 1; pending->retained = 0; + topo_recovery_changed(group); topo_group_connect_changed(group); + } else topo_group_pending_remove(pending); + return; } bool found_in_list = false; @@ -983,9 +1004,13 @@ void topo_group_remove_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn, en while (e) { struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data; if (item->conn == conn) { - queue_remove_data(group->senders_list, e); - topo_group_peer_free(item); - queue_entry_free(e); + if (reason == TOPO_REMOVE_TRANSPORT_DOWN && item->requests && node_conn_direct_get_conn(item->handle) == conn) { + topo_tx_clear(item); item->conn = NULL; item->transport_down = 1; item->retained = 0; + item->accepted = 0; item->table_sent = 0; item->table_received = 0; + topo_recovery_changed(group); topo_group_connect_changed(group); + } else { + queue_remove_data(group->senders_list, e); topo_group_peer_free(item); queue_entry_free(e); + } DEBUG_INFO(DEBUG_CATEGORY_BGP, "topo_group_remove_conn: REMOVED conn=%p peer=%016llx from senders_list grp=%016llx ch=%s", (void*)conn, (unsigned long long)conn->peer_node_id, (unsigned long long)group->group_id, group->channel_id); @@ -1163,6 +1188,7 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from struct TOPO_GROUP_NODE rejected = { .subnets = new_subnets }; topo_nodeq_free_group_fields(group->instance->topo_groups, &rejected); u_free(new_hop_list); + topo_group_peer_progressed(group, from->peer_node_id); return 0; } uint64_t incoming_timestamp = new_ni->timestamp; @@ -1388,6 +1414,7 @@ static void topo_tx_fail(struct TOPO_GROUP_CONN_ITEM* peer) { (unsigned long long)peer->group->group_id, (unsigned long long)peer->node_id, (unsigned long long)peer->local_epoch); topo_tx_clear(peer); peer->tx_failed = 1; peer->table_sent = 0; topo_recovery_changed(peer->group); + topo_group_connect_changed(peer->group); } static void topo_tx_wait(void* arg); @@ -1694,7 +1721,7 @@ static void topo_group_handle_join_group(struct TOPO_GROUP* group, struct ETCP_C static void topo_group_handle_resync(struct UTUN_INSTANCE* instance, struct ETCP_CONN* conn) { if (!instance || !conn || !conn->instance || !instance->topo_groups) return; /* общий сигнал от пассивной стороны: рестартуем обмен, если мы для этого пира VPN-клиент */ - if (!topo_config_handle(instance, conn->peer_node_id)) return; + if (!topo_configured(instance, conn->peer_node_id)) return; DEBUG_INFO(DEBUG_CATEGORY_BGP, "handle_resync: from %s — re-announce groups", conn->log_name); struct ll_entry* fe = instance->topo_groups->group_list->head; while (fe) { diff --git a/src/routing_layer/topo_group.h b/src/routing_layer/topo_group.h index 70aa12ce..ee1d9c06 100644 --- a/src/routing_layer/topo_group.h +++ b/src/routing_layer/topo_group.h @@ -179,7 +179,7 @@ struct TOPO_GROUP_CONN_ITEM { uint8_t retained; // участие запрошено постоянным владельцем или достигло READY uint64_t local_epoch, peer_epoch; // поколения сторон текущего присоединения uint8_t accepted; // JOIN_ACCEPT подтвердил local_epoch - uint8_t tx_failed; + uint8_t tx_failed, transport_down; struct topo_tx_item *tx_head, *tx_tail; struct queue_waiter_handle tx_waiter; void* tx_wake; @@ -306,11 +306,16 @@ enum topo_peer_phase { TOPO_PEER_FAILED, TOPO_PEER_CONNECTING, TOPO_PEER_SYNCING * владеет транспортом. Проверка CHAT-членства общая с входящим JOIN_GROUP. * close отменяет только свой запрос; последний запрос отменяет незавершённую * попытку, если её не удерживает другой инициатор. READY сохраняется в группе. + * DOWN/TIMEOUT сообщает FAILED, но живой запрос сохраняет интерес к NCD: поздний + * UP возобновляет сессию. Закрытие последнего запроса отменяет неудачную попытку. + * CLOSED/leave/удаление группы отсоединяет запрос окончательно. * Запрос остаётся валиден после удаления группы (phase=FAILED), caller обязан * закрыть его. Все операции выполняются в потоке единственного uasync. * progress — timebase последнего принятого NODEINFO/шага обмена, не polling time. */ int topo_group_peer_open(struct TOPO_GROUP* group, uint64_t node_id, struct TOPO_PEER_REQUEST** out); void topo_group_peer_close(struct TOPO_PEER_REQUEST* request); +struct ETCP_CONN* topo_group_peer_conn(const struct TOPO_PEER_REQUEST* request); /* borrowed, including CONNECTING */ +void topo_group_peer_leave(struct TOPO_GROUP* group, uint64_t node_id); enum topo_peer_phase topo_group_peer_phase(const struct TOPO_PEER_REQUEST* request); uint64_t topo_group_peer_progress(const struct TOPO_PEER_REQUEST* request); diff --git a/src/routing_layer/topo_group_connect.c b/src/routing_layer/topo_group_connect.c index 8d8bd207..d0266045 100644 --- a/src/routing_layer/topo_group_connect.c +++ b/src/routing_layer/topo_group_connect.c @@ -1,679 +1,283 @@ -/* - * topo_group_connect.c — авто-подключение к узлам группы (бесконечный цикл до цели) - * - * Phase 1: одновременный запуск node_conn_direct_open для пиров с connected=1. - * Таймаут = TGC_DIRECT_TIMEOUT_MS. Если цель достигнута → done. - * - * Phase 2: последовательный перебор (supernode, затем public/EIM адреса). - * После каждого TIMEOUT — пауза перед следующей попыткой. - * - * Phase 3: последовательный перебор локальных/strict NAT адресов. - * После каждого TIMEOUT — пауза перед следующей попыткой. - * - * После исчерпания Phase 3 → пауза → cycle_restart → Phase 1 (бесконечно). - * - * Цель (tgc_goal_reached): подключены хотя бы к одному суперузлу (node_type==4) - * ЛИБО к 3+ не-мобильным клиентам (device_type != MOBILE, не суперузлам). - * Пока цель не достигнута, успешное подключение НЕ останавливает цикл — - * продолжается перебор следующих кандидатов. - * - * Пауза: активный режим — 1с; фоновый (Android standby) — standby_wait - * (burst → сразу, sleep → до следующего burst). - * При DOWN, если цель перестала выполняться — немедленный cycle_restart - * (уже установленные соединения не трогаются). - * При повторном topo_group_connect_init — destroy + fresh init. - * - * Счётчик active_conn_count — число уникальных peer-узлов с реальными соединениями. - * UP: инкремент только для первого соединения к peer_node_id. - * DOWN: декремент только если нет других conn к peer_node_id И нет indirect-путей. - */ - +/* Политика CHAT-подключений. Транспортом и JOIN владеет групповая сессия. + * Сначала параллельно восстанавливаем connected-пиров, затем последовательно + * пробуем supernode/public/local. Незавершённые попытки — отменяемые запросы. */ +#include +#include #include "topo_group_connect.h" #include "topo_group.h" -#include "topo_node.h" #include "topo_node_sqlite.h" +#include "topo_recovery.h" #include "../utun_instance.h" -#include "../transport_layer/node_conn_direct.h" #include "../lib/debug_config.h" #include "../lib/mem.h" -#include "../lib/u_async.h" -#include "../lib/ll_queue.h" -#include "etcp.h" -#include "etcp_connections.h" +#include "../lib/platform_compat.h" #include "../chat/chat_event.h" +#include "etcp.h" #ifdef UTUN_HAVE_STANDBY #include "standby.h" #endif -/* ─── внутренние константы ─── */ -#define TGC_ID "topo_group_connect" -#define TGC_DIRECT_TIMEOUT_MS 2000 -#define TGC_PHASE_ONE 0 -#define TGC_PHASE_TWO 1 -#define TGC_PHASE_THREE 2 -#define TGC_PHASE_DONE 3 -#define TGC_PHASE_PAUSE 4 -#define TGC_MAX_HANDLES 128 -#define TGC_PAUSE_ACTIVE_TB 10000 /* 1s пауза в активном режиме */ -#define TGC_GOAL_SUPERNODES 1 /* цель: подключены хотя бы к одному суперузлу */ -#define TGC_GOAL_NONMOBILE 3 /* цель: подключены к 3+ не-мобильным клиентам (не суперузлам) */ +#define TGC_PARALLEL 128 +#define TGC_PAUSE_TB 10000 + +struct tgc_candidate { + uint64_t node_id, started; + struct TOPO_PEER_REQUEST* request; + uint8_t priority, tried, manual; +}; struct TOPO_GROUP_CONNECT { struct TOPO_GROUP* group; - void* phase_timer; - void* pause_timer; /* uasync timeout handle (активный режим) */ - void* pause_wait; /* standby_wait handle (фоновый режим, Android) */ - uint8_t phase; - int pending; - int connected_count; - int active_conn_count; - int cursor; - int tried_super; - uint64_t* candidate_ids; - int candidate_count; - struct NODE_CONN_DIRECT* handles[TGC_MAX_HANDLES]; - int handle_count; - struct NODE_CONN_DIRECT* cur_handle; /* handle текущей попытки Phase 2/3 (для закрытия на TIMEOUT) */ + struct tgc_candidate* candidates; + size_t count; + void *wake, *timer, *pause_timer, *pause_wait; + uint8_t automatic, cycle_done, paused, reload_after_pause; }; -static void tgc_phase2_try_next(struct TOPO_GROUP_CONNECT* gc); -static void tgc_phase3_try_next(struct TOPO_GROUP_CONNECT* gc); -static void tgc_phase1_timeout(void* arg); -static void tgc_cycle_restart(struct TOPO_GROUP_CONNECT* gc); -static void tgc_start_pause(struct TOPO_GROUP_CONNECT* gc, void (*cb)(void*), const char* label); -static int tgc_goal_reached(struct TOPO_GROUP_CONNECT* gc); - -static void tgc_callback(struct NODE_CONN_DIRECT* h, enum ncd_event event, void* arg); - -/* ─── нотификация GUI о списке узлов в процессе подключения ─── */ - -/** - * Отправляет в GUI событие CHAT_EVT_CONNECTING_NODES — список узлов, к которым - * сейчас идёт попытка подключения. GUI по нему рисует жёлтый кружок у этих узлов. - * Пустой список (count=0) означает «сейчас ни к кому не подключаемся» — - * все жёлтые кружки снимаются. - */ -static void tgc_notify_connecting(struct TOPO_GROUP_CONNECT* gc, const uint64_t* ids, int count) { - size_t cl = strlen(gc->group->channel_id); if (cl > 255) cl = 255; - size_t sz = 1 + cl + 2 + (size_t)count * 8; - uint8_t* buf = u_malloc(sz); if (!buf) return; - uint8_t* p = buf; - *p++ = (uint8_t)cl; memcpy(p, gc->group->channel_id, cl); p += cl; - uint16_t c = (uint16_t)count; memcpy(p, &c, 2); p += 2; - for (int i = 0; i < count; i++) { memcpy(p, &ids[i], 8); p += 8; } - chat_event_post(gc->group->instance, CHAT_EVT_CONNECTING_NODES, buf, (int)sz); - u_free(buf); - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: notify_connecting ch=%s count=%d", - TGC_ID, gc->group->channel_id, count); +static void tgc_step(void* arg); + +static void tgc_notify(struct TOPO_GROUP_CONNECT* gc) { + size_t length = strlen(gc->group->channel_id), count = 0; + for (size_t i = 0; i < gc->count; i++) if (gc->candidates[i].request) count++; + size_t size = 1 + length + 2 + count * 8; + uint8_t* data = u_malloc(size); + if (!data) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "connecting notification allocation failed"); return; } + data[0] = (uint8_t)length; memcpy(data + 1, gc->group->channel_id, length); + uint16_t n = (uint16_t)count; memcpy(data + 1 + length, &n, 2); + uint8_t* out = data + 3 + length; + for (size_t i = 0; i < gc->count; i++) if (gc->candidates[i].request) { + memcpy(out, &gc->candidates[i].node_id, 8); out += 8; + } + chat_event_post(gc->group->instance, CHAT_EVT_CONNECTING_NODES, data, (int)size); u_free(data); } -/* ═══════════════════════════════════════════════════════════════════════ - * Условие достижения цели авто-подключения - * ══════════════════════════════════════════════════════════════════════ */ - -/* Тип устройства пира из handshake (peer_device_type на первом UP-линке). */ -static uint8_t tgc_peer_device_type(struct ETCP_CONN* conn) { - if (!conn || !conn->links) return CLIENT_TYPE_SERVER; - return conn->links->peer_device_type; +int topo_group_connect_active_count(struct TOPO_GROUP* group) { + int count = 0; + for (struct ll_entry* e = group && group->senders_list ? group->senders_list->head : NULL; e; e = e->next) { + struct TOPO_GROUP_CONN_ITEM* peer = (struct TOPO_GROUP_CONN_ITEM*)e->data; + if (topo_group_peer_ready(group, peer->node_id)) count++; + } + return count; } -/** - * Достигнута ли цель подключения к группе: - * - подключены хотя бы к одному суперузлу (peers node_type==4), ЛИБО - * - подключены к 3+ не-мобильным клиентам (device_type != MOBILE и не суперузел). - * - * Сканирует senders_list (активные BGP-пиры) с дедупом по node_id. - * Состояние не хранится — пересчитывается на лету в редких точках принятия - * решения (таймаут фазы, каждый успех, каждый down). - */ -static int tgc_goal_reached(struct TOPO_GROUP_CONNECT* gc) { - if (!gc || !gc->group || !gc->group->senders_list) return 0; +static int tgc_goal(struct TOPO_GROUP_CONNECT* gc) { + int super = 0, other = 0; struct TOPO_GROUP* group = gc->group; - sqlite3* db = group->instance->topo_sqlite_db; - - uint64_t seen[TGC_MAX_HANDLES]; int seen_count = 0; - int supernode_count = 0, nonmobile_count = 0; - - struct ll_entry* e = group->senders_list->head; - while (e) { - struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data; - struct ETCP_CONN* conn = item ? item->conn : NULL; - if (conn) { - int dup = 0; - for (int i = 0; i < seen_count; i++) if (seen[i] == conn->peer_node_id) { dup = 1; break; } - if (!dup && seen_count < TGC_MAX_HANDLES) { - seen[seen_count++] = conn->peer_node_id; - int is_super = db && topo_node_sqlite_get_node_type(db, group->channel_id, conn->peer_node_id) == 4; - int is_mobile = tgc_peer_device_type(conn) == CLIENT_TYPE_MOBILE; - if (is_super) supernode_count++; - else if (!is_mobile) nonmobile_count++; - } - } - e = e->next; + for (struct ll_entry* e = group->senders_list->head; e; e = e->next) { + struct TOPO_GROUP_CONN_ITEM* peer = (struct TOPO_GROUP_CONN_ITEM*)e->data; + if (!topo_group_peer_ready(group, peer->node_id)) continue; + if (topo_node_sqlite_get_node_type(group->instance->topo_sqlite_db, group->channel_id, peer->node_id) == 4) super++; + else if (!peer->conn->links || peer->conn->links->peer_device_type != CLIENT_TYPE_MOBILE) other++; } - - int reached = (supernode_count >= TGC_GOAL_SUPERNODES) || (nonmobile_count >= TGC_GOAL_NONMOBILE); - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: goal ch=%s super=%d/%d nonmobile=%d/%d → %s", - TGC_ID, group->channel_id, supernode_count, TGC_GOAL_SUPERNODES, - nonmobile_count, TGC_GOAL_NONMOBILE, reached ? "reached" : "not reached"); - return reached; + return super >= 1 || other >= 3; } -/* ═══════════════════════════════════════════════════════════════════════ - * Жизненный цикл - * ══════════════════════════════════════════════════════════════════════ */ - -/** - * Запуск авто-подключения к узлам CHAT-группы. - * Если модуль уже запущен — останавливает его и стартует заново с нуля. - * Начинает с Phase 1: параллельные попытки ко всем пирам, что были connected - * в прошлой сессии (из БД). Если таких нет — сразу переходит к Phase 2. - */ -int topo_group_connect_init(struct TOPO_GROUP* group) { - if (!group || group->group_type != TOPO_GROUP_TYPE_CHAT || !group->channel_id[0]) - return -1; - if (group->connect) topo_group_connect_destroy(group); - struct TOPO_GROUP_CONNECT* gc = u_calloc(1, sizeof(*gc)); - if (!gc) return -1; - gc->group = group; gc->tried_super = 0; - group->connect = gc; +static struct tgc_candidate* tgc_add(struct TOPO_GROUP_CONNECT* gc, uint64_t id, uint8_t priority, int manual) { + for (size_t i = 0; i < gc->count; i++) if (gc->candidates[i].node_id == id) return &gc->candidates[i]; + struct tgc_candidate* list = u_realloc(gc->candidates, (gc->count + 1) * sizeof(*list)); + if (!list) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group candidate allocation failed"); return NULL; } + gc->candidates = list; + list[gc->count] = (struct tgc_candidate){ .node_id = id, .priority = priority, .manual = manual }; + return &list[gc->count++]; +} - uint64_t* ids = NULL; int count = 0; - sqlite3* db = group->instance->topo_sqlite_db; - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: init ch=%s db=%p", TGC_ID, group->channel_id, (void*)db); - if (!db || topo_node_sqlite_get_connected_peers(db, group->channel_id, &ids, &count) != 0 || count == 0) { - if (ids) { u_free(ids); ids = NULL; } - gc->phase = TGC_PHASE_TWO; - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: ch=%s no connected peers → Phase 2", TGC_ID, group->channel_id); - tgc_phase2_try_next(gc); - return 0; +static void tgc_load(struct TOPO_GROUP_CONNECT* gc) { + size_t retained = 0; + for (size_t i = 0; i < gc->count; i++) { + struct tgc_candidate* c = &gc->candidates[i]; + if (c->manual && c->request) gc->candidates[retained++] = *c; + else topo_group_peer_close(c->request); } - - gc->phase = TGC_PHASE_ONE; - int launched = 0; - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase 1 launching %d connects ch=%s grp=%016llx", TGC_ID, count, group->channel_id, (unsigned long long)group->group_id); - for (int i = 0; i < count && gc->handle_count < TGC_MAX_HANDLES; i++) { - if (ids[i] == group->instance->node_id) { DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: Phase1 skip self 0x%016llx", TGC_ID, (unsigned long long)ids[i]); continue; } - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: Phase1 connect to 0x%016llx (%d/%d)", TGC_ID, (unsigned long long)ids[i], i + 1, count); - node_conn_direct_open(group->instance, ids[i], tgc_callback, gc, - &gc->handles[gc->handle_count++], NULL); - launched++; + gc->count = retained; gc->cycle_done = 0; + sqlite3* db = gc->group->instance->topo_sqlite_db; + for (int priority = 0; priority < 4; priority++) { + uint64_t* ids = NULL; int count = 0, rc; + if (priority == 0) rc = topo_node_sqlite_get_connected_peers(db, gc->group->channel_id, &ids, &count); + else if (priority == 1) rc = topo_node_sqlite_get_supernode_peers(db, gc->group->channel_id, &ids, &count); + else if (priority == 2) rc = topo_node_sqlite_get_public_peers(db, gc->group->channel_id, &ids, &count, gc->group->instance->node_id); + else rc = topo_node_sqlite_get_local_peers(db, gc->group->channel_id, &ids, &count, gc->group->instance->node_id); + if (rc < 0) DEBUG_WARN(DEBUG_CATEGORY_BGP, "cannot load group candidates group=%016llx priority=%d", + (unsigned long long)gc->group->group_id, priority); + else for (int i = 0; i < count; i++) + if (ids[i] != gc->group->instance->node_id && !tgc_add(gc, ids[i], priority, 0)) break; + u_free(ids); } - gc->pending = launched; - tgc_notify_connecting(gc, ids, gc->handle_count); - u_free(ids); - - gc->phase_timer = uasync_set_timeout(group->instance->ua, - TGC_DIRECT_TIMEOUT_MS * 10, gc, tgc_phase1_timeout, "tgc_phase1"); - return 0; + DEBUG_INFO(DEBUG_CATEGORY_BGP, "group connect cycle: group=%016llx candidates=%zu", + (unsigned long long)gc->group->group_id, gc->count); } -/** - * Полная остановка авто-подключения: отменяет все таймеры, закрывает все - * открытые соединения и освобождает память. Вызывается при удалении группы - * и при перезапуске (init/restart). - */ -void topo_group_connect_destroy(struct TOPO_GROUP* group) { - struct TOPO_GROUP_CONNECT* gc = group->connect; - if (!gc) return; - group->connect = NULL; - if (gc->phase_timer) { uasync_cancel_timeout(group->instance->ua, gc->phase_timer); gc->phase_timer = NULL; } - if (gc->pause_timer) { uasync_cancel_timeout(group->instance->ua, gc->pause_timer); gc->pause_timer = NULL; } +static void tgc_cancel_pause(struct TOPO_GROUP_CONNECT* gc) { + if (gc->pause_timer) uasync_cancel_timeout(gc->group->instance->ua, gc->pause_timer); #ifdef UTUN_HAVE_STANDBY - if (gc->pause_wait) { standby_wait_cancel(gc->pause_wait); gc->pause_wait = NULL; } + if (gc->pause_wait) standby_wait_cancel(gc->pause_wait); #endif - for (int i = 0; i < gc->handle_count; i++) node_conn_direct_close(gc->handles[i]); - if (gc->cur_handle) { node_conn_direct_close(gc->cur_handle); gc->cur_handle = NULL; } - u_free(gc->candidate_ids); - u_free(gc); -} - -/** - * Внешний перезапуск авто-подключения (destroy + init). Используется, например, - * при смене сети, когда надо начать поиск заново. Не трогает модуль, если уже - * есть живые соединения. - */ -void topo_group_connect_restart(struct TOPO_GROUP* group) { - struct TOPO_GROUP_CONNECT* gc = group->connect; - if (!gc) { DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: restart ch=%s no gc → init", TGC_ID, group->channel_id); topo_group_connect_init(group); return; } - if (gc->active_conn_count > 0) { DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: restart ch=%s skip active=%d", TGC_ID, group->channel_id, gc->active_conn_count); return; } - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: restart ch=%s phase=%d", TGC_ID, group->channel_id, gc->phase); - topo_group_connect_destroy(group); - topo_group_connect_init(group); + gc->pause_timer = NULL; gc->pause_wait = NULL; gc->paused = 0; } -/** - * Сколько сейчас живых соединений с пирами группы. Нужно внешней логике - * (например, chat_sync) для решения, запускать ли переподключение. - */ -int topo_group_connect_active_count(struct TOPO_GROUP* group) { - struct TOPO_GROUP_CONNECT* gc = group->connect; - return gc ? gc->active_conn_count : 0; +static void tgc_resume(void* arg) { + struct TOPO_GROUP_CONNECT* gc = arg; + gc->pause_timer = NULL; gc->pause_wait = NULL; gc->paused = 0; + if (gc->reload_after_pause) tgc_load(gc); + topo_group_connect_changed(gc->group); } -/* ═══════════════════════════════════════════════════════════════════════ - * ON UP / DOWN - * ══════════════════════════════════════════════════════════════════════ */ - -/** - * Обработка появления соединения с пиром. Увеличивает счётчик активных и - * помечает connected=1 в БД (чтобы при следующем старте узел попал в Phase 1). - * Пропускает, если пир уже подключён через другое соединение. - */ -void topo_group_connect_on_up(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { - if (!group || group->group_type != TOPO_GROUP_TYPE_CHAT || !group->channel_id[0] || !conn) return; - uint64_t peer = conn->peer_node_id; - if (peer == group->instance->node_id) return; - struct TOPO_GROUP_CONNECT* gc = group->connect; - if (!gc) return; - - struct ll_entry* e = group->senders_list ? group->senders_list->head : NULL; - while (e) { - struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data; - if (item->conn && item->conn != conn && item->conn->peer_node_id == peer) { - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: UP peer=0x%016llx already connected via other conn, skip count", - TGC_ID, (unsigned long long)peer); - return; - } - e = e->next; +static void tgc_pause(struct TOPO_GROUP_CONNECT* gc, int reload) { + gc->paused = 1; gc->reload_after_pause = reload; +#ifdef UTUN_HAVE_STANDBY + if (standby_is_enabled()) { + gc->pause_wait = standby_wait(gc, tgc_resume); + if (gc->pause_wait) return; + DEBUG_WARN(DEBUG_CATEGORY_BGP, "cannot wait for standby; using group retry timer"); } - - gc->active_conn_count++; - sqlite3* db = group->instance->topo_sqlite_db; - if (db) topo_node_sqlite_set_connected(db, group->channel_id, peer, 1); - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: UP peer=0x%016llx ch=%s grp=%016llx active=%d", - TGC_ID, (unsigned long long)peer, group->channel_id, (unsigned long long)group->group_id, gc->active_conn_count); +#endif + gc->pause_timer = uasync_set_timeout(gc->group->instance->ua, TGC_PAUSE_TB, gc, tgc_resume, "group_connect_pause"); + if (!gc->pause_timer) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group retry timer allocation failed"); gc->paused = 0; } } -/** - * Обработка обрыва соединения. Уменьшает счётчик активных и помечает - * connected=0 в БД. Если живых соединений не осталось — немедленно - * перезапускает цикл поиска с Phase 1. - */ -void topo_group_connect_on_down(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { - if (!group || group->group_type != TOPO_GROUP_TYPE_CHAT || !group->channel_id[0] || !conn) return; - uint64_t peer = conn->peer_node_id; - if (peer == group->instance->node_id) return; - struct TOPO_GROUP_CONNECT* gc = group->connect; - if (!gc) return; - - struct ll_entry* e = group->senders_list ? group->senders_list->head : NULL; - while (e) { - struct TOPO_GROUP_CONN_ITEM* item = (struct TOPO_GROUP_CONN_ITEM*)e->data; - if (item->conn && item->conn != conn && item->conn->peer_node_id == peer) { - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: DOWN peer=0x%016llx still has other direct conn, skip", - TGC_ID, (unsigned long long)peer); - return; - } - e = e->next; - } - - struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(group, peer); - if (nq && nq->paths) { - struct ll_entry* pe = nq->paths->head; - while (pe) { - struct TOPO_NODEPATH* path = (struct TOPO_NODEPATH*)pe; - if (path->conn && path->conn != conn) { - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: DOWN peer=0x%016llx has indirect path, skip", - TGC_ID, (unsigned long long)peer); - return; - } - pe = pe->next; - } - } - - gc->active_conn_count--; - sqlite3* db = group->instance->topo_sqlite_db; - if (db) topo_node_sqlite_set_connected(db, group->channel_id, peer, 0); - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: DOWN peer=0x%016llx ch=%s active=%d", - TGC_ID, (unsigned long long)peer, group->channel_id, gc->active_conn_count); - if (!tgc_goal_reached(gc)) { - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: DOWN goal not reached — cycle restart ch=%s", TGC_ID, group->channel_id); - tgc_cycle_restart(gc); - } +static int tgc_open(struct TOPO_GROUP_CONNECT* gc, struct tgc_candidate* c) { + c->tried = 1; c->started = get_time_tb(); + int result = topo_group_peer_open(gc->group, c->node_id, &c->request); + DEBUG_INFO(DEBUG_CATEGORY_BGP, "group connect attempt: group=%016llx node=%016llx priority=%u manual=%u result=%d", + (unsigned long long)gc->group->group_id, (unsigned long long)c->node_id, c->priority, c->manual, result); + if (result == 0) topo_group_connect_changed(gc->group); + return result; } -/* ═══════════════════════════════════════════════════════════════════════ - * Единый коллбэк для всех фаз - * ══════════════════════════════════════════════════════════════════════ */ - -/** - * Единый коллбэк на результат каждой попытки подключения (успех/таймаут/обрыв). - * В Phase 1 — просто считает результаты (итог подводит tgc_phase1_timeout). - * В Phase 2/3 — успех проверяет цель: достигнута → done, иначе следующий кандидат; - * провал закрывает handle (чтобы не оставить залипшее состояние) и ставит паузу. - */ -/* Передача владения: после topo_group_new_conn (который взял свой NCD-handle) - * закрываем connect-фазу handle, чтобы владельцем conn остался один BGP-хэндл - * (в senders_list), без двойного счёта. */ -static void tgc_release_handle(struct TOPO_GROUP_CONNECT* gc, struct NODE_CONN_DIRECT* h) { - if (!gc || !h) return; - for (int i = 0; i < gc->handle_count; i++) { - if (gc->handles[i] == h) { gc->handles[i] = NULL; node_conn_direct_close(h); return; } - } - if (gc->cur_handle == h) { gc->cur_handle = NULL; node_conn_direct_close(h); return; } - node_conn_direct_close(h); +static uint64_t tgc_deadline(struct tgc_candidate* c, enum topo_peer_phase phase) { + return phase == TOPO_PEER_CONNECTING ? c->started + TOPO_RECOVERY_CONNECT_TIMEOUT_MS * 10ULL : + topo_group_peer_progress(c->request) + TOPO_RECOVERY_SYNC_TIMEOUT_MS * 10ULL; } -static void tgc_callback(struct NODE_CONN_DIRECT* h, enum ncd_event event, void* arg) { - struct TOPO_GROUP_CONNECT* gc = (struct TOPO_GROUP_CONNECT*)arg; - struct ETCP_CONN* conn = node_conn_direct_get_conn(h); - uint64_t node_id = node_conn_direct_node_id(h); - - int ok = (event == NCD_EVENT_UP); +static void tgc_timeout(void* arg) { + struct TOPO_GROUP_CONNECT* gc = arg; gc->timer = NULL; tgc_step(gc); +} - if (ok) { - /* запуск BGP: узел соединён с группой (topo_group_new_conn берёт владение NCD-handle) */ - if (conn) topo_group_new_conn(gc->group, conn); - tgc_release_handle(gc, h); - } else { - // Каждая попытка даёт один результат. Поздние DOWN/CLOSED не должны уменьшать pending повторно. - struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(gc->group, node_id); - if (nq) nq->handle = NULL; - tgc_release_handle(gc, h); +static void tgc_step(void* arg) { + struct TOPO_GROUP_CONNECT* gc = arg; + if (gc->wake) { uasync_call_soon_cancel(gc->group->instance->ua, gc->wake); gc->wake = NULL; } + if (gc->timer) { uasync_cancel_timeout(gc->group->instance->ua, gc->timer); gc->timer = NULL; } + uint64_t now = get_time_tb(); + int changed = 0, active = 0, failed = 0; + for (size_t i = 0; i < gc->count; i++) { + struct tgc_candidate* c = &gc->candidates[i]; + if (!c->request) continue; + enum topo_peer_phase phase = topo_group_peer_phase(c->request); + if (phase == TOPO_PEER_READY || phase == TOPO_PEER_FAILED || now >= tgc_deadline(c, phase)) { + if (phase != TOPO_PEER_READY) { + DEBUG_WARN(DEBUG_CATEGORY_BGP, "group connect failed: group=%016llx node=%016llx phase=%d", + (unsigned long long)gc->group->group_id, (unsigned long long)c->node_id, phase); + if (!c->manual && c->priority) failed = 1; + } else DEBUG_INFO(DEBUG_CATEGORY_BGP, "group connect READY: group=%016llx node=%016llx", + (unsigned long long)gc->group->group_id, (unsigned long long)c->node_id); + struct TOPO_PEER_REQUEST* request = c->request; c->request = NULL; + topo_group_peer_close(request); changed = 1; + } else if (!c->manual) active++; } - - switch (gc->phase) { - case TGC_PHASE_ONE: - gc->pending--; - if (ok) gc->connected_count++; - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: Phase1 result node=0x%016llx %s pending=%d connected=%d", - TGC_ID, (unsigned long long)node_id, ok ? "OK" : "FAIL", gc->pending, gc->connected_count); - break; - case TGC_PHASE_TWO: - if (ok) { - if (tgc_goal_reached(gc)) { - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 connected to 0x%016llx — goal reached, done", TGC_ID, (unsigned long long)node_id); - gc->phase = TGC_PHASE_DONE; - tgc_notify_connecting(gc, NULL, 0); - } else { - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 connected to 0x%016llx — goal not reached, try next", TGC_ID, (unsigned long long)node_id); - tgc_phase2_try_next(gc); + if (gc->automatic && tgc_goal(gc)) { + tgc_cancel_pause(gc); + if (!gc->cycle_done) { + DEBUG_INFO(DEBUG_CATEGORY_BGP, "group connect goal reached: group=%016llx", (unsigned long long)gc->group->group_id); + for (size_t i = 0; i < gc->count; i++) if (!gc->candidates[i].manual && gc->candidates[i].request) { + struct TOPO_PEER_REQUEST* request = gc->candidates[i].request; gc->candidates[i].request = NULL; + topo_group_peer_close(request); changed = 1; } - } else { - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 0x%016llx FAIL — pause then try next", TGC_ID, (unsigned long long)node_id); - if (gc->cur_handle) { node_conn_direct_close(gc->cur_handle); gc->cur_handle = NULL; } - tgc_start_pause(gc, (void(*)(void*))tgc_phase2_try_next, "tgc_p2_try"); + gc->cycle_done = 1; } - break; - case TGC_PHASE_THREE: - if (ok) { - if (tgc_goal_reached(gc)) { - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase3 connected to 0x%016llx — goal reached, done", TGC_ID, (unsigned long long)node_id); - gc->phase = TGC_PHASE_DONE; - tgc_notify_connecting(gc, NULL, 0); - } else { - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase3 connected to 0x%016llx — goal not reached, try next", TGC_ID, (unsigned long long)node_id); - tgc_phase3_try_next(gc); + } else if (gc->automatic && !gc->paused) { + if (gc->cycle_done) { tgc_load(gc); active = 0; } + if (failed) tgc_pause(gc, 0); + else { + /* Priority 0 — parallel historical peers; all other attempts are sequential. */ + int available = 0; + for (size_t i = 0; i < gc->count; i++) { + struct tgc_candidate* c = &gc->candidates[i]; + if (c->manual || c->tried) continue; + available = 1; + if (active && (c->priority || active >= TGC_PARALLEL)) break; + if (tgc_open(gc, c) == 0) { active++; changed = 1; } } - } else { - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase3 0x%016llx FAIL — pause then try next", TGC_ID, (unsigned long long)node_id); - if (gc->cur_handle) { node_conn_direct_close(gc->cur_handle); gc->cur_handle = NULL; } - tgc_start_pause(gc, (void(*)(void*))tgc_phase3_try_next, "tgc_p3_try"); + if (!available && !active) tgc_pause(gc, 1); } - break; } -} - -/* ═══════════════════════════════════════════════════════════════════════ - * Пауза: 1s ACTIVE / 30s STANDBY - * ══════════════════════════════════════════════════════════════════════ */ - -/** - * Пауза между попытками подключения. Длительность зависит от активности - * приложения: 1 секунда когда приложение активно, 30 секунд когда в фоне - * (экран выключен). После паузы вызывает cb (следующая попытка или перезапуск - * цикла). На время паузы уведомляет GUI пустым списком (жёлтые кружки гаснут). - */ -static void tgc_start_pause(struct TOPO_GROUP_CONNECT* gc, void (*cb)(void*), const char* label) { - gc->phase = TGC_PHASE_PAUSE; - tgc_notify_connecting(gc, NULL, 0); -#ifdef UTUN_HAVE_STANDBY - if (standby_is_enabled()) { - /* фоновый режим: burst → standby_wait будит сразу, sleep → до следующего burst */ - gc->pause_wait = standby_wait(gc, cb); - gc->pause_timer = NULL; - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: pause (standby) for ch=%s", TGC_ID, gc->group->channel_id); - return; + uint64_t earliest = UINT64_MAX; + for (size_t i = 0; i < gc->count; i++) if (gc->candidates[i].request) { + uint64_t deadline = tgc_deadline(&gc->candidates[i], topo_group_peer_phase(gc->candidates[i].request)); + if (deadline < earliest) earliest = deadline; } -#endif - gc->pause_timer = uasync_set_timeout(gc->group->instance->ua, TGC_PAUSE_ACTIVE_TB, gc, cb, label); - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: pause %dms (active) for ch=%s", - TGC_ID, TGC_PAUSE_ACTIVE_TB / 10, gc->group->channel_id); -} - -/* ═══════════════════════════════════════════════════════════════════════ - * Phase 1 timeout - * ══════════════════════════════════════════════════════════════════════ */ - -/** - * Таймаут параллельных попыток Phase 1 (2 секунды). Закрывает неудачные - * handles и сбрасывает connected в БД. Если цель достигнута — цикл завершён; - * иначе переходим к Phase 2 (дополнительные попытки). - */ -static void tgc_phase1_timeout(void* arg) { - struct TOPO_GROUP_CONNECT* gc = (struct TOPO_GROUP_CONNECT*)arg; - gc->phase_timer = NULL; - if (gc->phase != TGC_PHASE_ONE) return; - - sqlite3* db = gc->group->instance->topo_sqlite_db; - uint64_t* ids = NULL; int count = 0; - topo_node_sqlite_get_connected_peers(db, gc->group->channel_id, &ids, &count); - for (int i = 0; i < count; i++) { - int connected = 0; - for (int j = 0; j < gc->handle_count; j++) { - struct NODE_CONN_DIRECT* h = gc->handles[j]; - if (h && node_conn_direct_get_conn(h) != NULL) { connected = 1; break; } - } - if (!connected) topo_node_sqlite_set_connected(db, gc->group->channel_id, ids[i], 0); - } - u_free(ids); - - int remaining = 0; - for (int i = 0; i < gc->handle_count; i++) { - struct NODE_CONN_DIRECT* h = gc->handles[i]; - if (h && node_conn_direct_get_conn(h) != NULL) { - gc->handles[remaining++] = h; - } else { - if (h) node_conn_direct_close(h); - } - } - gc->handle_count = remaining; - - if (tgc_goal_reached(gc)) { - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase1 done — goal reached (%d connected), skipping Phase 2", TGC_ID, gc->connected_count); - gc->phase = TGC_PHASE_DONE; - tgc_notify_connecting(gc, NULL, 0); - return; + if (earliest != UINT64_MAX) { + uint64_t delay = earliest > now ? earliest - now : 1; + gc->timer = uasync_set_timeout(gc->group->instance->ua, delay > INT_MAX ? INT_MAX : (int)delay, gc, tgc_timeout, "group_connect"); + if (!gc->timer) DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group connect timer allocation failed"); } - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase1 done — goal not reached (%d connected) → Phase 2", TGC_ID, gc->connected_count); - gc->phase = TGC_PHASE_TWO; - tgc_phase2_try_next(gc); + if (changed) tgc_notify(gc); } -/* ═══════════════════════════════════════════════════════════════════════ - * Phase 2 - * ══════════════════════════════════════════════════════════════════════ */ - -/** - * Последовательный перебор узлов с публичными/EIM адресами (сначала supernode, - * затем обычные публичные). Пробует следующий узел и ставит ему жёлтый кружок. - * Когда все перебраны — переходит к Phase 3. - */ -static void tgc_phase2_try_next(struct TOPO_GROUP_CONNECT* gc) { - gc->pause_timer = NULL; gc->pause_wait = NULL; - if (gc->phase == TGC_PHASE_PAUSE) gc->phase = TGC_PHASE_TWO; - if (gc->phase != TGC_PHASE_TWO) return; - while (1) { - if (gc->candidate_count == 0) { - sqlite3* db = gc->group->instance->topo_sqlite_db; - if (gc->tried_super == 0) { - topo_node_sqlite_get_supernode_peers(db, gc->group->channel_id, &gc->candidate_ids, &gc->candidate_count); - gc->tried_super = 1; gc->cursor = 0; - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 supernode round — loaded %d for ch=%s", - TGC_ID, gc->candidate_count, gc->group->channel_id); - } else if (gc->tried_super == 1) { - topo_node_sqlite_get_public_peers(db, gc->group->channel_id, &gc->candidate_ids, &gc->candidate_count, - gc->group->instance->node_id); - gc->tried_super = 2; gc->cursor = 0; - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 public round — loaded %d for ch=%s", - TGC_ID, gc->candidate_count, gc->group->channel_id); - } else { - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 exhausted — no connections, ch=%s", TGC_ID, gc->group->channel_id); - gc->phase = TGC_PHASE_THREE; gc->candidate_count = 0; - tgc_phase3_try_next(gc); - return; - } - } - if (gc->cursor < gc->candidate_count) { - uint64_t nid = gc->candidate_ids[gc->cursor++]; - if (nid == gc->group->instance->node_id) continue; - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase2 trying 0x%016llx (%d/%d) super_round=%d", TGC_ID, - (unsigned long long)nid, gc->cursor, gc->candidate_count, gc->tried_super); - if (node_conn_direct_open(gc->group->instance, nid, tgc_callback, gc, &gc->cur_handle, NULL) < 0) continue; - tgc_notify_connecting(gc, &nid, 1); - return; - } - if (gc->candidate_ids) { u_free(gc->candidate_ids); gc->candidate_ids = NULL; } - gc->candidate_count = 0; - } +static void tgc_wake(void* arg) { + struct TOPO_GROUP_CONNECT* gc = arg; gc->wake = NULL; tgc_step(gc); } -/* ═══════════════════════════════════════════════════════════════════════ - * Phase 3 — локальные/strict NAT адреса - * ══════════════════════════════════════════════════════════════════════ */ - -/** - * Последовательный перебор узлов с локальными/strict NAT адресами (как Phase 2, - * но для узлов, доступных только в локальной сети). Когда все перебраны — - * пауза и полный перезапуск цикла с Phase 1 (бесконечно, пока не подключимся). - */ -static void tgc_phase3_try_next(struct TOPO_GROUP_CONNECT* gc) { - gc->pause_timer = NULL; gc->pause_wait = NULL; - if (gc->phase == TGC_PHASE_PAUSE) gc->phase = TGC_PHASE_THREE; - if (gc->phase != TGC_PHASE_THREE) return; - if (gc->candidate_count == 0) { - sqlite3* db = gc->group->instance->topo_sqlite_db; - topo_node_sqlite_get_local_peers(db, gc->group->channel_id, &gc->candidate_ids, &gc->candidate_count, - gc->group->instance->node_id); - gc->cursor = 0; - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase3 loaded %d local peers for ch=%s", - TGC_ID, gc->candidate_count, gc->group->channel_id); - } - while (gc->cursor < gc->candidate_count) { - uint64_t nid = gc->candidate_ids[gc->cursor++]; - if (nid == gc->group->instance->node_id) continue; - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase3 trying 0x%016llx (%d/%d)", TGC_ID, - (unsigned long long)nid, gc->cursor, gc->candidate_count); - if (node_conn_direct_open(gc->group->instance, nid, tgc_callback, gc, &gc->cur_handle, NULL) < 0) continue; - tgc_notify_connecting(gc, &nid, 1); - return; - } - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: Phase3 exhausted — pause then cycle restart, ch=%s", TGC_ID, gc->group->channel_id); - tgc_start_pause(gc, (void(*)(void*))tgc_cycle_restart, "tgc_p3_restart"); +void topo_group_connect_changed(struct TOPO_GROUP* group) { + struct TOPO_GROUP_CONNECT* gc = group ? group->connect : NULL; + if (!gc || group->stopping || gc->wake) return; + gc->wake = uasync_call_soon(group->instance->ua, gc, tgc_wake); + if (!gc->wake) DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group connect wake allocation failed"); } -/* ═══════════════════════════════════════════════════════════════════════ - * Полный перезапуск цикла с Phase 1 (без destroy/init) - * ══════════════════════════════════════════════════════════════════════ */ - -/** - * Полный перезапуск цикла с Phase 1 без уничтожения модуля. Закрывает все - * открытые соединения, сбрасывает состояние и заново запускает параллельные - * попытки ко всем connected-пирам из БД. Вызывается при обрыве всех соединений - * и после исчерпания Phase 3. - */ -static void tgc_cycle_restart(struct TOPO_GROUP_CONNECT* gc) { - gc->pause_timer = NULL; gc->pause_wait = NULL; - if (gc->phase_timer) { uasync_cancel_timeout(gc->group->instance->ua, gc->phase_timer); gc->phase_timer = NULL; } - for (int i = 0; i < gc->handle_count; i++) node_conn_direct_close(gc->handles[i]); - if (gc->cur_handle) { node_conn_direct_close(gc->cur_handle); gc->cur_handle = NULL; } - gc->handle_count = 0; gc->pending = 0; gc->connected_count = 0; - gc->cursor = 0; gc->tried_super = 0; - u_free(gc->candidate_ids); gc->candidate_ids = NULL; gc->candidate_count = 0; - - uint64_t* ids = NULL; int count = 0; - sqlite3* db = gc->group->instance->topo_sqlite_db; - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: cycle_restart ch=%s active=%d", - TGC_ID, gc->group->channel_id, gc->active_conn_count); - if (!db || topo_node_sqlite_get_connected_peers(db, gc->group->channel_id, &ids, &count) != 0 || count == 0) { - if (ids) { u_free(ids); ids = NULL; } - gc->phase = TGC_PHASE_TWO; - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: cycle_restart ch=%s no connected peers → Phase 2", TGC_ID, gc->group->channel_id); - tgc_phase2_try_next(gc); - return; - } +void topo_group_connect_destroy(struct TOPO_GROUP* group) { + struct TOPO_GROUP_CONNECT* gc = group ? group->connect : NULL; + if (!gc) return; + group->connect = NULL; + if (gc->wake) uasync_call_soon_cancel(group->instance->ua, gc->wake); + if (gc->timer) uasync_cancel_timeout(group->instance->ua, gc->timer); + tgc_cancel_pause(gc); + for (size_t i = 0; i < gc->count; i++) topo_group_peer_close(gc->candidates[i].request); + u_free(gc->candidates); u_free(gc); +} - gc->phase = TGC_PHASE_ONE; - int launched = 0; - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: cycle_restart Phase 1 launching %d connects ch=%s", - TGC_ID, count, gc->group->channel_id); - for (int i = 0; i < count && gc->handle_count < TGC_MAX_HANDLES; i++) { - if (ids[i] == gc->group->instance->node_id) { DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: Phase1 skip self 0x%016llx", TGC_ID, (unsigned long long)ids[i]); continue; } - DEBUG_DEBUG(DEBUG_CATEGORY_BGP, "%s: Phase1 connect to 0x%016llx (%d/%d)", TGC_ID, (unsigned long long)ids[i], i + 1, count); - node_conn_direct_open(gc->group->instance, ids[i], tgc_callback, gc, - &gc->handles[gc->handle_count++], NULL); - launched++; +int topo_group_connect_init(struct TOPO_GROUP* group) { + if (!group || group->stopping || group->group_type != TOPO_GROUP_TYPE_CHAT || !group->channel_id[0]) { + DEBUG_WARN(DEBUG_CATEGORY_BGP, "cannot start group connect without a CHAT group"); return -1; } - gc->pending = launched; - tgc_notify_connecting(gc, ids, gc->handle_count); - u_free(ids); - - gc->phase_timer = uasync_set_timeout(gc->group->instance->ua, - TGC_DIRECT_TIMEOUT_MS * 10, gc, tgc_phase1_timeout, "tgc_phase1"); + topo_group_connect_destroy(group); + struct TOPO_GROUP_CONNECT* gc = u_calloc(1, sizeof(*gc)); + if (!gc) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "group connect allocation failed"); return -1; } + gc->group = group; gc->automatic = 1; group->connect = gc; + tgc_load(gc); topo_group_connect_changed(group); return 0; } -/* ═══════════════════════════════════════════════════════════════════════ - * Однократное подключение к конкретному узлу (ручной "Connect" из GUI) - * ══════════════════════════════════════════════════════════════════════ */ - -struct tgc_once_ctx { - struct TOPO_GROUP* group; - uint64_t node_id; - struct NODE_CONN_DIRECT* h; -}; +void topo_group_connect_restart(struct TOPO_GROUP* group) { + if (!group) return; + if (topo_group_connect_active_count(group)) return; + topo_group_connect_init(group); +} -static void tgc_once_cb(struct NODE_CONN_DIRECT* h, enum ncd_event ev, void* arg) { - struct tgc_once_ctx* c = (struct tgc_once_ctx*)arg; - struct ETCP_CONN* conn = node_conn_direct_get_conn(h); +void topo_group_connect_on_up(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { + if (!group || group->group_type != TOPO_GROUP_TYPE_CHAT || !conn) return; + if (group->instance->topo_sqlite_db && topo_group_peer_ready(group, conn->peer_node_id)) + topo_node_sqlite_set_connected(group->instance->topo_sqlite_db, group->channel_id, conn->peer_node_id, 1); + topo_group_connect_changed(group); +} - if (ev == NCD_EVENT_UP && conn) { - node_conn_direct_set_callback(h, NULL, NULL); - topo_group_new_conn(c->group, conn); - node_conn_direct_close(h); /* владение перешло к topo_group_new_conn */ - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: once-connect UP node=0x%016llx ch=%s", - TGC_ID, (unsigned long long)c->node_id, c->group->channel_id); - } else if (ev == NCD_EVENT_TIMEOUT) { - DEBUG_WARN(DEBUG_CATEGORY_BGP, "%s: once-connect TIMEOUT node=0x%016llx ch=%s", - TGC_ID, (unsigned long long)c->node_id, c->group->channel_id); - node_conn_direct_close(h); - } else if (ev == NCD_EVENT_DOWN || ev == NCD_EVENT_CLOSED) { - DEBUG_WARN(DEBUG_CATEGORY_BGP, "%s: once-connect DOWN node=0x%016llx ch=%s", - TGC_ID, (unsigned long long)c->node_id, c->group->channel_id); - node_conn_direct_close(h); - } - u_free(c); +void topo_group_connect_on_down(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { + if (!group || group->group_type != TOPO_GROUP_TYPE_CHAT || !conn) return; + if (group->instance->topo_sqlite_db) + topo_node_sqlite_set_connected(group->instance->topo_sqlite_db, group->channel_id, conn->peer_node_id, 0); + topo_group_connect_changed(group); } int topo_group_connect_node_once(struct TOPO_GROUP* group, uint64_t node_id) { - if (!group || group->group_type != TOPO_GROUP_TYPE_CHAT || !node_id) return -1; - if (node_id == group->instance->node_id) return -1; - - struct ETCP_CONN* ex = instance_find_conn(group->instance, node_id); - if (ex && ex->links_up && ex->initialized) { - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: once-connect node=0x%016llx already connected — add to group ch=%s", - TGC_ID, (unsigned long long)node_id, group->channel_id); - topo_group_new_conn(group, ex); - return 0; + if (!group || group->stopping || group->group_type != TOPO_GROUP_TYPE_CHAT || !node_id || node_id == group->instance->node_id) { + DEBUG_WARN(DEBUG_CATEGORY_BGP, "invalid manual group connect"); return -1; } - - struct tgc_once_ctx* c = u_calloc(1, sizeof(*c)); + if (!group->connect) { + group->connect = u_calloc(1, sizeof(*group->connect)); + if (!group->connect) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "manual group connect allocation failed"); return -1; } + group->connect->group = group; + } + struct TOPO_GROUP_CONNECT* gc = group->connect; + struct tgc_candidate* c = tgc_add(gc, node_id, 0, 1); if (!c) return -1; - c->group = group; c->node_id = node_id; - - int r = node_conn_direct_open(group->instance, node_id, tgc_once_cb, c, &c->h, NULL); - if (r == NCD_ERR || !c->h) { u_free(c); return -1; } - DEBUG_INFO(DEBUG_CATEGORY_BGP, "%s: once-connect node=0x%016llx ch=%s started (rc=%d)", - TGC_ID, (unsigned long long)node_id, group->channel_id, r); - return 0; + c->manual = 1; + if (c->request) return 0; + int result = tgc_open(gc, c); tgc_notify(gc); return result; } diff --git a/src/routing_layer/topo_group_connect.h b/src/routing_layer/topo_group_connect.h index e638e9bb..273bcfc4 100644 --- a/src/routing_layer/topo_group_connect.h +++ b/src/routing_layer/topo_group_connect.h @@ -1,26 +1,11 @@ -/** - * @file topo_group_connect.h - * @brief Авто-подключение к узлам группы — бесконечный цикл до цели. - * - * При старте CHAT-группы (topo_group_connect_init): - * Phase 1 — параллельный запуск ко всем пирам с connected=1 (таймаут 2с). - * Если цель достигнута → done. - * Phase 2 — последовательный перебор (supernode, затем public/EIM). - * После каждого TIMEOUT — пауза перед следующей попыткой. - * Phase 3 — последовательный перебор локальных/strict NAT. - * После каждого TIMEOUT — пауза перед следующей попыткой. - * После исчерпания Phase 3 → пауза → cycle_restart → Phase 1 (бесконечно). - * - * Цель (tgc_goal_reached): подключены хотя бы к одному суперузлу (node_type==4) - * ЛИБО к 3+ не-мобильным клиентам (device_type != MOBILE, не суперузлам). - * Пока цель не достигнута, успех НЕ останавливает цикл. - * - * Пауза: активный режим — 1с; фоновый (Android standby) — standby_wait - * (burst → сразу, sleep → до следующего burst). - * При DOWN, если цель перестала выполняться — немедленный cycle_restart. - * При повторном вызове topo_group_connect_init — destroy + fresh init. - * - * active_conn_count — количество уникальных peer-узлов с живыми соединениями. +/* CHAT connection policy over cancellable group requests. + * Phase 1: historical connected peers in parallel (up to 128). + * Phase 2: sequential supernode/public peers. Phase 3: local peers. + * Each node is tried once per cycle. CONNECTING: 2s; SYNCING: 5s without progress. + * Goal: one READY supernode or three READY nonmobile peers. After exhaustion, + * retry after 1s (or standby burst). Existing READY group sessions survive + * planner restart/stop. Manual attempts are tracked and cancelled with the group. + * active_count is derived from the current group sessions, not a separate counter. */ #ifndef TOPO_GROUP_CONNECT_H #define TOPO_GROUP_CONNECT_H @@ -32,14 +17,13 @@ struct ETCP_CONN; int topo_group_connect_init(struct TOPO_GROUP* group); void topo_group_connect_destroy(struct TOPO_GROUP* group); +void topo_group_connect_changed(struct TOPO_GROUP* group); void topo_group_connect_on_up(struct TOPO_GROUP* group, struct ETCP_CONN* conn); void topo_group_connect_on_down(struct TOPO_GROUP* group, struct ETCP_CONN* conn); int topo_group_connect_active_count(struct TOPO_GROUP* group); void topo_group_connect_restart(struct TOPO_GROUP* group); -/* Однократная попытка подключения к конкретному узлу группы (без бесконечного - * цикла). При успехе соединение добавляется в BGP-группу (topo_group_new_conn), - * handle сохраняется в nq->handle. Используется для ручного "Connect" из GUI. */ +/* One manual request; shares the group session, no NCD ownership handoff on UP. */ int topo_group_connect_node_once(struct TOPO_GROUP* group, uint64_t node_id); #endif diff --git a/src/routing_layer/topo_group_invite.c b/src/routing_layer/topo_group_invite.c index d144d570..b0931ba7 100644 --- a/src/routing_layer/topo_group_invite.c +++ b/src/routing_layer/topo_group_invite.c @@ -1,3 +1,4 @@ +#include "topo_group_connect.h" /** * @file topo_group_invite.c * @brief Invite/join к каналу через прямое (ncd) подключение + проверка членства. @@ -187,7 +188,7 @@ static void tgi_handle_info_req(struct ETCP_CONN* conn, const uint8_t* data, siz { struct TOPO_GROUP* bgp_grp = topo_groups_find(inst->topo_groups, req->group_id); if (bgp_grp && bgp_grp->group_type == TOPO_GROUP_TYPE_CHAT) - topo_group_new_conn(bgp_grp, conn); + topo_group_connect_node_once(bgp_grp, conn->peer_node_id); } /* ── сериализуем ответ ── */ @@ -342,11 +343,10 @@ static void tgi_handle_info_resp(struct ETCP_CONN* conn, const uint8_t* data, si { struct TOPO_GROUP* bgp_grp = topo_groups_find(inv->inst->topo_groups, inv->group_id); if (bgp_grp && bgp_grp->group_type == TOPO_GROUP_TYPE_CHAT) - topo_group_new_conn(bgp_grp, conn); + topo_group_connect_node_once(bgp_grp, conn->peer_node_id); } - /* ── успех: владение conn перешло к topo_group_new_conn (item->handle), - закрываем connect-фазу handle и доставляем JOIN ── */ + /* Bootstrap приглашения завершён; запрос групповой сессии уже держит свой интерес. */ { uint64_t node_id = inv->node_id, group_id = inv->group_id; tgi_cb_t cb = inv->cb; void* cb_arg = inv->cb_arg; @@ -509,23 +509,6 @@ int topo_group_invite_join(struct UTUN_INSTANCE* inst, uint64_t group_id, /* ═══════════ инвайтера сторона (direct, без INVITE_INFO) ═══════════ */ -struct tgi_conn_ctx { struct UTUN_INSTANCE* inst; uint64_t node_id, group_id; }; - -static void tgi_to_channel_cb(struct NODE_CONN_DIRECT* h, enum ncd_event ev, void* arg) { - struct tgi_conn_ctx* c = (struct tgi_conn_ctx*)arg; - if (ev == NCD_EVENT_UP) { - struct ETCP_CONN* conn = node_conn_direct_get_conn(h); - struct TOPO_GROUP* g = topo_groups_find(c->inst->topo_groups, c->group_id); - /* владение переходит к topo_group_new_conn (item->handle), connect-фазу закрываем */ - if (g && conn) topo_group_new_conn(g, conn); - node_conn_direct_close(h); - u_free(c); - } else if (ev == NCD_EVENT_TIMEOUT || ev == NCD_EVENT_CLOSED) { - node_conn_direct_force_close(h); - u_free(c); - } -} - int topo_group_invite_to_channel(struct UTUN_INSTANCE* inst, uint64_t group_id, struct TOPO_NODE* ni, uint64_t node_id) { if (!inst || !inst->topo_groups) return -1; @@ -537,23 +520,11 @@ int topo_group_invite_to_channel(struct UTUN_INSTANCE* inst, uint64_t group_id, } if (!node_id || node_id == inst->node_id) return -1; - struct TOPO_GROUP_NODE* nq = topo_node_find_by_id(group, node_id); - if (!nq) { - struct ll_entry* qe = queue_entry_new(sizeof(struct TOPO_GROUP_NODE)); - if (!qe) return -1; - nq = (struct TOPO_GROUP_NODE*)qe; - memset((uint8_t*)nq + sizeof(struct ll_entry), 0, sizeof(*nq) - sizeof(struct ll_entry)); - nq->node_id = node_id; - queue_data_put_with_index(group->nodes, &nq->ll); + struct NODE_CONN_DIRECT* seed = NULL; + if (ni && node_conn_direct_open_node(inst, node_id, NULL, NULL, &seed, ni, NULL) == NCD_ERR) { + DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot seed invited peer addresses node=%016llx", (unsigned long long)node_id); return -1; } - - struct tgi_conn_ctx* c = u_calloc(1, sizeof(*c)); - if (!c) return -1; - c->inst = inst; c->node_id = node_id; c->group_id = group_id; - - struct NODE_CONN_DIRECT* h = NULL; - int r = ni ? node_conn_direct_open_node(inst, node_id, tgi_to_channel_cb, c, &h, ni, NULL) - : node_conn_direct_open(inst, node_id, tgi_to_channel_cb, c, &h, NULL); - if (r == NCD_ERR || !h) { u_free(c); return -1; } - return 0; + int result = topo_group_connect_node_once(group, node_id); + node_conn_direct_close(seed); + return result; } diff --git a/src/transport_layer/etcp_connections.c b/src/transport_layer/etcp_connections.c index ba2b1599..9b44c16a 100644 --- a/src/transport_layer/etcp_connections.c +++ b/src/transport_layer/etcp_connections.c @@ -2770,20 +2770,8 @@ int init_sockets(struct UTUN_INSTANCE* instance) { return 0; // All OK } -/* Конфигурация владеет только NCD handles; транспортные линки создаёт NCD. */ -static void config_client_event(struct NODE_CONN_DIRECT* handle, enum ncd_event event, void* arg) { - struct UTUN_INSTANCE* instance = arg; - if (event == NCD_EVENT_UP) { - struct TOPO_GROUP* group = topo_groups_get_default(instance->topo_groups); - if (group && topo_group_join_peer(group, handle) < 0) - DEBUG_ERROR(DEBUG_CATEGORY_BGP, "configured peer join failed node=%016llx", - (unsigned long long)node_conn_direct_node_id(handle)); - } else if (event == NCD_EVENT_TIMEOUT || event == NCD_EVENT_CLOSED) { - DEBUG_WARN(DEBUG_CATEGORY_CONNECTION, "configured peer node=%016llx event=%d", - (unsigned long long)node_conn_direct_node_id(handle), event); - } -} - +/* Адреса конфигурации добавляются в NCD при подготовке. Постоянным владельцем + * транспорта сразу становится группа; адаптер держит только запрос участия. */ static struct NODE_CONN_DIRECT* config_client_open(struct UTUN_INSTANCE* instance, struct CFG_CLIENT* client) { struct TOPO_NODE ni = { .node_name = client->name }; if (sc_hex_to_binary(client->peer_public_key_hex, ni.public_key, SC_PUBKEY_SIZE) != SC_OK) { @@ -2839,41 +2827,42 @@ fail: static void config_handles_close(struct CONFIG_CONN_HANDLE* head) { while (head) { struct CONFIG_CONN_HANDLE* next = head->next; - node_conn_direct_close(head->handle); u_free(head); head = next; + topo_group_peer_close(head->request); u_free(head); head = next; } } int init_client_connections(struct UTUN_INSTANCE* instance, struct CFG_CLIENT* clients) { if (!instance) { DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "configure clients without instance"); return -1; } + struct TOPO_GROUP* group = topo_groups_get_default(instance->topo_groups); + if (clients && !group) { DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "UTUN group is not running"); return -1; } struct CONFIG_CONN_HANDLE* updated = NULL; int count = 0; for (struct CFG_CLIENT* client = clients; client; client = client->next) { struct CONFIG_CONN_HANDLE* owner = u_calloc(1, sizeof(*owner)); if (!owner) { DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "client handle allocation failed"); goto fail; } - owner->handle = config_client_open(instance, client); - if (!owner->handle) { u_free(owner); goto fail; } - owner->node_id = node_conn_direct_node_id(owner->handle); + struct NODE_CONN_DIRECT* seed = config_client_open(instance, client); + if (!seed) { u_free(owner); goto fail; } + owner->node_id = node_conn_direct_node_id(seed); + int opened = topo_group_peer_open(group, owner->node_id, &owner->request); + node_conn_direct_close(seed); + if (opened < 0) { u_free(owner); goto fail; } snprintf(owner->name, sizeof(owner->name), "%s", client->name); owner->next = updated; updated = owner; count++; } struct CONFIG_CONN_HANDLE* old = instance->config_conn_handles; instance->config_conn_handles = updated; - struct TOPO_GROUP* group = topo_groups_get_default(instance->topo_groups); for (struct CONFIG_CONN_HANDLE* owner = old; owner; owner = owner->next) { int retained = 0; for (struct CONFIG_CONN_HANDLE* next = updated; next; next = next->next) if (next->node_id == owner->node_id) { retained = 1; break; } - if (!retained && group && node_conn_direct_get_conn(owner->handle)) topo_group_leave_peer(group, owner->handle); - } - for (struct CONFIG_CONN_HANDLE* owner = updated; owner; owner = owner->next) { - node_conn_direct_set_callback(owner->handle, config_client_event, instance); + if (!retained && group) topo_group_peer_leave(group, owner->node_id); } config_handles_close(old); - DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "configured %d NCD client handles", count); + DEBUG_INFO(DEBUG_CATEGORY_CONNECTION, "configured %d group client requests", count); return 0; fail: config_handles_close(updated); - DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "NCD client configuration rejected; previous handles retained"); + DEBUG_ERROR(DEBUG_CATEGORY_CONFIG, "group client configuration rejected; previous requests retained"); return -1; } diff --git a/src/utun_instance.c b/src/utun_instance.c index b8a2baf8..12cd7f7b 100644 --- a/src/utun_instance.c +++ b/src/utun_instance.c @@ -446,7 +446,7 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) { int ch_count = 0; while (ch) { struct CONFIG_CONN_HANDLE* next = ch->next; - node_conn_direct_close(ch->handle); + topo_group_peer_close(ch->request); u_free(ch); ch_count++; ch = next; } diff --git a/src/utun_instance.h b/src/utun_instance.h index 5338214c..5b1e4564 100644 --- a/src/utun_instance.h +++ b/src/utun_instance.h @@ -79,7 +79,7 @@ struct conn_queue_entry { struct CONFIG_CONN_HANDLE { uint64_t node_id; char name[MAX_CONN_NAME_LEN]; - struct NODE_CONN_DIRECT* handle; + struct TOPO_PEER_REQUEST* request; struct CONFIG_CONN_HANDLE* next; }; diff --git a/tests/bbr_integration/test_bbr_integration.c b/tests/bbr_integration/test_bbr_integration.c index d1ce82bb..40dc02e7 100644 --- a/tests/bbr_integration/test_bbr_integration.c +++ b/tests/bbr_integration/test_bbr_integration.c @@ -108,19 +108,8 @@ static uint64_t now_us(void) { /* ===== Instance creation ===== */ static struct UTUN_INSTANCE* create_instance(struct UASYNC* u, uint64_t node_id, const char* priv_hex, const char* pub_hex) { - struct UTUN_INSTANCE* inst = u_calloc(1, sizeof(*inst)); - if (!inst) return NULL; - inst->ua = u; - inst->node_id = node_id; - if (sc_init_local_keys(&inst->my_keys, pub_hex, priv_hex) != SC_OK) { u_free(inst); return NULL; } - inst->ack_pool = memory_pool_init(sizeof(struct ACK_PACKET), "ack_pool"); - inst->data_pool = memory_pool_init(PACKET_DATA_SIZE, "data_pool"); - inst->pkt_pool = memory_pool_init(sizeof(struct ETCP_DGRAM) + PACKET_DATA_SIZE, "pkt_pool"); - if (!inst->ack_pool || !inst->data_pool || !inst->pkt_pool) { u_free(inst); return NULL; } - inst->connections = queue_new(u, 256, 0, 8, "connections"); - if (!inst->connections) { u_free(inst); return NULL; } struct utun_config* cfg = u_calloc(1, sizeof(*cfg)); - if (!cfg) { u_free(inst); return NULL; } + if (!cfg) return NULL; strncpy(cfg->global.my_public_key_hex, pub_hex, MAX_KEY_LEN - 1); strncpy(cfg->global.my_private_key_hex, priv_hex, MAX_KEY_LEN - 1); cfg->global.my_node_id = node_id; @@ -129,8 +118,7 @@ static struct UTUN_INSTANCE* create_instance(struct UASYNC* u, uint64_t node_id, cfg->global.keepalive_interval = 500; cfg->global.allowed_keys_allow_all = 1; cfg->global.bbr_max_cwnd = 100000; - inst->config = cfg; - return inst; + return utun_instance_create_from_config(u, cfg); } static int add_server(struct UTUN_INSTANCE* inst, const char* name, int port) { diff --git a/tests/test_group_recovery.c b/tests/test_group_recovery.c index e427f7da..4c1ff85f 100644 --- a/tests/test_group_recovery.c +++ b/tests/test_group_recovery.c @@ -3,6 +3,8 @@ #include "utun_instance.h" #include "topo_group.h" #include "topo_recovery.h" +#include "topo_group_connect.h" +#include "topo_node_sqlite.h" #include "etcp.h" #include "etcp_api.h" #include "node_conn_direct.h" @@ -265,9 +267,45 @@ static void sender_backpressure(void) { poll_events(&f); destroy(&f); } +static void persistent_request(void) { + struct fixture f; create(&f); + struct TOPO_PEER_REQUEST* request = NULL; + assert(topo_group_peer_open(f.group, f.ids[0], &request) == 0); + struct NODE_CONN_DIRECT* ownership = peer(&f, 0)->handle; + up(&f, 0); ready(&f, 0); + f.conns[0]->links_up = 0; + etcp_cbk_fire(f.conns[0], ETCP_CBK_EVENT_DOWN); poll_events(&f); + assert(topo_group_peer_phase(request) == TOPO_PEER_FAILED); + assert(peer(&f, 0)->handle == ownership && topo_group_peer_conn(request) == f.conns[0]); + up(&f, 0); ready(&f, 0); + assert(topo_group_peer_phase(request) == TOPO_PEER_READY && peer(&f, 0)->handle == ownership); + topo_group_peer_close(request); destroy(&f); +} + +static void manual_planner_cancel(void) { + struct fixture f; create(&f); + assert(sqlite3_open(":memory:", &f.inst->topo_sqlite_db) == SQLITE_OK); + assert(topo_node_sqlite_init(f.inst->topo_sqlite_db) == 0); + f.group->group_type = TOPO_GROUP_TYPE_CHAT; strcpy(f.group->channel_id, "42"); + /* The membership table is the only authorization source. */ + uint8_t key[32] = {1}, signature[64] = {1}; + assert(topo_node_sqlite_channel_put(f.inst->topo_sqlite_db, "42", "test", 0, key, NULL, key, NULL, signature) == 0); + assert(topo_node_sqlite_member_block_put(f.inst->topo_sqlite_db, "42", f.ids[0], signature, 1, signature, 1, + key, key, "{}", NULL, 0, signature) == 0); + assert(topo_group_connect_node_once(f.group, f.ids[0]) == 0); + struct NODE_CONN_DIRECT* ownership = peer(&f, 0)->handle; + assert(topo_group_connect_node_once(f.group, f.ids[0]) == 0 && peer(&f, 0)->handle == ownership); + up(&f, 0); poll_events(&f); + assert(topo_group_connect_active_count(f.group) == 0); /* UP does not finish a group attempt. */ + topo_groups_remove_group(f.inst->topo_groups, f.group->group_id); f.group = NULL; + poll_events(&f); destroy(&f); +} + int main(void) { debug_config_init(); debug_set_level(DEBUG_LEVEL_INFO); utun_instance_set_tun_init_enabled(0); + persistent_request(); + manual_planner_cancel(); sender_backpressure(); partial_and_external_routes(); stalled_and_exhausted(); diff --git a/tests/test_ncd_config.c b/tests/test_ncd_config.c index d492e2fb..c8ea2572 100644 --- a/tests/test_ncd_config.c +++ b/tests/test_ncd_config.c @@ -3,6 +3,7 @@ #include #include "utun_instance.h" #include "config_parser.h" +#include "routing_layer/topo_group.h" #include "transport_layer/etcp.h" #include "transport_layer/etcp_connections.h" #include "transport_layer/node_conn_direct.h" @@ -46,12 +47,12 @@ int main(void) { first.next = &second; second.next = &third; client.links = &first; assert(init_client_connections(inst, &client) == 0); assert(inst->config_conn_handles && !inst->config_conn_handles->next); - struct ETCP_CONN* conn = node_conn_direct_get_conn(inst->config_conn_handles->handle); + struct ETCP_CONN* conn = topo_group_peer_conn(inst->config_conn_handles->request); assert(conn && links(conn) == 3 && !conn->fin_wait); struct NODE_CONN_DIRECT* chat = NULL; assert(node_conn_direct_open(inst, conn->peer_node_id, NULL, NULL, &chat, NULL) == NCD_REUSED); assert(init_client_connections(inst, &client) == 0); - assert(node_conn_direct_get_conn(inst->config_conn_handles->handle) == conn && links(conn) == 3 && !conn->fin_wait); + assert(topo_group_peer_conn(inst->config_conn_handles->request) == conn && links(conn) == 3 && !conn->fin_wait); /* Ошибка нового набора не освобождает предыдущие handles. */ struct CFG_CLIENT invalid = { .name = "invalid", .peer_public_key_hex = "not-a-key" }; @@ -85,7 +86,7 @@ int main(void) { } third.local_srv = NULL; assert(init_client_connections(inst, &client) == 0); - conn = node_conn_direct_get_conn(inst->config_conn_handles->handle); + conn = topo_group_peer_conn(inst->config_conn_handles->request); assert(conn && links(conn) == 1 && conn->links->is_tcp && !conn->links->conn); assert(conn->links->reality_set && !strcmp(conn->links->reality.server_name, client.reality.server_name)); assert(init_client_connections(inst, NULL) == 0);