Browse Source

conn_mgr: keep demand alive across bounded recovery attempts

master
evgeny 2 days ago
parent
commit
8f5f2b089c
  1. 157
      src/routing_layer/conn_mgr_core.c
  2. 34
      src/routing_layer/conn_mgr_indirect.c
  3. 7
      src/routing_layer/conn_mgr_priv.h

157
src/routing_layer/conn_mgr_core.c

@ -32,6 +32,7 @@
void conn_mgr_router_recv_handler(struct ETCP_CONN* conn, struct ll_entry* entry);
static void cm_attempt_tick(void* arg);
static void cm_start_attempt(struct CONN_MGR_ENTRY* entry, int recovery);
static const char* cm_path_name(uint8_t type) {
static const char* names[] = { "NONE", "DIRECT", "REVERSE", "INDIRECT" };
@ -260,9 +261,15 @@ void cm_refresh_path(struct CONN_MGR_ENTRY* entry) {
uint8_t type = CONN_TYPE_NONE;
struct ETCP_CONN* direct = cm_direct_conn(entry->mgr, entry->node_id);
if (direct) type = conn_mgr_connection_type(direct, entry->node_id);
else if (entry->indirect_confirmed) {
if (entry->indirect_confirmed) {
int reserve_up = 0;
for (uint8_t i = 0; i < entry->intermediariy_count; i++)
if (cm_intermediary_conn(entry, entry->intermediaries[i])) { type = CONN_TYPE_INDIRECT; break; }
if (cm_intermediary_conn(entry, entry->intermediaries[i])) { reserve_up = 1; break; }
if (reserve_up && !direct) type = CONN_TYPE_INDIRECT;
if (!reserve_up) {
entry->indirect_confirmed = 0;
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: reserve unavailable peer=%016llx", (unsigned long long)entry->node_id);
}
}
int was_up = entry->state == CONN_MGR_STATE_CONNECTED;
if (type == entry->conn_type && was_up == (type != CONN_TYPE_NONE)) return;
@ -271,8 +278,13 @@ void cm_refresh_path(struct CONN_MGR_ENTRY* entry) {
(unsigned long long)entry->mgr->group->group_id, (unsigned long long)entry->node_id,
cm_path_name(entry->conn_type), cm_path_name(type), (unsigned long long)(next ? next->peer_node_id : 0),
(unsigned long long)((get_time_tb() - entry->start_tb) / 10));
int lost_direct = entry->conn_type == CONN_TYPE_DIRECT || entry->conn_type == CONN_TYPE_REVERSE;
lost_direct &= type != CONN_TYPE_DIRECT && type != CONN_TYPE_REVERSE;
entry->conn_type = type;
entry->state = type ? CONN_MGR_STATE_CONNECTED : CONN_MGR_STATE_CONNECTING;
if (entry->handles && (lost_direct || (was_up && !type))) {
entry->attempt_active = 0; entry->retry_tb = get_time_tb(); entry->retry_delay_ms = CONN_MGR_RETRY_MIN_MS;
}
cm_update_nodeinfo(entry);
if (entry->state == CONN_MGR_STATE_DISCONNECTED) return;
if (!was_up && type) { entry->ever_up = 1; cm_deliver_event(entry, CONN_EVENT_UP); }
@ -281,6 +293,13 @@ void cm_refresh_path(struct CONN_MGR_ENTRY* entry) {
static void cm_refresh_cb(void* arg) {
struct CONN_MGR_ENTRY* entry = arg; entry->refresh_token = NULL;
struct TOPO_NODE* ni = topo_node_registry_find(entry->mgr->instance->topo_groups, entry->node_id);
if (entry->handles && ni && ni->timestamp != entry->node_timestamp && !cm_direct_conn(entry->mgr, entry->node_id)) {
entry->attempt_active = 0; entry->retry_tb = get_time_tb(); entry->retry_delay_ms = CONN_MGR_RETRY_MIN_MS;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: addresses changed peer=%016llx version=%llu; retry now",
(unsigned long long)entry->node_id, (unsigned long long)ni->timestamp);
entry->node_timestamp = ni->timestamp;
}
cm_refresh_path(entry);
if (entry->state != CONN_MGR_STATE_DISCONNECTED) cm_arm_timer(entry);
}
@ -297,6 +316,19 @@ void conn_mgr_routes_changed(struct CONN_MGR* mgr) {
}
}
/* Изменение собственных сокетов обходит backoff, но не прерывает рабочую доставку. */
void conn_mgr_network_changed(struct CONN_MGR* mgr) {
if (!mgr || !mgr->initialized) return;
for (struct ll_entry* e = mgr->entries->head; e; e = e->next) {
struct CONN_MGR_ENTRY* entry = (struct CONN_MGR_ENTRY*)e;
if (!entry->handles || cm_direct_conn(mgr, entry->node_id)) continue;
entry->attempt_active = 0; entry->retry_tb = get_time_tb(); entry->retry_delay_ms = CONN_MGR_RETRY_MIN_MS;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: local network changed peer=%016llx; retry now",
(unsigned long long)entry->node_id);
}
conn_mgr_routes_changed(mgr);
}
/* TIMEOUT первого NCD не закрывает транспорт: поздний UP всё ещё может улучшить путь. */
void cm_ncd_callback(struct NODE_CONN_DIRECT* ncd_h, enum ncd_event ev, void* arg) {
struct CONN_MGR_ENTRY* entry = arg;
@ -319,6 +351,9 @@ void cm_ncd_callback(struct NODE_CONN_DIRECT* ncd_h, enum ncd_event ev, void* ar
if (ev == NCD_EVENT_CLOSED) { entry->ncd_handle = NULL; node_conn_direct_close(ncd_h); }
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: direct transport lost peer=%016llx event=%d",
(unsigned long long)entry->node_id, ev);
if (entry->handles) {
entry->attempt_active = 0; entry->retry_tb = get_time_tb(); entry->retry_delay_ms = CONN_MGR_RETRY_MIN_MS;
}
}
cm_refresh_path(entry);
if (entry->state != CONN_MGR_STATE_DISCONNECTED) cm_arm_timer(entry);
@ -335,18 +370,29 @@ int cm_send_control(struct CONN_MGR* mgr, uint64_t node_id, const void* packet,
}
void cm_arm_timer(struct CONN_MGR_ENTRY* entry) {
if (entry->state == CONN_MGR_STATE_DISCONNECTED || entry->timer) return;
uint64_t now = get_time_tb(), when = entry->deadline_tb;
if (entry->state == CONN_MGR_STATE_DISCONNECTED) return;
if (entry->timer) { uasync_cancel_timeout(entry->mgr->instance->ua, entry->timer); entry->timer = NULL; }
uint64_t now = get_time_tb(), when = UINT64_MAX;
int direct_ready = entry->conn_type == CONN_TYPE_DIRECT || entry->conn_type == CONN_TYPE_REVERSE;
if (entry->timeout_notified || (entry->handles && direct_ready && (!entry->indirect_started || entry->indirect_confirmed))) return;
if (entry->handles && !direct_ready) {
if (!entry->handles) when = entry->deadline_tb;
else if (direct_ready && (!entry->indirect_started || entry->indirect_confirmed)) {
entry->attempt_active = 0; entry->retry_tb = 0; entry->retry_delay_ms = CONN_MGR_RETRY_MIN_MS;
} else if (entry->attempt_active) when = entry->deadline_tb;
else if (!direct_ready) {
if (!entry->retry_tb) entry->retry_tb = now;
when = entry->retry_tb;
}
if (entry->handles && !entry->ever_up && !entry->timeout_notified && entry->initial_deadline_tb < when)
when = entry->initial_deadline_tb;
if (entry->handles && entry->attempt_active && !direct_ready) {
uint64_t reverse = entry->reverse_started ? entry->reverse_next_tb : entry->reverse_start_tb;
if (reverse < when) when = reverse;
}
if (entry->handles && !entry->indirect_confirmed && (!direct_ready || entry->indirect_started)) {
if (entry->handles && entry->attempt_active && !entry->indirect_confirmed && (!direct_ready || entry->indirect_started)) {
uint64_t indirect = entry->indirect_started ? entry->indirect_next_tb : entry->indirect_start_tb;
if (indirect < when) when = indirect;
}
if (when == UINT64_MAX) return;
entry->timer = uasync_set_timeout(entry->mgr->instance->ua, when > now ? (int)(when - now) : 1,
entry, cm_attempt_tick, "conn_mgr_attempt");
if (!entry->timer) DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "conn_mgr: timer allocation failed peer=%016llx",
@ -359,25 +405,64 @@ static void cm_attempt_tick(void* arg) {
if (entry->state == CONN_MGR_STATE_DISCONNECTED) return;
uint64_t now = get_time_tb();
int direct_ready = entry->conn_type == CONN_TYPE_DIRECT || entry->conn_type == CONN_TYPE_REVERSE;
if (entry->handles && now < entry->deadline_tb) {
if (!entry->handles) {
if (now >= entry->deadline_tb) cm_entry_cleanup(entry);
else cm_arm_timer(entry);
return;
}
if (!entry->ever_up && !entry->timeout_notified && now >= entry->initial_deadline_tb) {
entry->timeout_notified = 1;
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: initial connection timeout peer=%016llx; retaining demand",
(unsigned long long)entry->node_id);
cm_deliver_event(entry, CONN_EVENT_TIMEOUT);
if (entry->state == CONN_MGR_STATE_DISCONNECTED) return;
}
if (!direct_ready && !entry->attempt_active && now >= entry->retry_tb) cm_start_attempt(entry, 1);
if (entry->state == CONN_MGR_STATE_DISCONNECTED) return;
if (entry->attempt_active && now < entry->deadline_tb) {
if (!entry->ncd_handle && !direct_ready) cm_start_phase_direct(entry);
if (!direct_ready && now >= entry->reverse_start_tb && now >= entry->reverse_next_tb) cm_start_phase_reverse(entry);
if ((!direct_ready || entry->indirect_started) && !entry->indirect_confirmed &&
now >= entry->indirect_start_tb && now >= entry->indirect_next_tb)
cm_start_phase_indirect(entry);
}
if (now >= entry->deadline_tb && !entry->timeout_notified) {
entry->timeout_notified = 1;
if (!entry->handles) { cm_entry_cleanup(entry); return; }
if (!entry->ever_up && entry->state != CONN_MGR_STATE_CONNECTED) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: all paths unavailable peer=%016llx elapsed=%llums",
(unsigned long long)entry->node_id, (unsigned long long)((now - entry->start_tb) / 10));
cm_deliver_event(entry, CONN_EVENT_TIMEOUT);
}
if (entry->state == CONN_MGR_STATE_DISCONNECTED) return;
if (entry->attempt_active && now >= entry->deadline_tb) {
entry->attempt_active = 0;
if (!entry->retry_delay_ms) entry->retry_delay_ms = CONN_MGR_RETRY_MIN_MS;
entry->retry_tb = now + entry->retry_delay_ms * 10;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: attempt exhausted peer=%016llx path=%s retry=%ums",
(unsigned long long)entry->node_id, cm_path_name(entry->conn_type), entry->retry_delay_ms);
unsigned next = entry->retry_delay_ms * 2;
entry->retry_delay_ms = next > CONN_MGR_RETRY_MAX_MS ? CONN_MGR_RETRY_MAX_MS : next;
}
cm_arm_timer(entry);
}
/* Новый цикл меняет только пробы; NCD, подтверждённый резерв и очереди router остаются живы. */
static void cm_start_attempt(struct CONN_MGR_ENTRY* entry, int recovery) {
struct CONN_MGR* mgr = entry->mgr;
struct TOPO_GROUP_NODE* target = topo_node_find_by_id(mgr->group, entry->node_id);
struct TOPO_NODE* ni = topo_node_registry_find(mgr->instance->topo_groups, entry->node_id);
entry->start_tb = get_time_tb(); entry->attempt_active = 1; entry->retry_tb = 0;
entry->node_timestamp = ni ? ni->timestamp : 0;
entry->reverse_started = 0; entry->reverse_next_tb = 0; entry->indirect_started = 0; entry->indirect_next_tb = 0;
entry->exchange_request_id = 0;
memset(&entry->selection, 0, sizeof(entry->selection));
int local_public = cm_has_direct_ip(mgr, mgr->group->local_node), peer_public = cm_has_direct_ip(mgr, target);
entry->reverse_start_tb = entry->start_tb + (local_public && !peer_public ? 0 : CONN_MGR_REVERSE_DELAY_MS * 10);
entry->indirect_start_tb = entry->start_tb + (recovery || (!local_public && !peer_public) ? 0 : CONN_MGR_INDIRECT_DELAY_MS * 10);
entry->deadline_tb = entry->indirect_start_tb + CONN_MGR_INTERM_EXCHANGE_TIMEOUT_MS * 10;
if (!recovery) entry->initial_deadline_tb = entry->deadline_tb;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: %s group=%016llx peer=%016llx public=%d/%d reverse=%llums indirect=%llums reserve=%d",
recovery ? "recover" : "start", (unsigned long long)mgr->group->group_id, (unsigned long long)entry->node_id,
local_public, peer_public, (unsigned long long)((entry->reverse_start_tb - entry->start_tb) / 10),
(unsigned long long)((entry->indirect_start_tb - entry->start_tb) / 10), entry->indirect_confirmed);
cm_start_phase_direct(entry);
if (entry->state == CONN_MGR_STATE_DISCONNECTED) return;
// После начального бюджета ждём поздний NCD/READY и изменения путей через callbacks.
if (now < entry->deadline_tb) cm_arm_timer(entry);
if (!local_public && !peer_public) cm_start_local_scan(entry);
if (entry->reverse_start_tb == entry->start_tb) cm_start_phase_reverse(entry);
if (!entry->indirect_confirmed && entry->indirect_start_tb == entry->start_tb) cm_start_phase_indirect(entry);
}
/* ═══════ init / destroy ═══════ */
@ -451,15 +536,14 @@ int cm_open(struct CONN_MGR* mgr, uint64_t node_id,
struct CONN_MGR_ENTRY* entry = cm_find_entry(mgr, node_id);
if (entry && entry->state == CONN_MGR_STATE_CONNECTED && entry->handles) {
if (entry && entry->handles) {
struct CONN_MGR_HANDLE* h = cm_handle_new(entry, gid, cb, cb_arg);
if (!h) return -1; if (out_handle) *out_handle = h;
h->deliver_up_token = uasync_call_soon(mgr->instance->ua, h, cm_deliver_up_cb); return 0;
}
if (entry && entry->state == CONN_MGR_STATE_CONNECTING && entry->handles && !entry->timeout_notified) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "node 0x%016llx already connecting, adding handle", (unsigned long long)node_id);
struct CONN_MGR_HANDLE* h = cm_handle_new(entry, gid, cb, cb_arg);
if (!h) return -1; if (out_handle) *out_handle = h; return 0;
if (entry->state == CONN_MGR_STATE_CONNECTED)
h->deliver_up_token = uasync_call_soon(mgr->instance->ua, h, cm_deliver_up_cb);
DEBUG_DEBUG(DEBUG_CATEGORY_GENERAL, "conn_mgr: joined demand peer=%016llx state=%u retry=%llu",
(unsigned long long)node_id, entry->state, (unsigned long long)entry->retry_tb);
return 0;
}
entry = cm_ensure_entry(mgr, node_id);
@ -471,24 +555,12 @@ int cm_open(struct CONN_MGR* mgr, uint64_t node_id,
// Первый владелец входящего пути начинает свои пробы, сохраняя рабочую доставку.
if (entry->timer) { uasync_cancel_timeout(mgr->instance->ua, entry->timer); entry->timer = NULL; }
if (already_up) h->deliver_up_token = uasync_call_soon(mgr->instance->ua, h, cm_deliver_up_cb);
entry->last_traffic_tb = entry->start_tb = get_time_tb();
entry->last_traffic_tb = get_time_tb();
entry->timeout_notified = 0;
entry->reverse_started = 0; entry->reverse_next_tb = 0; entry->indirect_started = 0; entry->indirect_next_tb = 0;
memset(&entry->selection, 0, sizeof(entry->selection));
entry->retry_delay_ms = CONN_MGR_RETRY_MIN_MS;
entry->db_loaded |= (uint8_t)loaded_from_db;
int local_public = cm_has_direct_ip(mgr, group->local_node), peer_public = cm_has_direct_ip(mgr, target);
entry->reverse_start_tb = entry->start_tb + (local_public && !peer_public ? 0 : CONN_MGR_REVERSE_DELAY_MS * 10);
entry->indirect_start_tb = entry->start_tb + (!local_public && !peer_public ? 0 : CONN_MGR_INDIRECT_DELAY_MS * 10);
entry->deadline_tb = entry->indirect_start_tb + CONN_MGR_INTERM_EXCHANGE_TIMEOUT_MS * 10;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: start group=%016llx peer=%016llx public=%d/%d reverse=%llums indirect=%llums",
(unsigned long long)gid, (unsigned long long)node_id, local_public, peer_public,
(unsigned long long)((entry->reverse_start_tb - entry->start_tb) / 10),
(unsigned long long)((entry->indirect_start_tb - entry->start_tb) / 10));
cm_start_phase_direct(entry);
cm_start_attempt(entry, 0);
if (entry->state == CONN_MGR_STATE_DISCONNECTED) return 0;
if (!local_public && !peer_public) cm_start_local_scan(entry);
if (entry->reverse_start_tb == entry->start_tb) cm_start_phase_reverse(entry);
if (!entry->indirect_confirmed && entry->indirect_start_tb == entry->start_tb) cm_start_phase_indirect(entry);
cm_arm_timer(entry);
return 0;
}
@ -589,7 +661,7 @@ void cm_handle_direct_req(uint64_t src, struct CONN_MGR* mgr, const uint8_t* dat
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: malformed DIRECT_REQ peer=%016llx len=%zu", (unsigned long long)src, len);
return;
}
(void)data;
struct CM_DIRECT_REQ req; memcpy(&req, data, sizeof(req));
struct TOPO_NODE* ni = topo_node_registry_find(mgr->instance->topo_groups, src);
if (!ni) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: REVERSE unknown peer=%016llx", (unsigned long long)src); return; }
struct CONN_MGR_ENTRY* entry = cm_ensure_entry(mgr, src);
@ -598,6 +670,9 @@ void cm_handle_direct_req(uint64_t src, struct CONN_MGR* mgr, const uint8_t* dat
entry->state = CONN_MGR_STATE_CONNECTING; entry->start_tb = get_time_tb();
entry->deadline_tb = entry->start_tb + CONN_MGR_INTERM_EXCHANGE_TIMEOUT_MS * 10;
}
if (entry->peer_reverse_request_id == req.request_id) return;
entry->peer_reverse_request_id = req.request_id;
if (!entry->handles) entry->deadline_tb = get_time_tb() + CONN_MGR_INTERM_EXCHANGE_TIMEOUT_MS * 10;
if (entry->direct_up && cm_direct_conn(mgr, src)) return;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: REVERSE request peer=%016llx path=%u existing_transport=%d",
(unsigned long long)src, entry->conn_type, entry->ncd_handle != NULL);

