Browse Source

conn_mgr: parallel reverse+indirect, retransmit, and call teardown fix

- call_on_conn_status: only end call when its own data conn drops
  (INDIRECT call was killed by unrelated topo_group_connect conn DELETE)
- conn_mgr: run REVERSE and INDIRECT in parallel after DIRECT fails;
  reverse is preferred (indirect defers via indirect_ready until reverse
  resolves). REVERSE timeout 15s -> 3s.
- retransmit DIRECT_REQ / EXCHANGE_REQ every 1s (idempotent, same request_id)
  to survive packet loss (router RST dropping stale-channel requests)
- replace main.phase with reverse_active/indirect_active/indirect_ready flags
- add diagnostics to trace the parallel flow
proxy
evgeny 1 week ago
parent
commit
570b421ba0
  1. 7
      src/call/call.c
  2. 150
      src/routing_layer/conn_mgr_core.c
  3. 83
      src/routing_layer/conn_mgr_indirect.c
  4. 17
      src/routing_layer/conn_mgr_priv.h

7
src/call/call.c

@ -55,6 +55,7 @@ struct call_session {
uint8_t reason; /* CALL_REASON_* (при ENDED) */
struct call_ctx* ctx; /* обратная ссылка (для коллбэков conn_mgr/таймеров) */
struct CONN_MGR_HANDLE* cm_handle;
struct ETCP_CONN* data_conn; /* прямой data-коннект звонка (NULL для INDIRECT) */
void* ring_timer; /* caller: 45с ожидания ответа */
void* traffic_timer; /* ACTIVE: 20с без медиа (запускается по первому медиа) */
void* teardown_timer; /* отложенный разрыв conn_mgr (500мс после завершения) */
@ -323,6 +324,7 @@ static void call_conn_cb(struct CONN_MGR_HANDLE* h, uint64_t node_id, uint64_t g
if (event == CONN_EVENT_UP) {
s->cm_handle = h;
s->data_conn = conn_mgr_get_conn(h); /* NULL для INDIRECT */
s->state = CALL_ST_ACTIVE;
DEBUG_INFO(DEBUG_CATEGORY_CALL, "%s: conn UP id=%016llx peer=0x%016llx -> ACTIVE",
CALL_ID, (unsigned long long)s->call_id, (unsigned long long)node_id);
@ -432,9 +434,14 @@ static void call_on_conn_status(struct ETCP_CONN* conn, int status, void* arg) {
struct call_session* s = (struct call_session*)e->data;
if (s->peer_node_id == peer) {
if (s->state == CALL_ST_ACTIVE) {
/* Завершаем только если упал именно data-коннект звонка (DIRECT/REVERSE).
* Для INDIRECT data_conn==NULL — посторонний коннект к тому же пиру
* (topo_group_connect) не должен рвать работающий звонок. */
if (s->data_conn && s->data_conn == conn) {
DEBUG_WARN(DEBUG_CATEGORY_CALL, "%s: peer conn lost id=%016llx peer=0x%016llx -> end",
CALL_ID, (unsigned long long)s->call_id, (unsigned long long)peer);
call_end(ctx, s, CALL_REASON_REMOTE_HANGUP);
}
} else if (s->state == CALL_ST_ENDED) {
/* удалённая сторона уже порвала (DISCONNECT чистит entry без коллбэка) */
s->cm_handle = NULL;

150
src/routing_layer/conn_mgr_core.c

@ -30,6 +30,7 @@
/* ═══════ forward-декларации ═══════ */
void conn_mgr_router_recv_handler(struct ETCP_CONN* conn, struct ll_entry* entry);
static void cm_start_parallel(struct CONN_MGR_ENTRY* entry);
/* ═══════ утилиты ═══════ */
@ -181,9 +182,11 @@ struct CONN_MGR_ENTRY* cm_ensure_entry(struct CONN_MGR* mgr, uint64_t node_id) {
void cm_entry_cleanup(struct CONN_MGR_ENTRY* entry) {
if (!entry || entry->state == CONN_MGR_STATE_DISCONNECTED) return;
if (entry->main.timer) { uasync_cancel_timeout(entry->mgr->instance->ua, entry->main.timer); entry->main.timer = NULL; }
if (entry->reverse_timer) { uasync_cancel_timeout(entry->mgr->instance->ua, entry->reverse_timer); entry->reverse_timer = NULL; }
entry->reverse_active = 0; entry->indirect_active = 0; entry->indirect_ready = 0;
{ struct cm_reverse_pending** pp = &entry->mgr->reverse_pending;
while (*pp) {
if ((*pp)->entry == entry) { struct cm_reverse_pending* rp = *pp; *pp = rp->next; u_free(rp); }
if ((*pp)->entry == entry) { struct cm_reverse_pending* rp = *pp; *pp = rp->next; if (rp->pkt) u_free(rp->pkt); u_free(rp); }
else pp = &(*pp)->next;
}
}
@ -192,7 +195,7 @@ void cm_entry_cleanup(struct CONN_MGR_ENTRY* entry) {
if ((*pp)->entry == entry) {
if ((*pp)->timer && (*pp)->timer != entry->main.timer)
{ uasync_cancel_timeout(entry->mgr->instance->ua, (*pp)->timer); (*pp)->timer = NULL; }
struct cm_exchange_pending* ep = *pp; *pp = ep->next; u_free(ep);
struct cm_exchange_pending* ep = *pp; *pp = ep->next; if (ep->pkt) u_free(ep->pkt); u_free(ep);
} else pp = &(*pp)->next;
}
}
@ -251,14 +254,22 @@ void cm_ncd_callback(struct NODE_CONN_DIRECT* ncd_h, enum ncd_event ncd_ev, void
switch (ncd_ev) {
case NCD_EVENT_UP:
if (entry->state == CONN_MGR_STATE_CONNECTED) return;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "link: ETCP handshake OK with 0x%016llx — connection ESTABLISHED",
(unsigned long long)entry->node_id);
{ int won = entry->reverse_active ? CONN_TYPE_REVERSE : CONN_TYPE_DIRECT;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "link: ETCP handshake OK 0x%016llx — %s (reverse_active=%d indirect_active=%d)",
(unsigned long long)entry->node_id,
won == CONN_TYPE_REVERSE ? "REVERSE WON (indirect cancelled)" : "DIRECT",
entry->reverse_active, entry->indirect_active); }
if (entry->main.timer) { uasync_cancel_timeout(entry->mgr->instance->ua, entry->main.timer); entry->main.timer = NULL; }
if (entry->reverse_timer) { uasync_cancel_timeout(entry->mgr->instance->ua, entry->reverse_timer); entry->reverse_timer = NULL; }
entry->main_connect_state = CM_TRY_OK;
entry->conn_type = (entry->main.phase == 2) ? CONN_TYPE_REVERSE : CONN_TYPE_DIRECT;
entry->conn_type = entry->reverse_active ? CONN_TYPE_REVERSE : CONN_TYPE_DIRECT;
entry->state = CONN_MGR_STATE_CONNECTED;
entry->reverse_active = 0; entry->indirect_active = 0; entry->indirect_ready = 0;
{ struct cm_reverse_pending** pp = &entry->mgr->reverse_pending;
while (*pp) { if ((*pp)->entry == entry) { struct cm_reverse_pending* rp = *pp; *pp = rp->next; u_free(rp); break; } pp = &(*pp)->next; }
while (*pp) { if ((*pp)->entry == entry) { struct cm_reverse_pending* rp = *pp; *pp = rp->next; if (rp->pkt) u_free(rp->pkt); u_free(rp); break; } pp = &(*pp)->next; }
}
{ struct cm_exchange_pending** pp = &entry->mgr->exchange_pending;
while (*pp) { if ((*pp)->entry == entry) { struct cm_exchange_pending* ep = *pp; *pp = ep->next; if (ep->pkt) u_free(ep->pkt); u_free(ep); break; } pp = &(*pp)->next; }
}
cm_update_nodeinfo(entry);
{ struct ETCP_CONN* c = node_conn_direct_get_conn(ncd_h);
@ -287,19 +298,12 @@ void cm_ncd_callback(struct NODE_CONN_DIRECT* ncd_h, enum ncd_event ncd_ev, void
cm_entry_cleanup(entry);
return;
}
/* Инициатор с handle'ами: выбор следующей фазы делаем только из DIRECT (phase==0).
* В REVERSE (phase==2) эскалацией владеет собственный 15с таймер, а повторные
/* Инициатор с handle'ами: параллельный запуск REVERSE + INDIRECT. Повторные
* NCD-таймауты (шаримый NCD оживляется другими модулями) игнорируем — иначе
* cm_start_phase_reverse утечёт старыми таймерами и выстрелит в освобождённый entry. */
if (entry->main.phase == 2) return;
{ struct TOPO_GROUP* g = entry->mgr->group;
struct TOPO_GROUP_NODE* t = topo_node_find_by_id(g, entry->node_id);
if (cm_has_direct_ip(entry->mgr, g->local_node) && !(t && cm_has_direct_ip(entry->mgr, t)))
{ cm_start_phase_reverse(entry); return; } /* keep ncd_handle: target will reuse this conn for incoming INIT */
{ struct NODE_CONN_DIRECT* ncd = entry->ncd_handle; entry->ncd_handle = NULL;
if (ncd) node_conn_direct_force_close(ncd); }
cm_start_phase_indirect(entry); return;
}
* протекут таймеры и выстрелят в освобождённый entry. */
if (entry->reverse_active || entry->indirect_active) return;
cm_start_parallel(entry);
break;
case NCD_EVENT_DOWN:
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "link: ETCP link DROPPED to 0x%016llx",
(unsigned long long)entry->node_id);
@ -624,11 +628,7 @@ static void cm_direct_add_links(struct ETCP_CONN* conn, struct CONN_MGR_ENTRY* e
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: db_node 0x%016llx direct failed (no compatible addr/socket)", (unsigned long long)entry->node_id);
cm_cleanup_db_node(entry); cm_deliver_event(entry, CONN_EVENT_TIMEOUT);
} else {
struct TOPO_GROUP* g = entry->mgr->group;
struct TOPO_GROUP_NODE* t = topo_node_find_by_id(g, entry->node_id);
if (cm_has_direct_ip(entry->mgr, g->local_node) && !(t && cm_has_direct_ip(entry->mgr, t)))
{ cm_start_phase_reverse(entry); return; }
cm_start_phase_indirect(entry);
cm_start_parallel(entry);
}
}
@ -652,27 +652,33 @@ void cm_start_phase_direct(struct CONN_MGR_ENTRY* entry) {
cm_direct_add_links(node_conn_direct_get_conn(entry->ncd_handle), entry, ni, entry->db_loaded);
}
/* ═══════ REVERSE фаза ═══════ */
/* ═══════ REVERSE + INDIRECT (параллельный запуск) ═══════ */
void cm_reverse_timeout_cb(void* arg) {
struct CONN_MGR_ENTRY* entry = (struct CONN_MGR_ENTRY*)arg; entry->main.timer = NULL;
{ struct cm_reverse_pending** pp = &entry->mgr->reverse_pending;
while (*pp) { if ((*pp)->entry == entry) { struct cm_reverse_pending* rp = *pp; *pp = rp->next; u_free(rp); break; } pp = &(*pp)->next; }
}
{ struct NODE_CONN_DIRECT* ncd = entry->ncd_handle; entry->ncd_handle = NULL;
if (ncd) node_conn_direct_force_close(ncd); }
cm_start_phase_indirect(entry);
}
/* REVERSE фаза: у нас прямой IP, у цели нет. Отправляем DIRECT_REQ с НАШИМИ
* адресами цели через BGP-маршрут ("подключись ко мне"). Ждём обратного INIT.
* Таймаут 15с — при провале переходим к INDIRECT. */
/* После провала DIRECT запускаем ОБА пути параллельно: INDIRECT (гарантированный
* fallback через релей) и REVERSE (оптимизация — прямой коннект, если пир дозвонится).
* Reverse строго приоритетнее: его INIT побеждает, indirect откладывается (indirect_ready). */
static void cm_start_parallel(struct CONN_MGR_ENTRY* entry) {
struct TOPO_GROUP* g = entry->mgr->group;
struct TOPO_GROUP_NODE* t = topo_node_find_by_id(g, entry->node_id);
int can_reverse = cm_has_direct_ip(entry->mgr, g->local_node) && !(t && cm_has_direct_ip(entry->mgr, t));
int can_indirect = entry->mgr->best_candidate_count > 0;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: direct failed 0x%016llx — parallel start reverse=%d indirect=%d (candidates=%u)",
(unsigned long long)entry->node_id, can_reverse, can_indirect, entry->mgr->best_candidate_count);
if (can_indirect) cm_start_phase_indirect(entry);
if (can_reverse) cm_start_phase_reverse(entry); /* keep ncd_handle: target will reuse this conn for incoming INIT */
else { struct NODE_CONN_DIRECT* ncd = entry->ncd_handle; entry->ncd_handle = NULL; if (ncd) node_conn_direct_force_close(ncd); }
if (!entry->reverse_active && !entry->indirect_active) cm_deliver_event(entry, CONN_EVENT_TIMEOUT);
}
/* REVERSE фаза: у нас прямой IP, у цели нет. Отправляем DIRECT_REQ с НАШИМИ адресами
* цели через BGP-маршрут ("подключись ко мне"). Ждём обратного INIT (дедлайн 3с),
* периодически ретрансмитим DIRECT_REQ (1с) — надёжность на случай потери пакета. */
void cm_start_phase_reverse(struct CONN_MGR_ENTRY* entry) {
struct TOPO_GROUP* group = entry->mgr->group;
struct TOPO_GROUP_NODE* local = group->local_node;
struct TOPO_NODE* lni = local ? topo_node_registry_find(entry->mgr->instance->topo_groups, local->node_id) : NULL;
if (!local || !lni) { cm_deliver_event(entry, CONN_EVENT_TIMEOUT); return; }
uint32_t req_id = ++entry->mgr->next_request_id; entry->main.request_id = req_id; entry->main.phase = 2;
if (!local || !lni) return;
uint32_t req_id = ++entry->mgr->next_request_id; entry->main.request_id = req_id;
uint8_t dc = 0; struct { uint8_t t; uint8_t ip[4]; uint16_t port; uint8_t sid; } addrs[8];
for (const struct TOPO_ADDR4* a = lni->v4_addrs; a && dc < 8; a = a->next) {
if (a->type != TOPO_ADDR_NAT && a->type != TOPO_ADDR_INTERFACE) continue;
@ -682,26 +688,72 @@ void cm_start_phase_reverse(struct CONN_MGR_ENTRY* entry) {
{ addrs[dc]=(typeof(addrs[0])){a->type}; memcpy(addrs[dc].ip,a->addr,4); addrs[dc].port=htons(a->port); addrs[dc].sid=a->socket_id; dc++; break; }
}
}
if (!dc) { cm_start_phase_indirect(entry); return; }
if (!dc) return;
size_t sz = CM_DIRECT_REQ_H_SIZE + (size_t)dc * 8; uint8_t* pkt = u_malloc(sz);
if (!pkt) { cm_deliver_event(entry, CONN_EVENT_TIMEOUT); return; }
if (!pkt) return;
struct CM_DIRECT_REQ* req = (struct CM_DIRECT_REQ*)pkt; memset(req,0,sizeof(*req));
req->cmd=ETCP_RT_ID_CONN_MGR; req->subcmd=CM_SUBCMD_DIRECT_REQ; req->request_id=req_id; req->addr_count=dc;
uint8_t* p = pkt + CM_DIRECT_REQ_H_SIZE;
for (uint8_t i=0;i<dc;i++) { *p++=addrs[i].t; memcpy(p,addrs[i].ip,4);p+=4; memcpy(p,&addrs[i].port,2);p+=2; *p++=addrs[i].sid; }
struct ll_entry* qe = queue_entry_new(0);
if (!qe) { u_free(pkt); cm_deliver_event(entry, CONN_EVENT_TIMEOUT); return; }
qe->dgram=u_malloc(1+sz); qe->dgram[0]=ETCP_RT_ID_CONN_MGR; memcpy(qe->dgram+1,pkt,sz); u_free(pkt); qe->len=(uint16_t)(1+sz);
etcp_route_send(entry->mgr->instance, group->group_id, entry->node_id, qe, 1, 0);
if (!qe) { u_free(pkt); return; }
qe->dgram=u_malloc(1+sz);
if (!qe->dgram) { queue_entry_free(qe); u_free(pkt); return; }
qe->dgram[0]=ETCP_RT_ID_CONN_MGR; memcpy(qe->dgram+1,pkt,sz); qe->len=(uint16_t)(1+sz);
/* копия для ретрансмита (etcp_route_send забирает владение qe) */
struct cm_reverse_pending* rp = u_calloc(1, sizeof(*rp));
if (rp) { rp->request_id=req_id; rp->entry=entry; rp->next=entry->mgr->reverse_pending; entry->mgr->reverse_pending=rp; }
if (entry->main.timer) { uasync_cancel_timeout(entry->mgr->instance->ua, entry->main.timer); entry->main.timer = NULL; }
entry->main.timer = uasync_set_timeout(entry->mgr->instance->ua, CONN_MGR_CONNECT_REVERSE_TIMEOUT_MS*10,
entry, cm_reverse_timeout_cb, "conn_mgr_reverse");
if (rp) {
rp->request_id=req_id; rp->entry=entry;
rp->pkt=u_malloc(1+sz);
if (rp->pkt) { memcpy(rp->pkt, qe->dgram, 1+sz); rp->pkt_len=1+sz; }
rp->next=entry->mgr->reverse_pending; entry->mgr->reverse_pending=rp;
}
etcp_route_send(entry->mgr->instance, group->group_id, entry->node_id, qe, 1, 0);
u_free(pkt);
entry->reverse_active = 1;
entry->reverse_deadline_tb = get_time_tb() + CONN_MGR_CONNECT_REVERSE_TIMEOUT_MS*10;
entry->reverse_timer = uasync_set_timeout(entry->mgr->instance->ua, CONN_MGR_REVERSE_RETRANSMIT_MS*10,
entry, cm_reverse_tick_cb, "conn_mgr_reverse");
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: phase 2 (reverse) sent DIRECT_REQ to 0x%016llx with %u addrs",
(unsigned long long)entry->node_id, dc);
}
/* Периодический тик REVERSE: ретрансмит DIRECT_REQ, по дедлайну — сдаёмся, закрываем
* NCD; если indirect уже завершил обмен (indirect_ready) — финализируем INDIRECT. */
void cm_reverse_tick_cb(void* arg) {
struct CONN_MGR_ENTRY* entry = (struct CONN_MGR_ENTRY*)arg;
entry->reverse_timer = NULL;
if (!entry->reverse_active) return;
if (get_time_tb() >= entry->reverse_deadline_tb) {
entry->reverse_active = 0;
{ struct cm_reverse_pending** pp = &entry->mgr->reverse_pending;
while (*pp) { if ((*pp)->entry == entry) { struct cm_reverse_pending* rp = *pp; *pp = rp->next; if (rp->pkt) u_free(rp->pkt); u_free(rp); break; } pp = &(*pp)->next; }
}
{ struct NODE_CONN_DIRECT* ncd = entry->ncd_handle; entry->ncd_handle = NULL;
if (ncd) node_conn_direct_force_close(ncd); }
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: reverse TIMEOUT 0x%016llx — indirect_ready=%d indirect_active=%d",
(unsigned long long)entry->node_id, entry->indirect_ready, entry->indirect_active);
if (entry->indirect_ready) cm_finalize_indirect(entry);
else if (!entry->indirect_active) cm_deliver_event(entry, CONN_EVENT_TIMEOUT);
return;
}
for (struct cm_reverse_pending* x = entry->mgr->reverse_pending; x; x = x->next) {
if (x->entry == entry && x->pkt) {
struct ll_entry* qe = queue_entry_new(0);
if (qe) {
qe->dgram = u_malloc(x->pkt_len);
if (qe->dgram) { memcpy(qe->dgram, x->pkt, x->pkt_len); qe->len = (uint16_t)x->pkt_len;
etcp_route_send(entry->mgr->instance, entry->mgr->group->group_id, entry->node_id, qe, 1, 0); }
else queue_entry_free(qe);
}
DEBUG_DEBUG(DEBUG_CATEGORY_GENERAL, "conn_mgr: reverse retransmit DIRECT_REQ 0x%016llx", (unsigned long long)entry->node_id);
break;
}
}
entry->reverse_timer = uasync_set_timeout(entry->mgr->instance->ua, CONN_MGR_REVERSE_RETRANSMIT_MS*10,
entry, cm_reverse_tick_cb, "conn_mgr_reverse");
}
/* ═══════ обработка DIRECT_REQ (REVERSE входящие) ═══════ */
/* Принимающая сторона REVERSE: получили DIRECT_REQ от инициатора — открываем
@ -725,7 +777,7 @@ void cm_handle_direct_req(uint64_t src, struct CONN_MGR* mgr, const uint8_t* dat
if (entry && entry->state == CONN_MGR_STATE_CONNECTING) return;
if (!entry) entry = cm_ensure_entry(mgr, src);
if (!entry) return;
entry->state = CONN_MGR_STATE_CONNECTING; entry->main.phase = 2;
entry->state = CONN_MGR_STATE_CONNECTING; entry->reverse_active = 1;
entry->main_connect_state = CM_TRY_PENDING; entry->db_loaded = 0; entry->conn_type = CONN_TYPE_NONE;
struct TOPO_NODE nbuf; memset(&nbuf,0,sizeof(nbuf)); nbuf.node_id=src; memcpy(nbuf.public_key,src_ni->public_key,SC_PUBKEY_SIZE);

83
src/routing_layer/conn_mgr_indirect.c

@ -12,35 +12,75 @@
/* INDIRECT фаза: отправляет EXCHANGE_REQ с нашими лучшими кандидатами (best_candidates)
* цели через BGP-маршрут. Ожидаем EXCHANGE_RESP — по RTT-замерам с обеих сторон
* выберем общего посредника с минимальной суммой RTT (cm_compute_intermediaries). */
* выберем общего посредника с минимальной суммой RTT (cm_compute_intermediaries).
* Периодически ретрансмитим EXCHANGE_REQ (1с) — надёжность на случай потери пакета. */
void cm_start_phase_indirect(struct CONN_MGR_ENTRY* entry) {
struct CONN_MGR* mgr = entry->mgr;
if (!mgr->best_candidate_count) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "conn_mgr: phase 3 (indirect) no candidates for 0x%016llx", (unsigned long long)entry->node_id);
cm_deliver_event(entry, CONN_EVENT_TIMEOUT); return;
}
entry->main.phase = 3; entry->main.request_id = ++mgr->next_request_id;
if (!mgr->best_candidate_count) return;
entry->main.request_id = ++mgr->next_request_id;
struct CM_EXCHANGE_REQ req; memset(&req,0,sizeof(req));
req.cmd=ETCP_RT_ID_CONN_MGR; req.subcmd=CM_SUBCMD_INTERM_EXCHANGE_REQ; req.request_id=entry->main.request_id;
req.candidate_count = mgr->best_candidate_count > 4 ? 4 : mgr->best_candidate_count;
for (uint8_t i=0;i<req.candidate_count;i++) req.candidates[i]=mgr->best_candidates[i];
struct ll_entry* qe = queue_entry_new(0); if (!qe) return;
qe->dgram=u_malloc(1+sizeof(req)); qe->dgram[0]=ETCP_RT_ID_CONN_MGR; memcpy(qe->dgram+1,&req,sizeof(req)); qe->len=(uint16_t)(1+sizeof(req));
/* копия для ретрансмита (etcp_route_send забирает владение qe) */
struct cm_exchange_pending* ep = u_calloc(1,sizeof(*ep));
if (ep) {
ep->request_id=entry->main.request_id; ep->entry=entry;
ep->pkt=u_malloc(1+sizeof(req));
if (ep->pkt) { memcpy(ep->pkt, qe->dgram, 1+sizeof(req)); ep->pkt_len=1+sizeof(req); }
ep->next=mgr->exchange_pending; mgr->exchange_pending=ep;
}
etcp_route_send(mgr->instance, mgr->group->group_id, entry->node_id, qe, 1, 0);
entry->indirect_active = 1;
entry->indirect_deadline_tb = get_time_tb() + CONN_MGR_INTERM_EXCHANGE_TIMEOUT_MS*10;
if (entry->main.timer) { uasync_cancel_timeout(mgr->instance->ua, entry->main.timer); entry->main.timer = NULL; }
entry->main.timer = uasync_set_timeout(mgr->instance->ua, CONN_MGR_INTERM_EXCHANGE_TIMEOUT_MS*10, entry, cm_exchange_timeout_cb, "conn_mgr_interm_exch");
struct cm_exchange_pending* ep = u_calloc(1,sizeof(*ep));
if (ep) { ep->request_id=entry->main.request_id; ep->entry=entry; ep->timer=entry->main.timer; ep->next=mgr->exchange_pending; mgr->exchange_pending=ep; }
entry->main.timer = uasync_set_timeout(mgr->instance->ua, CONN_MGR_EXCHANGE_RETRANSMIT_MS*10, entry, cm_exchange_tick_cb, "conn_mgr_interm_exch");
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: phase 3 sent EXCHANGE_REQ to 0x%016llx with %u candidates",
(unsigned long long)entry->node_id, req.candidate_count);
}
void cm_exchange_timeout_cb(void* arg) {
struct CONN_MGR_ENTRY* entry = (struct CONN_MGR_ENTRY*)arg; entry->main.timer=NULL;
/* Периодический тик INDIRECT: ретрансмит EXCHANGE_REQ, по дедлайну — таймаут обмена. */
void cm_exchange_tick_cb(void* arg) {
struct CONN_MGR_ENTRY* entry = (struct CONN_MGR_ENTRY*)arg;
entry->main.timer = NULL;
if (!entry->indirect_active) return;
if (get_time_tb() >= entry->indirect_deadline_tb) {
entry->indirect_active = 0;
{ struct cm_exchange_pending** pp=&entry->mgr->exchange_pending;
while(*pp){if((*pp)->entry==entry){struct cm_exchange_pending* ep=*pp; *pp=ep->next; u_free(ep); break;} pp=&(*pp)->next;} }
while(*pp){if((*pp)->entry==entry){struct cm_exchange_pending* ep=*pp; *pp=ep->next; if(ep->pkt)u_free(ep->pkt); u_free(ep); break;} pp=&(*pp)->next;} }
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: exchange timeout for 0x%016llx", (unsigned long long)entry->node_id);
cm_deliver_event(entry, CONN_EVENT_TIMEOUT);
if (!entry->reverse_active) cm_deliver_event(entry, CONN_EVENT_TIMEOUT);
return;
}
for (struct cm_exchange_pending* x = entry->mgr->exchange_pending; x; x = x->next) {
if (x->entry == entry && x->pkt) {
struct ll_entry* qe = queue_entry_new(0);
if (qe) {
qe->dgram = u_malloc(x->pkt_len);
if (qe->dgram) { memcpy(qe->dgram, x->pkt, x->pkt_len); qe->len = (uint16_t)x->pkt_len;
etcp_route_send(entry->mgr->instance, entry->mgr->group->group_id, entry->node_id, qe, 1, 0); }
else queue_entry_free(qe);
}
DEBUG_DEBUG(DEBUG_CATEGORY_GENERAL, "conn_mgr: exchange retransmit EXCHANGE_REQ 0x%016llx", (unsigned long long)entry->node_id);
break;
}
}
entry->main.timer = uasync_set_timeout(entry->mgr->instance->ua, CONN_MGR_EXCHANGE_RETRANSMIT_MS*10,
entry, cm_exchange_tick_cb, "conn_mgr_interm_exch");
}
/* Финализация INDIRECT: пометить CONNECTED и доставить UP. Вызывается либо сразу
* (reverse не активен), либо из cm_reverse_tick_cb после reverse-таймаута. */
void cm_finalize_indirect(struct CONN_MGR_ENTRY* entry) {
entry->conn_type = CONN_TYPE_INDIRECT;
entry->state = CONN_MGR_STATE_CONNECTED;
entry->indirect_active = 0; entry->indirect_ready = 0;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: INDIRECT CONNECTED 0x%016llx via %u intermediaries",
(unsigned long long)entry->node_id, entry->intermediariy_count);
cm_update_nodeinfo(entry);
cm_deliver_event(entry, CONN_EVENT_UP);
}
/* Вычисляет лучших общих посредников по данным из EXCHANGE_RESP:
@ -76,9 +116,11 @@ void cm_compute_intermediaries(struct CONN_MGR_ENTRY* entry, struct CM_EXCHANGE_
struct TOPO_GROUP_NODE* nq=topo_node_find_by_id(entry->mgr->group,all[i].nid); if(nq)pkt.selected[i].rtt=topo_get_chain_rtt(nq);}
struct ll_entry* qe=queue_entry_new(0);
if(qe){qe->dgram=u_malloc(1+sizeof(pkt));qe->dgram[0]=ETCP_RT_ID_CONN_MGR;memcpy(qe->dgram+1,&pkt,sizeof(pkt));qe->len=(uint16_t)(1+sizeof(pkt));etcp_route_send(entry->mgr->instance,entry->mgr->group->group_id,entry->node_id,qe,1,0);}
entry->conn_type=CONN_TYPE_INDIRECT; entry->state=CONN_MGR_STATE_CONNECTED;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL,"conn_mgr: indirect OK for 0x%016llx via %u intermediaries",(unsigned long long)entry->node_id,sel);
cm_update_nodeinfo(entry); cm_deliver_event(entry,CONN_EVENT_UP);
/* reverse приоритетнее: если он ещё идёт — откладываем финализацию (indirect_ready),
* иначе финализируем INDIRECT сразу. */
if (entry->reverse_active) { entry->indirect_ready = 1; DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: indirect deferred 0x%016llx (reverse active)", (unsigned long long)entry->node_id); return; }
cm_finalize_indirect(entry);
}
void cm_exchange_probe_retry_cb(void* arg) {
@ -89,7 +131,7 @@ void cm_exchange_probe_retry_cb(void* arg) {
if(!nq) continue; /* self / не в группе — не посредник, пропускаем */
if(!cm_is_rtt_fresh(nq,now)){need++;break;}
}
if(!need){cm_compute_intermediaries(ep->entry,&ep->cached);struct cm_exchange_pending**pp=&ep->entry->mgr->exchange_pending;while(*pp){if(*pp==ep){*pp=ep->next;break;}pp=&(*pp)->next;}u_free(ep);}
if(!need){cm_compute_intermediaries(ep->entry,&ep->cached);struct cm_exchange_pending**pp=&ep->entry->mgr->exchange_pending;while(*pp){if(*pp==ep){*pp=ep->next;break;}pp=&(*pp)->next;}if(ep->pkt)u_free(ep->pkt);u_free(ep);}
else ep->timer=uasync_set_timeout(ep->entry->mgr->instance->ua,2000,ep,cm_exchange_probe_retry_cb,"cm_exch_probe");
}
@ -114,7 +156,7 @@ void cm_handle_interm_exchange_resp(struct CONN_MGR* mgr, const uint8_t* data, s
if(!nq) continue; /* self / не в группе — не посредник, пропускаем */
if(!cm_is_rtt_fresh(nq,now)){need++;route_connectivity_probe_node(mgr->instance,mgr->group,nq);}
}
if(!need){ep->probed=1;cm_compute_intermediaries(ep->entry,&ep->cached);struct cm_exchange_pending**pp=&mgr->exchange_pending;while(*pp){if(*pp==ep){*pp=ep->next;break;}pp=&(*pp)->next;}u_free(ep);return;}
if(!need){ep->probed=1;cm_compute_intermediaries(ep->entry,&ep->cached);struct cm_exchange_pending**pp=&mgr->exchange_pending;while(*pp){if(*pp==ep){*pp=ep->next;break;}pp=&(*pp)->next;}if(ep->pkt)u_free(ep->pkt);u_free(ep);return;}
ep->timer=uasync_set_timeout(mgr->instance->ua,2000,ep,cm_exchange_probe_retry_cb,"cm_exch_probe");
}
@ -128,11 +170,12 @@ void cm_handle_interm_selected(struct CONN_MGR* mgr, uint64_t src, const uint8_t
struct CONN_MGR_ENTRY* entry = cm_find_entry(mgr, src);
if(!entry)return;
if(entry->state==CONN_MGR_STATE_CONNECTED)return; /* симметричный случай — UP уже доставлен инициаторским путём */
if(entry->main.phase!=3)return; /* не INDIRECT-фаза */
if(!entry->indirect_active && !entry->reverse_active)return; /* не INDIRECT/REVERSE фаза */
for(uint8_t i=0;i<c;i++)entry->intermediaries[i]=sel->selected[i].node_id;
entry->intermediariy_count=c; entry->rr_idx=0; entry->conn_type=CONN_TYPE_INDIRECT; entry->state=CONN_MGR_STATE_CONNECTED;
cm_update_nodeinfo(entry); cm_deliver_event(entry,CONN_EVENT_UP);
entry->intermediariy_count=c; entry->rr_idx=0;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL,"conn_mgr: got INTERM_SELECTED from 0x%016llx count=%u",(unsigned long long)src,c);
if (entry->reverse_active) { entry->indirect_ready = 1; DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "conn_mgr: indirect deferred 0x%016llx (reverse active)", (unsigned long long)entry->node_id); return; }
cm_finalize_indirect(entry);
}
/* Цель получила EXCHANGE_REQ: измеряет RTT до кандидатов инициатора (если

17
src/routing_layer/conn_mgr_priv.h

@ -28,8 +28,10 @@ struct stcp_link;
#define CONN_MGR_LOCAL_SCAN_ATTEMPTS 3
#define CONN_MGR_LOCAL_SCAN_TIMEOUT_MS 100
#define CONN_MGR_CONNECT_DIRECT_TIMEOUT_MS 2000
#define CONN_MGR_CONNECT_REVERSE_TIMEOUT_MS 15000
#define CONN_MGR_CONNECT_REVERSE_TIMEOUT_MS 3000
#define CONN_MGR_INTERM_EXCHANGE_TIMEOUT_MS 15000
#define CONN_MGR_REVERSE_RETRANSMIT_MS 1000
#define CONN_MGR_EXCHANGE_RETRANSMIT_MS 1000
#define CACHE_MS 20000 /* время жизни RTT-замера ~20с (через cm_is_rtt_fresh, мс) */
#define CONN_MGR_MAX_INTERMEDIARIES 3
@ -79,12 +81,13 @@ struct CM_DISCONNECT { uint8_t cmd, subcmd; uint64_t node_id; };
/* ═══════ внутренние структуры ═══════ */
struct cm_reverse_pending { struct cm_reverse_pending* next; uint32_t request_id; struct CONN_MGR_ENTRY* entry; };
struct cm_reverse_pending { struct cm_reverse_pending* next; uint32_t request_id; struct CONN_MGR_ENTRY* entry; uint8_t* pkt; size_t pkt_len; };
struct cm_exchange_pending {
struct cm_exchange_pending* next; uint32_t request_id;
struct CONN_MGR_ENTRY* entry; void* timer;
struct CM_EXCHANGE_RESP cached; uint8_t received, probed;
uint8_t* pkt; size_t pkt_len;
};
struct CONN_MGR_HANDLE {
@ -101,7 +104,10 @@ struct CONN_MGR_ENTRY {
struct CONN_MGR_HANDLE* handles; struct NODE_CONN_DIRECT* ncd_handle;
uint64_t intermediaries[3]; uint8_t intermediariy_count, rr_idx;
uint8_t local_scan_state, main_connect_state, db_loaded;
struct { uint8_t phase; void* timer; uint32_t request_id; void* conn_ctx; } main;
uint8_t reverse_active, indirect_active, indirect_ready;
void* reverse_timer;
uint64_t reverse_deadline_tb, indirect_deadline_tb;
struct { void* timer; uint32_t request_id; void* conn_ctx; } main;
struct CONN_MGR* mgr;
};
@ -139,7 +145,8 @@ void cm_start_phase_indirect(struct CONN_MGR_ENTRY* entry);
void cm_start_local_scan(struct CONN_MGR_ENTRY* entry);
void cm_handle_direct_req(uint64_t src, struct CONN_MGR* mgr, const uint8_t* data, size_t len);
void cm_handle_disconnect(struct CONN_MGR* mgr, uint64_t node_id);
void cm_reverse_timeout_cb(void* arg);
void cm_reverse_tick_cb(void* arg);
void cm_finalize_indirect(struct CONN_MGR_ENTRY* entry);
void cm_add_v4_link(struct ETCP_CONN* conn, const uint8_t* addr, uint16_t port, struct ETCP_SOCKET* s);
void cm_add_v6_link(struct ETCP_CONN* conn, const uint8_t addr[16], uint16_t port, struct ETCP_SOCKET* s);
int cm_has_direct_ip(struct CONN_MGR* mgr, struct TOPO_GROUP_NODE* nq);
@ -154,7 +161,7 @@ uint8_t cm_sock_v6_classify(const struct ETCP_SOCKET* s);
#define CM_V6_ANY 4
/* indirect */
void cm_exchange_timeout_cb(void* arg);
void cm_exchange_tick_cb(void* arg);
void cm_compute_intermediaries(struct CONN_MGR_ENTRY* entry, struct CM_EXCHANGE_RESP* resp);
void cm_exchange_probe_retry_cb(void* arg);
void cm_handle_interm_exchange_req(uint64_t src, struct CONN_MGR* mgr, struct CM_EXCHANGE_REQ* req);

Loading…
Cancel
Save