From a7632cfee1a3714f4325c4054657de5519793e6d Mon Sep 17 00:00:00 2001 From: evgeny Date: Fri, 2 Oct 2026 11:43:10 +0300 Subject: [PATCH] Select chat auto peers by device class and transport UP --- src/routing_layer/topo_group_connect.c | 308 ++++++++++++++++--------- src/routing_layer/topo_group_connect.h | 19 +- 2 files changed, 216 insertions(+), 111 deletions(-) diff --git a/src/routing_layer/topo_group_connect.c b/src/routing_layer/topo_group_connect.c index d0266045..cb1d0c0b 100644 --- a/src/routing_layer/topo_group_connect.c +++ b/src/routing_layer/topo_group_connect.c @@ -1,6 +1,6 @@ /* Политика CHAT-подключений. Транспортом и JOIN владеет групповая сессия. - * Сначала параллельно восстанавливаем connected-пиров, затем последовательно - * пробуем supernode/public/local. Незавершённые попытки — отменяемые запросы. */ + * Суперузел соединяется со всеми суперузлами, клиент — с двумя. Без UP + * суперузлов резервируем desktop, затем mobile, продолжая проверки старших классов. */ #include #include #include "topo_group_connect.h" @@ -12,40 +12,59 @@ #include "../lib/mem.h" #include "../lib/platform_compat.h" #include "../chat/chat_event.h" +#include "../transport_layer/node_conn_direct.h" #include "etcp.h" #ifdef UTUN_HAVE_STANDBY #include "standby.h" #endif -#define TGC_PARALLEL 128 +#define TGC_PARALLEL 2 +#define TGC_TARGET 2 #define TGC_PAUSE_TB 10000 +#define TGC_RETRY_MAX_TB 300000 + +enum tgc_class { TGC_UNKNOWN, TGC_SUPER, TGC_DESKTOP, TGC_MOBILE, TGC_CLASSES }; +static const char* const tgc_names[] = { "unknown", "supernode", "desktop", "mobile" }; struct tgc_candidate { - uint64_t node_id, started; + uint64_t node_id, started, retry_at, retry_delay; struct TOPO_PEER_REQUEST* request; - uint8_t priority, tried, manual; + uint8_t kind, tested, manual, listed; }; struct TOPO_GROUP_CONNECT { struct TOPO_GROUP* group; struct tgc_candidate* candidates; size_t count; - void *wake, *timer, *pause_timer, *pause_wait; - uint8_t automatic, cycle_done, paused, reload_after_pause; + void *wake, *timer, *pause_wait; + uint8_t automatic, paused, dirty, self_super, logged; + int last_up[TGC_CLASSES], last_allowed[TGC_CLASSES]; }; static void tgc_step(void* arg); +/* UP относится к удерживаемому группой транспорту, независимо от JOIN/обмена таблицами. */ +static int tgc_peer_up(struct TOPO_GROUP* group, uint64_t id) { + 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 (peer->node_id != id || peer->transport_down) continue; + struct ETCP_CONN* conn = node_conn_direct_get_conn(peer->handle); + return conn && conn->links_up && !conn->close_requested; + } + return 0; +} + 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++; + for (size_t i = 0; i < gc->count; i++) + if (gc->candidates[i].request && !tgc_peer_up(gc->group, gc->candidates[i].node_id)) 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) { + for (size_t i = 0; i < gc->count; i++) if (gc->candidates[i].request && !tgc_peer_up(gc->group, gc->candidates[i].node_id)) { 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); @@ -55,155 +74,217 @@ 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++; + if (tgc_peer_up(group, peer->node_id)) count++; } return count; } -static int tgc_goal(struct TOPO_GROUP_CONNECT* gc) { - int super = 0, other = 0; - struct TOPO_GROUP* group = gc->group; - 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++; - } - return super >= 1 || other >= 3; -} - -static struct tgc_candidate* tgc_add(struct TOPO_GROUP_CONNECT* gc, uint64_t id, uint8_t priority, int manual) { +/* Добавить кандидата один раз; живые запросы сохраняются при обновлении списка. */ +static struct tgc_candidate* tgc_add(struct TOPO_GROUP_CONNECT* gc, uint64_t id) { 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 }; + list[gc->count] = (struct tgc_candidate){ .node_id = id, .retry_delay = TGC_PAUSE_TB }; return &list[gc->count++]; } -static void tgc_load(struct TOPO_GROUP_CONNECT* gc) { +/* Роль не зависит от адресов; неизвестное устройство не подменяем десктопом. */ +static int tgc_load(struct TOPO_GROUP_CONNECT* gc) { + struct topo_connect_peer* peers = NULL; size_t count = 0; + struct TOPO_GROUP* group = gc->group; + if (topo_node_sqlite_get_connect_peers(group->instance->topo_sqlite_db, group->channel_id, &peers, &count) < 0) return -1; + for (size_t i = 0; i < gc->count; i++) gc->candidates[i].listed = 0; + int self_super = 0; + for (size_t i = 0; i < count; i++) { + if (peers[i].node_id == group->instance->node_id) { self_super = peers[i].supernode; continue; } + struct tgc_candidate* c = tgc_add(gc, peers[i].node_id); + if (!c) { u_free(peers); return -1; } + uint8_t device = peers[i].client_type; + struct TOPO_NODE* node = topo_node_registry_find(group->instance->topo_groups, c->node_id); + if (node && node->timestamp) device = node->client_type; + uint8_t kind = peers[i].supernode ? TGC_SUPER : device == CLIENT_TYPE_DESKTOP ? TGC_DESKTOP : + device == CLIENT_TYPE_MOBILE ? TGC_MOBILE : TGC_UNKNOWN; + if (c->kind != kind) { + DEBUG_INFO(DEBUG_CATEGORY_BGP, "group candidate class: group=%016llx node=%016llx %s->%s device=%u", + (unsigned long long)group->group_id, (unsigned long long)c->node_id, + tgc_names[c->kind], tgc_names[kind], device); + c->kind = kind; c->tested = 0; c->retry_at = 0; c->retry_delay = TGC_PAUSE_TB; + } + c->listed = 1; + } + u_free(peers); 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); + if (c->listed || c->manual) gc->candidates[retained++] = *c; + else { + DEBUG_INFO(DEBUG_CATEGORY_BGP, "group candidate removed: group=%016llx node=%016llx", + (unsigned long long)group->group_id, (unsigned long long)c->node_id); + topo_group_peer_close(c->request); + } + } + gc->count = retained; + if (gc->self_super != self_super) gc->logged = 0; + gc->self_super = self_super; gc->dirty = 0; return 0; +} + +/* Проверка старшего класса не останавливается после набора резервных пиров. */ +static void tgc_policy(struct TOPO_GROUP_CONNECT* gc, int* up, int* pending, int* allowed) { + int tested[TGC_CLASSES] = {1, 1, 1, 1}; + memset(up, 0, TGC_CLASSES * sizeof(*up)); memset(pending, 0, TGC_CLASSES * sizeof(*pending)); + memset(allowed, 0, TGC_CLASSES * sizeof(*allowed)); + for (size_t i = 0; i < gc->count; i++) { + struct tgc_candidate* c = &gc->candidates[i]; + if (!c->listed) continue; + if (tgc_peer_up(gc->group, c->node_id)) { up[c->kind]++; c->tested = 1; } + else { + if (!c->tested) tested[c->kind] = 0; + if (c->request) pending[c->kind]++; + } } - 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); + allowed[TGC_SUPER] = gc->self_super || up[TGC_SUPER] < TGC_TARGET; + allowed[TGC_DESKTOP] = !gc->self_super && !up[TGC_SUPER] && tested[TGC_SUPER] && up[TGC_DESKTOP] < TGC_TARGET; + allowed[TGC_MOBILE] = !gc->self_super && !up[TGC_SUPER] && tested[TGC_SUPER] && + !up[TGC_DESKTOP] && tested[TGC_DESKTOP] && up[TGC_MOBILE] < TGC_TARGET; + if (!gc->logged || memcmp(up, gc->last_up, sizeof(gc->last_up)) || memcmp(allowed, gc->last_allowed, sizeof(gc->last_allowed))) { + DEBUG_INFO(DEBUG_CATEGORY_BGP, "group connect policy: group=%016llx self=%s UP=super:%d desktop:%d mobile:%d" + " probe=super:%d desktop:%d mobile:%d scanned=super:%d desktop:%d", + (unsigned long long)gc->group->group_id, gc->self_super ? "supernode" : + gc->group->instance->client_type == CLIENT_TYPE_MOBILE ? "mobile" : "desktop", + up[TGC_SUPER], up[TGC_DESKTOP], up[TGC_MOBILE], allowed[TGC_SUPER], allowed[TGC_DESKTOP], + allowed[TGC_MOBILE], tested[TGC_SUPER], tested[TGC_DESKTOP]); + memcpy(gc->last_up, up, sizeof(gc->last_up)); memcpy(gc->last_allowed, allowed, sizeof(gc->last_allowed)); gc->logged = 1; } - DEBUG_INFO(DEBUG_CATEGORY_BGP, "group connect cycle: group=%016llx candidates=%zu", - (unsigned long long)gc->group->group_id, gc->count); } 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); #endif - gc->pause_timer = NULL; gc->pause_wait = NULL; gc->paused = 0; + gc->pause_wait = NULL; gc->paused = 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); + gc->pause_wait = NULL; gc->paused = 0; topo_group_connect_changed(gc->group); } -static void tgc_pause(struct TOPO_GROUP_CONNECT* gc, int reload) { - gc->paused = 1; gc->reload_after_pause = reload; +static int tgc_pause(struct TOPO_GROUP_CONNECT* gc) { #ifdef UTUN_HAVE_STANDBY if (standby_is_enabled()) { gc->pause_wait = standby_wait(gc, tgc_resume); - if (gc->pause_wait) return; + if (gc->pause_wait) { gc->paused = 1; return 1; } DEBUG_WARN(DEBUG_CATEGORY_BGP, "cannot wait for standby; using group retry timer"); } #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; } + return 0; +} + +/* Неудачные попытки повторяем с ограниченным backoff; новые кандидаты идут раньше повторных. */ +static void tgc_retry(struct tgc_candidate* c, uint64_t now) { + c->tested = 1; c->retry_at = now + c->retry_delay; + c->retry_delay = c->retry_delay < TGC_RETRY_MAX_TB / 2 ? c->retry_delay * 2 : TGC_RETRY_MAX_TB; } static int tgc_open(struct TOPO_GROUP_CONNECT* gc, struct tgc_candidate* c) { - c->tried = 1; c->started = get_time_tb(); + 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); + DEBUG_INFO(DEBUG_CATEGORY_BGP, "group connect attempt: group=%016llx node=%016llx class=%s manual=%u result=%d", + (unsigned long long)gc->group->group_id, (unsigned long long)c->node_id, tgc_names[c->kind], c->manual, result); if (result == 0) topo_group_connect_changed(gc->group); + else { tgc_retry(c, c->started); c->manual = 0; } return result; } -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 uint64_t tgc_deadline(struct tgc_candidate* c) { + return c->started + TOPO_RECOVERY_CONNECT_TIMEOUT_MS * 10ULL; } static void tgc_timeout(void* arg) { - struct TOPO_GROUP_CONNECT* gc = arg; gc->timer = NULL; tgc_step(gc); + struct TOPO_GROUP_CONNECT* gc = arg; gc->timer = NULL; gc->dirty = 1; tgc_step(gc); } +/* UP закрепляем за группой до освобождения временного запроса; READY занимается сама группа. */ 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; + int changed = 0; + int loaded = !gc->automatic || !gc->dirty || tgc_load(gc) == 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); + if (tgc_peer_up(gc->group, c->node_id)) { + struct ETCP_CONN* conn = topo_group_peer_conn(c->request); + if (topo_group_new_conn(gc->group, conn) < 0) { + DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot retain UP group peer=%016llx group=%016llx", + (unsigned long long)c->node_id, (unsigned long long)gc->group->group_id); + tgc_retry(c, now); + } else { + c->tested = 1; c->retry_at = 0; c->retry_delay = TGC_PAUSE_TB; + DEBUG_INFO(DEBUG_CATEGORY_BGP, "group connect UP: group=%016llx node=%016llx class=%s ready=%d", + (unsigned long long)gc->group->group_id, (unsigned long long)c->node_id, + tgc_names[c->kind], topo_group_peer_ready(gc->group, 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++; + topo_group_peer_close(request); c->manual = 0; changed = 1; + } else if (phase == TOPO_PEER_FAILED || now >= tgc_deadline(c)) { + tgc_retry(c, now); + DEBUG_WARN(DEBUG_CATEGORY_BGP, "group connect failed: group=%016llx node=%016llx class=%s phase=%d" + " elapsed_ms=%llu retry_ms=%llu", (unsigned long long)gc->group->group_id, + (unsigned long long)c->node_id, tgc_names[c->kind], phase, + (unsigned long long)((now - c->started) / 10), (unsigned long long)((c->retry_at - now) / 10)); + struct TOPO_PEER_REQUEST* request = c->request; c->request = NULL; + topo_group_peer_close(request); c->manual = 0; changed = 1; + } } - 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; + int up[TGC_CLASSES], pending[TGC_CLASSES], allowed[TGC_CLASSES]; + tgc_policy(gc, up, pending, allowed); + if (gc->automatic && loaded) { + for (size_t i = 0; i < gc->count; i++) { + struct tgc_candidate* c = &gc->candidates[i]; + if (c->manual || !c->request || tgc_peer_up(gc->group, c->node_id)) continue; + int limit = gc->self_super && c->kind == TGC_SUPER ? TGC_PARALLEL : TGC_TARGET - up[c->kind]; + if (!allowed[c->kind] || pending[c->kind] > limit) { + DEBUG_INFO(DEBUG_CATEGORY_BGP, "group connect cancel: group=%016llx node=%016llx class=%s" + " reason=%s super_UP=%d desktop_UP=%d", (unsigned long long)gc->group->group_id, + (unsigned long long)c->node_id, tgc_names[c->kind], allowed[c->kind] ? "target" : "class_policy", + up[TGC_SUPER], up[TGC_DESKTOP]); + struct TOPO_PEER_REQUEST* request = c->request; c->request = NULL; + topo_group_peer_close(request); pending[c->kind]--; c->retry_at = now; changed = 1; } - gc->cycle_done = 1; } - } 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; } + for (int kind = TGC_SUPER; !gc->paused && kind < TGC_CLASSES; kind++) { + if (!allowed[kind]) continue; + int limit = gc->self_super && kind == TGC_SUPER ? TGC_PARALLEL : TGC_TARGET - up[kind]; + while (pending[kind] < limit) { + struct tgc_candidate* best = NULL; + for (size_t i = 0; i < gc->count; i++) { + struct tgc_candidate* c = &gc->candidates[i]; + if (!c->listed || c->kind != kind || c->request || c->retry_at > now || tgc_peer_up(gc->group, c->node_id)) continue; + if (!best || c->tested < best->tested || (c->tested == best->tested && c->retry_at < best->retry_at)) best = c; + } + if (!best) break; + if (tgc_open(gc, best) == 0) pending[kind]++; + changed = 1; } - if (!available && !active) tgc_pause(gc, 1); } } uint64_t earliest = UINT64_MAX; + int connecting = 0; 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)); + uint64_t deadline = tgc_deadline(&gc->candidates[i]); connecting++; if (deadline < earliest) earliest = deadline; } + int retry = gc->automatic && (!loaded || (!gc->self_super && up[TGC_SUPER] < TGC_TARGET)); + if (gc->automatic && gc->self_super) + for (size_t i = 0; i < gc->count; i++) + if (gc->candidates[i].kind == TGC_SUPER && !tgc_peer_up(gc->group, gc->candidates[i].node_id)) retry = 1; + if (retry && !gc->paused && (connecting || !tgc_pause(gc)) && now + TGC_PAUSE_TB < earliest) earliest = now + TGC_PAUSE_TB; 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"); @@ -218,7 +299,9 @@ static void tgc_wake(void* arg) { 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; + if (!gc || group->stopping) return; + gc->dirty = 1; + if (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"); } @@ -226,11 +309,19 @@ void topo_group_connect_changed(struct TOPO_GROUP* group) { void topo_group_connect_destroy(struct TOPO_GROUP* group) { struct TOPO_GROUP_CONNECT* gc = group ? group->connect : NULL; if (!gc) return; + DEBUG_INFO(DEBUG_CATEGORY_BGP, "group connect stop: group=%016llx automatic=%u candidates=%zu UP=%d", + (unsigned long long)group->group_id, gc->automatic, gc->count, topo_group_connect_active_count(group)); 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); + for (size_t i = 0; i < gc->count; i++) { + struct tgc_candidate* c = &gc->candidates[i]; + if (!group->stopping && c->request && tgc_peer_up(group, c->node_id) && + topo_group_new_conn(group, topo_group_peer_conn(c->request)) < 0) + DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot retain UP peer while stopping planner node=%016llx", (unsigned long long)c->node_id); + topo_group_peer_close(c->request); + } u_free(gc->candidates); u_free(gc); } @@ -238,30 +329,41 @@ 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; } - topo_group_connect_destroy(group); - struct TOPO_GROUP_CONNECT* gc = u_calloc(1, sizeof(*gc)); + struct TOPO_GROUP_CONNECT* gc = group->connect; + if (!gc) 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; + gc->group = group; gc->automatic = 1; gc->dirty = 1; gc->logged = 0; group->connect = gc; + DEBUG_INFO(DEBUG_CATEGORY_BGP, "group connect start: group=%016llx device=%u target=%d parallel_per_class=%d", + (unsigned long long)group->group_id, group->instance->client_type, TGC_TARGET, TGC_PARALLEL); + topo_group_connect_changed(group); return 0; } 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); + struct TOPO_GROUP_CONNECT* gc = group->connect; + if (!gc || !gc->automatic) return; + tgc_cancel_pause(gc); + for (size_t i = 0; i < gc->count; i++) { + gc->candidates[i].retry_at = 0; gc->candidates[i].retry_delay = TGC_PAUSE_TB; + } + DEBUG_INFO(DEBUG_CATEGORY_BGP, "group connect network retry: group=%016llx UP=%d", + (unsigned long long)group->group_id, topo_group_connect_active_count(group)); + topo_group_connect_changed(group); } 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); + if (group->instance->topo_sqlite_db && conn->links_up && !conn->close_requested && + topo_node_sqlite_set_connected(group->instance->topo_sqlite_db, group->channel_id, conn->peer_node_id, 1) < 0) + DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot save group peer UP node=%016llx", (unsigned long long)conn->peer_node_id); topo_group_connect_changed(group); } 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); + if (group->instance->topo_sqlite_db && + topo_node_sqlite_set_connected(group->instance->topo_sqlite_db, group->channel_id, conn->peer_node_id, 0) < 0) + DEBUG_ERROR(DEBUG_CATEGORY_BGP, "cannot save group peer DOWN node=%016llx", (unsigned long long)conn->peer_node_id); topo_group_connect_changed(group); } @@ -275,7 +377,7 @@ int topo_group_connect_node_once(struct TOPO_GROUP* group, uint64_t node_id) { group->connect->group = group; } struct TOPO_GROUP_CONNECT* gc = group->connect; - struct tgc_candidate* c = tgc_add(gc, node_id, 0, 1); + struct tgc_candidate* c = tgc_add(gc, node_id); if (!c) return -1; c->manual = 1; if (c->request) return 0; diff --git a/src/routing_layer/topo_group_connect.h b/src/routing_layer/topo_group_connect.h index 461dd985..c3ea44fd 100644 --- a/src/routing_layer/topo_group_connect.h +++ b/src/routing_layer/topo_group_connect.h @@ -1,8 +1,11 @@ /* topo_group_connect — выбор пиров и попытки подключения CHAT-группы. - * Сначала ранее подключённые пиры (до 128 параллельно), затем по одному - * суперузлы/публичные узлы, затем локальные. Цель — READY суперузел или три - * READY немобильных пира. CONNECTING: 2 s; SYNCING: 5 s без прогресса. - * После исчерпания кандидатов новый цикл через 1 s или следующий standby burst. + * Суперузел подключается ко всем суперузлам канала. Desktop/mobile — к двум; + * без UP суперузлов резервируют два desktop, без них — два mobile. Проверки + * старших классов продолжаются независимо от резервных связей. Роль supernode + * не зависит от адресов. История/RTT упорядочивают кандидатов внутри класса. + * Успех — транспортный UP; группа удерживает сессию до DOWN, обмен таблицами + * продолжается отдельно. На класс не более двух попыток, CONNECTING: 2 s. + * Повторы с backoff 1..30 s или в standby burst. READY/recovery не изменяются. * Планировщик владеет запросами, группа — NCD. Остановка отменяет попытки, * включая ручные, но сохраняет готовые сессии. Все вызовы — в потоке uasync. */ @@ -14,19 +17,19 @@ struct TOPO_GROUP; struct ETCP_CONN; -/* Запустить новый автоматический цикл вместо прежнего. 0/-1. */ +/* Включить/обновить автоматический подбор, сохранив текущие запросы. 0/-1. */ 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); -/* Учесть появление пира; для READY записать connected в БД. */ +/* Учесть появление пира; для UP записать connected в БД. */ 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); -/* Число текущих READY-пиров группы. */ +/* Число текущих UP-пиров группы, независимо от READY. */ int topo_group_connect_active_count(struct TOPO_GROUP* group); -/* Перезапустить автоматический цикл, только если нет READY-пиров. */ +/* Сбросить backoff после изменения сети, сохранив UP и незавершённые запросы. */ void topo_group_connect_restart(struct TOPO_GROUP* group); /* Одна ручная попытка, общий запрос к сессии группы. 0=принято, -1=ошибка. */