Browse Source

Select chat auto peers by device class and transport UP

master
evgeny 15 hours ago
parent
commit
a7632cfee1
  1. 308
      src/routing_layer/topo_group_connect.c
  2. 19
      src/routing_layer/topo_group_connect.h

308
src/routing_layer/topo_group_connect.c

@ -1,6 +1,6 @@
/* Политика CHAT-подключений. Транспортом и JOIN владеет групповая сессия.
* Сначала параллельно восстанавливаем connected-пиров, затем последовательно
* пробуем supernode/public/local. Незавершённые попытки — отменяемые запросы. */
* Суперузел соединяется со всеми суперузлами, клиент — с двумя. Без UP
* суперузлов резервируем desktop, затем mobile, продолжая проверки старших классов. */
#include <string.h>
#include <limits.h>
#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;

19
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=ошибка. */

Loading…
Cancel
Save