You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
652 lines
27 KiB
652 lines
27 KiB
/* |
|
* chat_conn_mgr.c — менеджер ETCP-подключений (1 на инстанс) |
|
* |
|
* Управляет подключениями к узлам: перебор адресов, случайный выбор, |
|
* ротация при падении соединения. До CM_MAX_PARALLEL одновременных |
|
* подключений на узел. Фоновый пинг для измерения RTT. |
|
*/ |
|
|
|
#include "chat_conn_mgr.h" |
|
#include "chat_core.h" |
|
#include "gui_bridge.h" |
|
|
|
#include "../../../src/utun_instance.h" |
|
#include "../../../src/etcp_api.h" |
|
#include "../../../src/etcp.h" |
|
#include "../../../src/etcp_connections.h" |
|
#include "../../../src/secure_channel.h" |
|
#include "../../../lib/u_async.h" |
|
#include "../../../lib/ll_queue.h" |
|
#include "../../../lib/mem.h" |
|
#include "../../../lib/debug_config.h" |
|
|
|
#include <sqlite3.h> |
|
#include <string.h> |
|
#include <stdlib.h> |
|
#include "../../../lib/platform_compat.h" |
|
|
|
#define CM_ID "chat_conn_mgr" |
|
|
|
/* ─── константы ─── */ |
|
#define CM_MAX_PARALLEL 10 |
|
#define CM_MAX_ADDRS 16 |
|
#define CM_MAX_NODES 32 |
|
#define CM_INIT_TIMEOUT_MS 3000 |
|
#define CM_RETRY_MS 3000 |
|
#define CM_PING_MS 5000 |
|
#define CM_PING_TIMEOUT_MS 2000 |
|
#define CM_INIT_TIMEOUT_TB (CM_INIT_TIMEOUT_MS * 10) |
|
#define CM_RETRY_TB (CM_RETRY_MS * 10) |
|
#define CM_PING_TB (CM_PING_MS * 10) |
|
|
|
/* ─── состояния слота ─── */ |
|
#define CM_SLOT_FREE 0 |
|
#define CM_SLOT_CONNECTING 1 |
|
#define CM_SLOT_ACTIVE 2 |
|
#define CM_SLOT_DEAD 3 |
|
|
|
/* ─── типы ─── */ |
|
|
|
struct cm_cb { |
|
void (*cb)(int result, uint64_t node_id, void* arg); |
|
void* arg; |
|
struct cm_cb* next; |
|
}; |
|
|
|
struct cm_slot { |
|
uint8_t addr_idx; |
|
uint8_t state; |
|
uint16_t rtt; |
|
struct ETCP_CONN* conn; |
|
void* init_timer; |
|
}; |
|
|
|
struct cm_addr { |
|
struct sockaddr_storage sa; |
|
uint16_t rtt; |
|
}; |
|
|
|
struct cm_node; |
|
|
|
struct cm_slot_ctx { |
|
struct cm_node* node; |
|
int slot_idx; |
|
}; |
|
|
|
struct cm_node_changed_ctx { |
|
struct chat_conn_mgr* mgr; |
|
uint64_t node_id; |
|
}; |
|
|
|
struct cm_node { |
|
uint64_t node_id; |
|
uint8_t pubkey[SC_PUBKEY_SIZE]; |
|
struct cm_addr addrs[CM_MAX_ADDRS]; |
|
int addr_count; |
|
int addr_loaded; |
|
struct cm_slot slots[CM_MAX_PARALLEL]; |
|
struct cm_cb* cbs; |
|
int delivered; |
|
void* retry_timer; |
|
struct chat_conn_mgr* mgr; /* back-pointer */ |
|
}; |
|
|
|
struct chat_conn_mgr { |
|
struct UTUN_INSTANCE* inst; |
|
struct sqlite3* db; |
|
struct cm_node nodes[CM_MAX_NODES]; |
|
int node_count; |
|
void* ping_timer; |
|
}; |
|
|
|
/* ─── forward ─── */ |
|
static struct cm_node* cm_find_node(struct chat_conn_mgr* mgr, uint64_t node_id); |
|
static struct cm_node* cm_ensure_node(struct chat_conn_mgr* mgr, uint64_t node_id); |
|
static void cm_load_addrs(struct cm_node* node); |
|
static int cm_count_active(struct cm_node* node); |
|
static int cm_addr_taken(struct cm_node* node, int idx); |
|
static int cm_external_conn_active(struct cm_node* node); |
|
static struct cm_slot* cm_find_free(struct cm_node* node); |
|
static struct cm_slot* cm_find_dead(struct cm_node* node); |
|
static struct cm_slot* cm_find_slot_by_conn(struct cm_node* node, struct ETCP_CONN* conn); |
|
static void cm_deliver_all(struct cm_node* node, int result); |
|
static void cm_rotation(struct cm_node* node); |
|
static void cm_init_cb(struct ETCP_CONN* conn, void* arg); |
|
static void cm_init_timeout(void* arg); |
|
static void cm_conn_down_cb(struct ETCP_CONN* conn, void* arg); |
|
static void cm_retry_cb(void* arg); |
|
static void cm_ping_timer_cb(void* arg); |
|
static void cm_ping_result_cb(int success, uint16_t rtt, void* arg, uint64_t nonce, const uint8_t* resp_data, size_t resp_data_len); |
|
static struct ETCP_SOCKET* cm_best_socket(struct chat_conn_mgr* mgr); |
|
static void cm_invite_result_cb(int result, uint64_t node_id, void* arg); |
|
static void cm_node_changed_cb(struct ETCP_CONN* conn, void* arg); |
|
static void cm_node_changed_deferred(void* arg); |
|
|
|
/* ─── init / destroy ─── */ |
|
|
|
struct chat_conn_mgr* chat_conn_mgr_init(struct UTUN_INSTANCE* inst, struct sqlite3* db) { |
|
if (!inst || !db) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: init bad args", CM_ID); return NULL; } |
|
struct chat_conn_mgr* mgr = u_calloc(1, sizeof(struct chat_conn_mgr)); |
|
if (!mgr) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: alloc failed", CM_ID); return NULL; } |
|
mgr->inst = inst; mgr->db = db; |
|
mgr->ping_timer = uasync_set_timeout(inst->ua, CM_PING_TB, mgr, cm_ping_timer_cb, "cm_ping"); |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: initialized", CM_ID); |
|
return mgr; |
|
} |
|
|
|
void chat_conn_mgr_destroy(struct chat_conn_mgr* mgr) { |
|
if (!mgr) return; |
|
if (mgr->ping_timer) { uasync_cancel_timeout(mgr->inst->ua, mgr->ping_timer); mgr->ping_timer = NULL; } |
|
for (int ni = 0; ni < mgr->node_count; ni++) { |
|
struct cm_node* node = &mgr->nodes[ni]; |
|
if (node->retry_timer) { uasync_cancel_timeout(mgr->inst->ua, node->retry_timer); node->retry_timer = NULL; } |
|
for (int si = 0; si < CM_MAX_PARALLEL; si++) { |
|
struct cm_slot* slot = &node->slots[si]; |
|
if (slot->init_timer) { uasync_cancel_timeout(mgr->inst->ua, slot->init_timer); slot->init_timer = NULL; } |
|
if (slot->conn) { etcp_connection_close(slot->conn); slot->conn = NULL; } |
|
} |
|
cm_deliver_all(node, CC_ERR_INTERNAL); |
|
} |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: destroyed", CM_ID); |
|
u_free(mgr); |
|
} |
|
|
|
struct sqlite3* chat_conn_mgr_get_db(struct chat_conn_mgr* mgr) { |
|
return mgr ? mgr->db : NULL; |
|
} |
|
|
|
/* ─── управление узлами ─── */ |
|
|
|
static struct cm_node* cm_find_node(struct chat_conn_mgr* mgr, uint64_t node_id) { |
|
for (int i = 0; i < mgr->node_count; i++) |
|
if (mgr->nodes[i].node_id == node_id) return &mgr->nodes[i]; |
|
return NULL; |
|
} |
|
|
|
static struct cm_node* cm_ensure_node(struct chat_conn_mgr* mgr, uint64_t node_id) { |
|
struct cm_node* node = cm_find_node(mgr, node_id); |
|
if (node) return node; |
|
if (mgr->node_count >= CM_MAX_NODES) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: too many nodes", CM_ID); return NULL; } |
|
node = &mgr->nodes[mgr->node_count++]; |
|
memset(node, 0, sizeof(*node)); |
|
node->node_id = node_id; node->mgr = mgr; |
|
for (int i = 0; i < CM_MAX_PARALLEL; i++) node->slots[i].addr_idx = 255; |
|
return node; |
|
} |
|
|
|
void chat_conn_mgr_add_node(struct chat_conn_mgr* mgr, uint64_t node_id) { |
|
if (!mgr) return; |
|
if (cm_find_node(mgr, node_id)) return; |
|
cm_ensure_node(mgr, node_id); |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: node 0x%016llx added", CM_ID, (unsigned long long)node_id); |
|
} |
|
|
|
/* ─── загрузка адресов ─── */ |
|
|
|
static void cm_load_addrs(struct cm_node* node) { |
|
if (!node || node->addr_loaded) return; |
|
struct sqlite3* db = node->mgr->db; |
|
struct { uint8_t a[4]; uint16_t p; } raw[CM_MAX_ADDRS]; |
|
int rc = 0; |
|
sqlite3_stmt* st = NULL; |
|
if (sqlite3_prepare_v2(db, "SELECT address,port FROM node_addresses WHERE node_id=? AND family=4", |
|
-1, &st, NULL) == SQLITE_OK) { |
|
sqlite3_bind_int64(st, 1, (sqlite3_int64)node->node_id); |
|
while (sqlite3_step(st) == SQLITE_ROW && rc < CM_MAX_ADDRS) { |
|
const void* a = sqlite3_column_blob(st, 0); |
|
if (a && sqlite3_column_bytes(st, 0) == 4) { memcpy(raw[rc].a, a, 4); raw[rc].p = (uint16_t)sqlite3_column_int(st, 1); rc++; } |
|
} |
|
sqlite3_finalize(st); |
|
} |
|
int uniq = 0; |
|
for (int i = 0; i < rc; i++) { |
|
int dup = 0; |
|
for (int j = 0; j < uniq; j++) if (memcmp(raw[i].a, raw[j].a, 4) == 0 && raw[i].p == raw[j].p) { dup = 1; break; } |
|
if (!dup) { if (i != uniq) raw[uniq] = raw[i]; uniq++; } |
|
} |
|
for (int i = 0; i < uniq; i++) { |
|
struct sockaddr_in sin; memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET; |
|
memcpy(&sin.sin_addr.s_addr, raw[i].a, 4); sin.sin_port = htons(raw[i].p); |
|
memcpy(&node->addrs[i].sa, &sin, sizeof(sin)); node->addrs[i].rtt = 65535; |
|
} |
|
node->addr_count = uniq; node->addr_loaded = 1; |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: node 0x%016llx loaded %d addrs", CM_ID, (unsigned long long)node->node_id, uniq); |
|
|
|
uint8_t pubkey[SC_PUBKEY_SIZE] = {0}; |
|
if (sqlite3_prepare_v2(db, "SELECT x25519_pubkey FROM nodes WHERE node_id=?", |
|
-1, &st, NULL) == SQLITE_OK) { |
|
sqlite3_bind_int64(st, 1, (sqlite3_int64)node->node_id); |
|
if (sqlite3_step(st) == SQLITE_ROW) { |
|
const void* pk = sqlite3_column_blob(st, 0); |
|
if (pk && sqlite3_column_bytes(st, 0) == SC_PUBKEY_SIZE) memcpy(pubkey, pk, SC_PUBKEY_SIZE); |
|
} |
|
sqlite3_finalize(st); |
|
} |
|
if (pubkey[0] || pubkey[1]) memcpy(node->pubkey, pubkey, SC_PUBKEY_SIZE); |
|
else DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: node 0x%016llx no pubkey in DB", CM_ID, (unsigned long long)node->node_id); |
|
} |
|
|
|
/* ─── helpers ─── */ |
|
|
|
static int cm_count_active(struct cm_node* node) { |
|
int a = 0; |
|
for (int i = 0; i < CM_MAX_PARALLEL; i++) if (node->slots[i].state == CM_SLOT_ACTIVE) a++; |
|
return a; |
|
} |
|
|
|
static int cm_addr_taken(struct cm_node* node, int idx) { |
|
for (int i = 0; i < CM_MAX_PARALLEL; i++) |
|
if (node->slots[i].state != CM_SLOT_FREE && node->slots[i].addr_idx == idx) return 1; |
|
return 0; |
|
} |
|
|
|
static int cm_external_conn_active(struct cm_node* node) { |
|
struct UTUN_INSTANCE* inst = node->mgr->inst; |
|
struct ll_entry* e = queue_find_data_by_index(inst->connections, (const uint8_t*)&node->node_id); |
|
if (!e) return 0; |
|
struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; |
|
return (ce && ce->conn && ce->conn->links_up && ce->conn->peer_node_id == node->node_id); |
|
} |
|
|
|
static struct cm_slot* cm_find_free(struct cm_node* node) { |
|
for (int i = 0; i < CM_MAX_PARALLEL; i++) |
|
if (node->slots[i].state == CM_SLOT_FREE) return &node->slots[i]; |
|
return NULL; |
|
} |
|
|
|
static struct cm_slot* cm_find_dead(struct cm_node* node) { |
|
for (int i = 0; i < CM_MAX_PARALLEL; i++) |
|
if (node->slots[i].state == CM_SLOT_DEAD) return &node->slots[i]; |
|
return NULL; |
|
} |
|
|
|
static struct cm_slot* cm_find_slot_by_conn(struct cm_node* node, struct ETCP_CONN* conn) { |
|
for (int i = 0; i < CM_MAX_PARALLEL; i++) |
|
if (node->slots[i].conn == conn) return &node->slots[i]; |
|
return NULL; |
|
} |
|
|
|
static void cm_deliver_all(struct cm_node* node, int result) { |
|
while (node->cbs) { struct cm_cb* cb = node->cbs; node->cbs = cb->next; cb->cb(result, node->node_id, cb->arg); u_free(cb); } |
|
} |
|
|
|
static struct ETCP_SOCKET* cm_best_socket(struct chat_conn_mgr* mgr) { |
|
struct ETCP_SOCKET* s = mgr->inst->etcp_sockets; |
|
while (s) { if (s->local_addr.ss_family == AF_INET) return s; s = s->next; } |
|
return NULL; |
|
} |
|
|
|
/* ─── ротация ─── */ |
|
|
|
static void cm_rotation(struct cm_node* node) { |
|
if (!node) return; |
|
struct chat_conn_mgr* mgr = node->mgr; |
|
|
|
int active = cm_count_active(node); |
|
|
|
if (cm_external_conn_active(node)) { |
|
if (!node->delivered) { node->delivered = 1; cm_deliver_all(node, CC_OK); } |
|
if (node->retry_timer) { uasync_cancel_timeout(mgr->inst->ua, node->retry_timer); node->retry_timer = NULL; } |
|
return; |
|
} |
|
|
|
if (active >= CM_MAX_PARALLEL) { |
|
if (!node->delivered) { node->delivered = 1; cm_deliver_all(node, CC_OK); } |
|
if (node->retry_timer) { uasync_cancel_timeout(mgr->inst->ua, node->retry_timer); node->retry_timer = NULL; } |
|
return; |
|
} |
|
|
|
struct cm_slot* slot = cm_find_free(node); |
|
if (!slot) slot = cm_find_dead(node); |
|
if (!slot) return; |
|
|
|
struct ETCP_SOCKET* sock = cm_best_socket(mgr); |
|
if (!sock) return; |
|
|
|
for (int attempt = 0; attempt < 3; attempt++) { |
|
int idx = rand() % node->addr_count; |
|
if (!cm_addr_taken(node, idx)) { |
|
slot->addr_idx = idx; slot->state = CM_SLOT_CONNECTING; |
|
struct ETCP_CONN* conn = etcp_connection_create(mgr->inst, NULL); |
|
if (!conn) { slot->state = CM_SLOT_FREE; return; } |
|
sc_init_ctx(&conn->crypto_ctx, &mgr->inst->my_keys); |
|
sc_set_peer_public_key(&conn->crypto_ctx, node->pubkey, 0); |
|
etcp_conn_set_peer_node_id(conn, node->node_id); |
|
struct cm_slot_ctx* ctx = u_calloc(1, sizeof(struct cm_slot_ctx)); |
|
if (!ctx) { etcp_connection_close(conn); slot->state = CM_SLOT_FREE; return; } |
|
ctx->node = node; ctx->slot_idx = (int)(slot - node->slots); |
|
etcp_conn_add_init_cbk(conn, cm_init_cb, ctx); |
|
etcp_conn_add_node_changed_cbk(conn, cm_node_changed_cb, mgr); |
|
if (!etcp_link_new(conn, sock, &node->addrs[idx].sa, 0)) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: link_new failed node=0x%016llx idx=%d", CM_ID, (unsigned long long)node->node_id, idx); |
|
etcp_connection_close(conn); u_free(ctx); slot->state = CM_SLOT_DEAD; |
|
cm_rotation(node); return; |
|
} |
|
slot->conn = conn; |
|
slot->init_timer = uasync_set_timeout(mgr->inst->ua, CM_INIT_TIMEOUT_TB, ctx, cm_init_timeout, "cm_init"); |
|
DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTIVITY, "%s: INIT node=0x%016llx slot=%d addr=%d", CM_ID, (unsigned long long)node->node_id, ctx->slot_idx, idx); |
|
return; |
|
} |
|
} |
|
|
|
if (slot->addr_idx != 255) { |
|
slot->state = CM_SLOT_CONNECTING; |
|
struct ETCP_CONN* conn = etcp_connection_create(mgr->inst, NULL); |
|
if (!conn) { slot->state = CM_SLOT_FREE; return; } |
|
sc_init_ctx(&conn->crypto_ctx, &mgr->inst->my_keys); |
|
sc_set_peer_public_key(&conn->crypto_ctx, node->pubkey, 0); |
|
etcp_conn_set_peer_node_id(conn, node->node_id); |
|
struct cm_slot_ctx* ctx = u_calloc(1, sizeof(struct cm_slot_ctx)); |
|
if (!ctx) { etcp_connection_close(conn); slot->state = CM_SLOT_FREE; return; } |
|
ctx->node = node; ctx->slot_idx = (int)(slot - node->slots); |
|
etcp_conn_add_init_cbk(conn, cm_init_cb, ctx); |
|
etcp_conn_add_node_changed_cbk(conn, cm_node_changed_cb, mgr); |
|
if (!etcp_link_new(conn, sock, &node->addrs[slot->addr_idx].sa, 0)) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: link_new(retry) failed node=0x%016llx idx=%d", CM_ID, (unsigned long long)node->node_id, slot->addr_idx); |
|
etcp_connection_close(conn); u_free(ctx); slot->state = CM_SLOT_DEAD; |
|
cm_rotation(node); return; |
|
} |
|
slot->conn = conn; |
|
slot->init_timer = uasync_set_timeout(mgr->inst->ua, CM_INIT_TIMEOUT_TB, ctx, cm_init_timeout, "cm_init"); |
|
DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTIVITY, "%s: INIT(retry) node=0x%016llx slot=%d addr=%d", CM_ID, (unsigned long long)node->node_id, ctx->slot_idx, slot->addr_idx); |
|
return; |
|
} |
|
|
|
slot->state = CM_SLOT_FREE; |
|
} |
|
|
|
/* ─── коллбэки ─── */ |
|
|
|
static void cm_init_cb(struct ETCP_CONN* conn, void* arg) { |
|
struct cm_slot_ctx* ctx = (struct cm_slot_ctx*)arg; |
|
struct cm_node* node = ctx->node; |
|
struct cm_slot* slot = &node->slots[ctx->slot_idx]; |
|
|
|
if (slot->init_timer) { uasync_cancel_timeout(node->mgr->inst->ua, slot->init_timer); slot->init_timer = NULL; } |
|
|
|
if (slot->conn && slot->conn != conn) { etcp_connection_close(slot->conn); slot->conn = NULL; } |
|
|
|
slot->state = CM_SLOT_ACTIVE; slot->conn = conn; slot->rtt = conn->rtt_last; |
|
etcp_conn_add_down_cbk(conn, cm_conn_down_cb, node); |
|
etcp_conn_remove_init_cbk(conn, cm_init_cb, ctx); |
|
u_free(ctx); |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: INIT OK node=0x%016llx slot=%d rtt=%u active=%d", |
|
CM_ID, (unsigned long long)node->node_id, (int)(slot - node->slots), slot->rtt, cm_count_active(node)); |
|
|
|
struct chat_conn_mgr* mgr = node->mgr; |
|
int active = cm_count_active(node); |
|
|
|
if (active > CM_MAX_PARALLEL) { |
|
int worst_idx = -1; uint16_t worst_rtt = 0; |
|
for (int i = 0; i < CM_MAX_PARALLEL; i++) { |
|
if (node->slots[i].state == CM_SLOT_ACTIVE && node->slots[i].rtt > worst_rtt) { |
|
worst_rtt = node->slots[i].rtt; worst_idx = i; |
|
} |
|
} |
|
if (worst_idx >= 0) { |
|
struct cm_slot* w = &node->slots[worst_idx]; |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: evict slot=%d rtt=%u node=0x%016llx", |
|
CM_ID, worst_idx, w->rtt, (unsigned long long)node->node_id); |
|
w->state = CM_SLOT_FREE; w->conn = NULL; w->addr_idx = 255; |
|
} |
|
active = cm_count_active(node); |
|
} |
|
|
|
if (active >= CM_MAX_PARALLEL) { |
|
if (!node->delivered) { node->delivered = 1; cm_deliver_all(node, CC_OK); } |
|
if (node->retry_timer) { uasync_cancel_timeout(mgr->inst->ua, node->retry_timer); node->retry_timer = NULL; } |
|
return; |
|
} |
|
|
|
cm_rotation(node); |
|
} |
|
|
|
static void cm_init_timeout(void* arg) { |
|
struct cm_slot_ctx* ctx = (struct cm_slot_ctx*)arg; |
|
struct cm_node* node = ctx->node; |
|
struct cm_slot* slot = &node->slots[ctx->slot_idx]; |
|
slot->init_timer = NULL; |
|
|
|
if (slot->state == CM_SLOT_CONNECTING) slot->state = CM_SLOT_DEAD; |
|
node->delivered = 0; |
|
|
|
DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTIVITY, "%s: INIT TIMEOUT node=0x%016llx slot=%d", |
|
CM_ID, (unsigned long long)node->node_id, ctx->slot_idx); |
|
|
|
cm_rotation(node); |
|
u_free(ctx); |
|
} |
|
|
|
static void cm_conn_down_cb(struct ETCP_CONN* conn, void* arg) { |
|
struct cm_node* node = (struct cm_node*)arg; |
|
struct cm_slot* slot = cm_find_slot_by_conn(node, conn); |
|
if (!slot) return; |
|
|
|
slot->state = CM_SLOT_FREE; slot->conn = NULL; |
|
node->delivered = 0; |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: DOWN node=0x%016llx slot=%d active=%d", |
|
CM_ID, (unsigned long long)node->node_id, (int)(slot - node->slots), cm_count_active(node)); |
|
|
|
cm_rotation(node); |
|
} |
|
|
|
static void cm_retry_cb(void* arg) { |
|
struct cm_node* node = (struct cm_node*)arg; |
|
node->retry_timer = NULL; |
|
|
|
for (int i = 0; i < CM_MAX_PARALLEL; i++) |
|
if (node->slots[i].state == CM_SLOT_DEAD) { node->slots[i].state = CM_SLOT_FREE; node->slots[i].addr_idx = 255; } |
|
|
|
DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTIVITY, "%s: retry node=0x%016llx active=%d", CM_ID, (unsigned long long)node->node_id, cm_count_active(node)); |
|
|
|
cm_rotation(node); |
|
|
|
if (cm_count_active(node) < CM_MAX_PARALLEL && !node->delivered) |
|
node->retry_timer = uasync_set_timeout(node->mgr->inst->ua, CM_RETRY_TB, node, cm_retry_cb, "cm_retry"); |
|
} |
|
|
|
/* ─── node_changed ─── */ |
|
|
|
static void cm_node_changed_deferred(void* arg) { |
|
struct cm_node_changed_ctx* ctx = (struct cm_node_changed_ctx*)arg; |
|
struct cm_node* node = cm_find_node(ctx->mgr, ctx->node_id); |
|
if (node) { |
|
for (int i = 0; i < CM_MAX_PARALLEL; i++) { |
|
struct cm_slot* slot = &node->slots[i]; |
|
if (slot->conn) { etcp_connection_close(slot->conn); slot->conn = NULL; } |
|
slot->state = CM_SLOT_FREE; slot->addr_idx = 255; |
|
} |
|
cm_deliver_all(node, CC_ERR_INTERNAL); |
|
node->delivered = 0; |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: NODE_CHANGED deferred close done node=0x%016llx", |
|
CM_ID, (unsigned long long)ctx->node_id); |
|
} |
|
u_free(ctx); |
|
} |
|
|
|
static void cm_node_changed_cb(struct ETCP_CONN* conn, void* arg) { |
|
struct chat_conn_mgr* mgr = (struct chat_conn_mgr*)arg; |
|
if (!mgr || !conn) return; |
|
uint64_t node_id = conn->peer_node_id; |
|
struct cm_node* node = cm_find_node(mgr, node_id); |
|
if (!node) return; |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: NODE_CHANGED node=0x%016llx — clearing addresses and closing connections", |
|
CM_ID, (unsigned long long)node_id); |
|
|
|
sqlite3_stmt* st = NULL; |
|
if (sqlite3_prepare_v2(mgr->db, "DELETE FROM node_addresses WHERE node_id=?", -1, &st, NULL) == SQLITE_OK) { |
|
sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); |
|
sqlite3_step(st); sqlite3_finalize(st); |
|
} |
|
|
|
node->addr_loaded = 0; |
|
node->addr_count = 0; |
|
|
|
if (node->retry_timer) { uasync_cancel_timeout(mgr->inst->ua, node->retry_timer); node->retry_timer = NULL; } |
|
|
|
for (int i = 0; i < CM_MAX_PARALLEL; i++) { |
|
struct cm_slot* slot = &node->slots[i]; |
|
if (slot->init_timer) { uasync_cancel_timeout(mgr->inst->ua, slot->init_timer); slot->init_timer = NULL; } |
|
if (slot->conn && slot->conn != conn) { etcp_connection_close(slot->conn); slot->conn = NULL; } |
|
if (slot->conn != conn) { slot->state = CM_SLOT_FREE; slot->addr_idx = 255; } |
|
} |
|
|
|
struct cm_node_changed_ctx* ctx = u_calloc(1, sizeof(*ctx)); |
|
if (ctx) { ctx->mgr = mgr; ctx->node_id = node_id; uasync_post(mgr->inst->ua, cm_node_changed_deferred, ctx); } |
|
|
|
uint8_t evt[8]; memcpy(evt, &node_id, 8); |
|
gui_bridge_post(GUI_EVT_NODE_CHANGED, evt, 8); |
|
} |
|
|
|
/* ─── фоновый пинг ─── */ |
|
|
|
static void cm_ping_timer_cb(void* arg) { |
|
struct chat_conn_mgr* mgr = (struct chat_conn_mgr*)arg; |
|
if (mgr->node_count == 0) { mgr->ping_timer = uasync_set_timeout(mgr->inst->ua, CM_PING_TB, mgr, cm_ping_timer_cb, "cm_ping"); return; } |
|
|
|
int ni = rand() % mgr->node_count; |
|
struct cm_node* node = &mgr->nodes[ni]; |
|
if (!node->addr_loaded || node->addr_count == 0) { mgr->ping_timer = uasync_set_timeout(mgr->inst->ua, CM_PING_TB, mgr, cm_ping_timer_cb, "cm_ping"); return; } |
|
|
|
int ai = rand() % node->addr_count; |
|
|
|
if (cm_addr_taken(node, ai)) { mgr->ping_timer = uasync_set_timeout(mgr->inst->ua, CM_PING_TB, mgr, cm_ping_timer_cb, "cm_ping"); return; } |
|
|
|
struct ETCP_SOCKET* sock = cm_best_socket(mgr); |
|
if (sock) etcp_send_ping_to_socket(mgr->inst, sock, node->pubkey, &node->addrs[ai].sa, CM_PING_TIMEOUT_MS, cm_ping_result_cb, &node->addrs[ai], NULL, 0); |
|
|
|
mgr->ping_timer = uasync_set_timeout(mgr->inst->ua, CM_PING_TB, mgr, cm_ping_timer_cb, "cm_ping"); |
|
} |
|
|
|
static void cm_ping_result_cb(int success, uint16_t rtt, void* arg, uint64_t nonce, const uint8_t* resp_data, size_t resp_data_len) { |
|
(void)nonce; (void)resp_data; (void)resp_data_len; |
|
struct cm_addr* addr = (struct cm_addr*)arg; |
|
addr->rtt = success ? rtt : 65535; |
|
} |
|
|
|
/* ─── публичное API ─── */ |
|
|
|
void chat_conn_mgr_connect(struct chat_conn_mgr* mgr, uint64_t node_id, |
|
void (*cb)(int result, uint64_t node_id, void* arg), |
|
void* arg) { |
|
if (!mgr) { if (cb) cb(CC_ERR_INTERNAL, node_id, arg); return; } |
|
|
|
struct cm_node* node = cm_ensure_node(mgr, node_id); |
|
if (!node) { if (cb) cb(CC_ERR_INTERNAL, node_id, arg); return; } |
|
|
|
if (!node->addr_loaded) cm_load_addrs(node); |
|
|
|
if (cb) { |
|
struct cm_cb* cb_node = u_calloc(1, sizeof(struct cm_cb)); |
|
if (cb_node) { cb_node->cb = cb; cb_node->arg = arg; cb_node->next = node->cbs; node->cbs = cb_node; } |
|
} |
|
|
|
if (node->pubkey[0] == 0 && node->pubkey[1] == 0) { |
|
DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: connect to 0x%016llx — no pubkey in DB", CM_ID, (unsigned long long)node_id); |
|
if (!node->delivered) { node->delivered = 1; cm_deliver_all(node, CC_ERR_NOT_FOUND); } |
|
return; |
|
} |
|
if (node->addr_count == 0) { |
|
DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: connect to 0x%016llx — no addresses", CM_ID, (unsigned long long)node_id); |
|
if (!node->delivered) { node->delivered = 1; cm_deliver_all(node, CC_ERR_NO_ADDRESSES); } |
|
return; |
|
} |
|
|
|
if (node->delivered) return; |
|
|
|
if (cm_external_conn_active(node)) { |
|
node->delivered = 1; cm_deliver_all(node, CC_OK); |
|
DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTIVITY, "%s: node 0x%016llx already connected externally", CM_ID, (unsigned long long)node_id); |
|
return; |
|
} |
|
|
|
if (cm_count_active(node) >= 1) { |
|
// at least one active already (might be from a previous rotation that didn't deliver yet) |
|
if (cm_count_active(node) >= CM_MAX_PARALLEL && !node->delivered) { node->delivered = 1; cm_deliver_all(node, CC_OK); } |
|
return; |
|
} |
|
|
|
cm_rotation(node); |
|
|
|
if (cm_count_active(node) == 0 && !node->retry_timer && !node->delivered) |
|
node->retry_timer = uasync_set_timeout(mgr->inst->ua, CM_RETRY_TB, node, cm_retry_cb, "cm_retry"); |
|
} |
|
|
|
void chat_conn_mgr_cancel(struct chat_conn_mgr* mgr, uint64_t node_id) { |
|
if (!mgr) return; |
|
struct cm_node* node = cm_find_node(mgr, node_id); |
|
if (!node) return; |
|
|
|
if (node->retry_timer) { uasync_cancel_timeout(mgr->inst->ua, node->retry_timer); node->retry_timer = NULL; } |
|
for (int i = 0; i < CM_MAX_PARALLEL; i++) { |
|
struct cm_slot* slot = &node->slots[i]; |
|
if (slot->init_timer) { uasync_cancel_timeout(mgr->inst->ua, slot->init_timer); slot->init_timer = NULL; } |
|
if (slot->conn) { etcp_connection_close(slot->conn); slot->conn = NULL; } |
|
slot->state = CM_SLOT_FREE; slot->addr_idx = 255; |
|
} |
|
cm_deliver_all(node, CC_ERR_INTERNAL); |
|
node->delivered = 0; |
|
DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTIVITY, "%s: cancelled node 0x%016llx", CM_ID, (unsigned long long)node_id); |
|
} |
|
|
|
void chat_conn_mgr_connect_from_invite(struct chat_conn_mgr* mgr, struct chat_invite* inv) { |
|
if (!mgr || !inv) return; |
|
uint64_t node_id = inv->node_id; |
|
|
|
/* save pubkey */ |
|
sqlite3_stmt* st = NULL; |
|
if (sqlite3_prepare_v2(mgr->db, "INSERT INTO nodes(node_id,x25519_pubkey) VALUES(?,?)" |
|
" ON CONFLICT(node_id) DO UPDATE SET x25519_pubkey=excluded.x25519_pubkey", |
|
-1, &st, NULL) == SQLITE_OK) { |
|
sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); |
|
sqlite3_bind_blob(st, 2, inv->pubkey, 32, SQLITE_STATIC); |
|
sqlite3_step(st); sqlite3_finalize(st); |
|
} |
|
|
|
/* save addresses (socket_id=0) */ |
|
sqlite3_stmt* ds = NULL; |
|
if (sqlite3_prepare_v2(mgr->db, "DELETE FROM node_addresses WHERE node_id=? AND socket_id=0", |
|
-1, &ds, NULL) == SQLITE_OK) { |
|
sqlite3_bind_int64(ds, 1, (sqlite3_int64)node_id); |
|
sqlite3_step(ds); sqlite3_finalize(ds); |
|
} |
|
sqlite3_stmt* is = NULL; |
|
if (sqlite3_prepare_v2(mgr->db, "INSERT INTO node_addresses(node_id,family,protocol,address,port,addr_type,socket_id)" |
|
" VALUES(?,?,1,?,?,0,?)", -1, &is, NULL) == SQLITE_OK) { |
|
const uint8_t* src = inv->addrs_data; |
|
for (int i = 0; i < inv->addr_count; i++) { |
|
uint8_t family = *src++; uint8_t sid = *src++; |
|
if (family == 4) { |
|
sqlite3_bind_int64(is, 1, (sqlite3_int64)node_id); |
|
sqlite3_bind_int(is, 2, 4); |
|
sqlite3_bind_blob(is, 3, src, 4, SQLITE_STATIC); src += 4; |
|
uint16_t port = ((uint16_t)src[0] << 8) | src[1]; src += 2; |
|
sqlite3_bind_int(is, 4, (int)port); |
|
sqlite3_bind_int(is, 5, (int)sid); |
|
sqlite3_step(is); sqlite3_reset(is); |
|
} else src += 18; |
|
} |
|
sqlite3_finalize(is); |
|
} |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: invite node=0x%016llx addrs=%d — saved, connecting", |
|
CM_ID, (unsigned long long)node_id, inv->addr_count); |
|
|
|
/* add node + reset addr_loaded to reload from DB */ |
|
struct cm_node* node = cm_ensure_node(mgr, node_id); |
|
if (node) { node->addr_loaded = 0; cm_load_addrs(node); } |
|
|
|
{ char ch_str[32]; snprintf(ch_str, sizeof(ch_str), "%llu", (unsigned long long)inv->channel_id); |
|
chat_core_ensure_channel_ready(ch_str); } |
|
|
|
chat_conn_mgr_connect(mgr, node_id, cm_invite_result_cb, &inv->channel_id); |
|
} |
|
|
|
static void cm_invite_result_cb(int result, uint64_t node_id, void* arg) { |
|
uint64_t channel_id = arg ? *(uint64_t*)arg : 0; |
|
uint8_t data[20]; memcpy(data, &node_id, 8); memcpy(data + 8, &result, 4); memcpy(data + 12, &channel_id, 8); |
|
gui_bridge_post(GUI_EVT_CONNECT_RESULT, data, 20); |
|
}
|
|
|