Browse Source

media: восстановить conn_mgr_open (DIR/REV/IND) + популяция best_candidates из RTT + watchdog не трогает CONN-блоки

v2
evgeny 3 weeks ago
parent
commit
2045d4a169
  1. 33
      src/media_delivery/media_download.c
  2. 3
      src/routing_layer/topo_node.c

33
src/media_delivery/media_download.c

@ -204,12 +204,28 @@ static int md_dl_request_block(struct media_download* dl, int bi) {
dl->inflight++; dl->inflight++;
struct UTUN_INSTANCE* inst = dl->inst; struct UTUN_INSTANCE* inst = dl->inst;
/* BLOCK_REQ шлём через роутер (etcp_route_send): при наличии прямого соединения if (instance_find_conn(inst, holder)) {
* пойдёт напрямую, иначе — транзитом. Держатель ответит чанками по лучшему пути, p->connected = 1;
* поэтому прямое соединение не требуется (NAT-узлы качают транзитом). */ b->state = MD_BLK_REQ;
p->connected = 1; md_dl_send_block_req(inst, holder, dl, bi);
b->state = MD_BLK_REQ; } else {
md_dl_send_block_req(inst, holder, dl, bi); struct TOPO_GROUP* grp = inst->topo_groups ? topo_groups_find(inst->topo_groups, dl->group_id) : NULL;
if (grp && grp->conn_mgr) {
int rc = conn_mgr_open(inst, grp->group_id, holder, md_dl_conn_cb, dl, &p->cm_handle);
if (rc != 0) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: conn_mgr_open failed rc=%d for 0x%016llx", MDL_ID, rc, (unsigned long long)holder);
dl->inflight--;
b->state = MD_BLK_IDLE;
return -1;
}
/* CONN — BLOCK_REQ отправится в md_dl_conn_cb по CONN_EVENT_UP */
} else {
/* нет conn_mgr — шлём напрямую best-effort */
p->connected = 1;
b->state = MD_BLK_REQ;
md_dl_send_block_req(inst, holder, dl, bi);
}
}
return 0; return 0;
} }
@ -262,10 +278,11 @@ int md_dl_check_stall(struct media_download* dl, uint64_t now_tb) {
struct UTUN_INSTANCE* inst = dl->inst; struct UTUN_INSTANCE* inst = dl->inst;
uint64_t stall = md_dl_stall_tb(inst); uint64_t stall = md_dl_stall_tb(inst);
/* зависшие блоки — фейловер */ /* зависшие блоки — фейловер (только REQ/RECV; CONN ждёт установку соединения
* через conn_mgr, которая заканчивается CONN_EVENT_UP/TIMEOUT в md_dl_conn_cb) */
for (int bi = 0; bi < dl->num_blocks; bi++) { for (int bi = 0; bi < dl->num_blocks; bi++) {
struct media_download_block* b = &dl->blocks[bi]; struct media_download_block* b = &dl->blocks[bi];
if (b->state != MD_BLK_CONN && b->state != MD_BLK_REQ && b->state != MD_BLK_RECV) continue; if (b->state != MD_BLK_REQ && b->state != MD_BLK_RECV) continue;
if (now_tb - b->last_progress_tb >= stall) { if (now_tb - b->last_progress_tb >= stall) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: STALL blk=%d state=%d no progress %llums", DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: STALL blk=%d state=%d no progress %llums",
MDL_ID, bi, b->state, (unsigned long long)(now_tb - b->last_progress_tb)); MDL_ID, bi, b->state, (unsigned long long)(now_tb - b->last_progress_tb));

3
src/routing_layer/topo_node.c

@ -11,6 +11,7 @@
#include "config_parser.h" #include "config_parser.h"
#include "topo_node.h" #include "topo_node.h"
#include "topo_group.h" #include "topo_group.h"
#include "conn_mgr.h"
#include "route_lib.h" #include "route_lib.h"
#include "etcp_debug.h" #include "etcp_debug.h"
#include "topo_node_sqlite.h" #include "topo_node_sqlite.h"
@ -580,6 +581,8 @@ void topo_node_ping_update_rtt(struct TOPO_GROUPS* groups, uint64_t node_id, uin
nq->connectivity.real_min_rtt = rtt; nq->connectivity.real_min_rtt = rtt;
DEBUG_TRACE(DEBUG_CATEGORY_BGP, "ping_update_rtt: node=%016llx rtt=%u (%llu) group=%016llx", DEBUG_TRACE(DEBUG_CATEGORY_BGP, "ping_update_rtt: node=%016llx rtt=%u (%llu) group=%016llx",
(unsigned long long)node_id, (unsigned)rtt, (unsigned long long)now, (unsigned long long)g->group_id); (unsigned long long)node_id, (unsigned)rtt, (unsigned long long)now, (unsigned long long)g->group_id);
if (g->conn_mgr)
conn_mgr_update_best_candidates(g->conn_mgr, node_id, rtt);
if (groups->instance && groups->instance->topo_sqlite_db) if (groups->instance && groups->instance->topo_sqlite_db)
topo_node_sqlite_update_rtt(groups->instance->topo_sqlite_db, node_id, rtt); topo_node_sqlite_update_rtt(groups->instance->topo_sqlite_db, node_id, rtt);
topo_fire_nodeinfo_cbk(groups->instance, g, nq); topo_fire_nodeinfo_cbk(groups->instance, g, nq);

Loading…
Cancel
Save