diff --git a/src/routing_layer/conn_mgr_core.c b/src/routing_layer/conn_mgr_core.c index 6af20d37..737f775f 100644 --- a/src/routing_layer/conn_mgr_core.c +++ b/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); diff --git a/src/routing_layer/conn_mgr_indirect.c b/src/routing_layer/conn_mgr_indirect.c index 723e10e4..35ae6305 100644 --- a/src/routing_layer/conn_mgr_indirect.c +++ b/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++) { diff --git a/src/routing_layer/conn_mgr_priv.h b/src/routing_layer/conn_mgr_priv.h index cc6aa054..67adcd9c 100644 --- a/src/routing_layer/conn_mgr_priv.h +++ b/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; };