34
src/routing_layer/conn_mgr_indirect.c

@ -47,6 +47,7 @@ void cm_finalize_indirect(struct CONN_MGR_ENTRY* entry) {
uint8_t old_type = entry->conn_type;
cm_refresh_path(entry); // Готовый прямой путь имеет приоритет; повторного UP при улучшении нет.
if (entry->state != CONN_MGR_STATE_DISCONNECTED && entry->conn_type == old_type) cm_update_nodeinfo(entry);
if (entry->state != CONN_MGR_STATE_DISCONNECTED) cm_arm_timer(entry);
}
/* Объединение кандидатов обеих сторон без дублей, выбор по сумме двух действительных RTT. */
@ -95,7 +96,7 @@ void cm_handle_interm_exchange_resp(struct CONN_MGR* mgr, uint64_t src, const ui
}
struct CM_EXCHANGE_RESP resp; memcpy(&resp, data, sizeof(resp));
struct CONN_MGR_ENTRY* entry = cm_find_entry(mgr, src);
if (!entry || !entry->indirect_started || entry->indirect_confirmed || entry->selection.count ||
if (!entry || !entry->attempt_active || !entry->indirect_started || entry->indirect_confirmed || entry->selection.count ||
entry->exchange_request_id != resp.request_id || get_time_tb() >= entry->deadline_tb) {
DEBUG_DEBUG(DEBUG_CATEGORY_GENERAL, "conn_mgr: stale EXCHANGE_RESP peer=%016llx request=%u", (unsigned long long)src, resp.request_id);
return;
@ -121,7 +122,7 @@ void cm_handle_interm_selected(struct CONN_MGR* mgr, uint64_t src, const uint8_t
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: stale selection peer=%016llx request=%u", (unsigned long long)src, sel.request_id);
return;
}
if (get_time_tb() >= entry->deadline_tb && !entry->indirect_confirmed) {
if (get_time_tb() >= entry->peer_exchange_deadline_tb) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: selection after deadline peer=%016llx", (unsigned long long)src);
return;
}
@ -150,7 +151,8 @@ void cm_handle_interm_ack(struct CONN_MGR* mgr, uint64_t src, const uint8_t* dat
}
struct CM_INTERM_SEL ack; memcpy(&ack, data, sizeof(ack));
struct CONN_MGR_ENTRY* entry = cm_find_entry(mgr, src);
if (!entry || entry->exchange_request_id != ack.request_id || !entry->selection.count || entry->indirect_confirmed ||
if (!entry || !entry->attempt_active || entry->exchange_request_id != ack.request_id ||
!entry->selection.count || entry->indirect_confirmed ||
get_time_tb() >= entry->deadline_tb) {
DEBUG_DEBUG(DEBUG_CATEGORY_GENERAL, "conn_mgr: stale INTERM_ACK peer=%016llx request=%u", (unsigned long long)src, ack.request_id);
return;
@ -158,23 +160,24 @@ void cm_handle_interm_ack(struct CONN_MGR* mgr, uint64_t src, const uint8_t* dat
if (ack.count > CONN_MGR_MAX_INTERMEDIARIES) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: invalid ACK count=%u peer=%016llx", ack.count, (unsigned long long)src); return;
}
entry->intermediariy_count = 0;
uint64_t confirmed[CONN_MGR_MAX_INTERMEDIARIES]; uint8_t count = 0;
for (uint8_t i = 0; i < ack.count; i++) {
uint64_t nid = ack.selected[i].node_id;
if (!cm_valid_candidate(mgr, src, nid, ack.selected[i].rtt)) continue;
for (uint8_t j = 0; j < entry->selection.count; j++) {
if (entry->selection.selected[j].node_id != nid) continue;
int duplicate = 0;
for (uint8_t k = 0; k < entry->intermediariy_count; k++) if (entry->intermediaries[k] == nid) duplicate = 1;
if (!duplicate) entry->intermediaries[entry->intermediariy_count++] = nid;
for (uint8_t k = 0; k < count; k++) if (confirmed[k] == nid) duplicate = 1;
if (!duplicate) confirmed[count++] = nid;
break;
}
}
if (!entry->intermediariy_count) {
if (!count) {
memset(&entry->selection, 0, sizeof(entry->selection));
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "conn_mgr: intermediary verification failed peer=%016llx; repeating exchange", (unsigned long long)src);
return;
}
memcpy(entry->intermediaries, confirmed, count * sizeof(*confirmed)); entry->intermediariy_count = count;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: INDIRECT confirmed peer=%016llx count=%u elapsed=%llums",
(unsigned long long)src, entry->intermediariy_count, (unsigned long long)((get_time_tb() - entry->start_tb) / 10));
cm_finalize_indirect(entry);
@ -188,12 +191,21 @@ void cm_handle_interm_exchange_req(uint64_t src, struct CONN_MGR* mgr, struct CM
}
struct CONN_MGR_ENTRY* entry = cm_ensure_entry(mgr, src);
if (!entry) return;
entry->peer_exchange_request_id = req->request_id;
uint64_t now = get_time_tb();
if (entry->peer_exchange_request_id != req->request_id || !entry->peer_exchange_deadline_tb) {
entry->peer_exchange_request_id = req->request_id;
entry->peer_exchange_deadline_tb = now + CONN_MGR_INTERM_EXCHANGE_TIMEOUT_MS * 10;
if (!entry->handles) entry->deadline_tb = entry->peer_exchange_deadline_tb;
}
if (now >= entry->peer_exchange_deadline_tb) {
DEBUG_DEBUG(DEBUG_CATEGORY_GENERAL, "conn_mgr: expired EXCHANGE_REQ peer=%016llx request=%u",
(unsigned long long)src, req->request_id);
return;
}
if (entry->state == CONN_MGR_STATE_DISCONNECTED) {
entry->state = CONN_MGR_STATE_CONNECTING; entry->start_tb = get_time_tb();
entry->deadline_tb = entry->start_tb + CONN_MGR_INTERM_EXCHANGE_TIMEOUT_MS * 10;
cm_arm_timer(entry);
entry->state = CONN_MGR_STATE_CONNECTING; entry->start_tb = now;
}
cm_arm_timer(entry);
struct CM_EXCHANGE_RESP resp = {0};
resp.cmd = ETCP_RT_ID_CONN_MGR; resp.subcmd = CM_SUBCMD_INTERM_EXCHANGE_RESP; resp.request_id = req->request_id;
for (uint8_t i = 0; i < mgr->best_candidate_count; i++) {

7
src/routing_layer/conn_mgr_priv.h

@ -32,6 +32,8 @@ struct stcp_client;
#define CONN_MGR_INTERM_EXCHANGE_TIMEOUT_MS 15000
#define CONN_MGR_REVERSE_RETRANSMIT_MS 1000
#define CONN_MGR_EXCHANGE_RETRANSMIT_MS 1000
#define CONN_MGR_RETRY_MIN_MS 1000
#define CONN_MGR_RETRY_MAX_MS 15000
#define CACHE_MS 20000 /* время жизни RTT-замера ~20с (через cm_is_rtt_fresh, мс) */
#define CONN_MGR_MAX_INTERMEDIARIES 3
@ -100,10 +102,13 @@ struct CONN_MGR_ENTRY {
uint64_t intermediaries[3]; uint8_t intermediariy_count;
uint8_t local_scan_state, main_connect_state, db_loaded;
uint8_t reverse_started, indirect_started, indirect_confirmed, direct_up, timeout_notified, ever_up;
uint8_t attempt_active;
uint16_t retry_delay_ms;
void* timer; void* refresh_token;
uint64_t start_tb, reverse_start_tb, indirect_start_tb, deadline_tb;
uint64_t initial_deadline_tb, retry_tb, node_timestamp, peer_exchange_deadline_tb;
uint64_t reverse_next_tb, reverse_deadline_tb, indirect_next_tb;
uint32_t reverse_request_id, exchange_request_id, peer_exchange_request_id;
uint32_t reverse_request_id, exchange_request_id, peer_exchange_request_id, peer_reverse_request_id;
struct CM_INTERM_SEL selection; /* Исходящий выбор до подтверждения пира. */
struct CONN_MGR* mgr;
};

Loading…
Cancel
Save