Browse Source

media_download: fix block_id OOB, add backpressure, close conn_mgr

- md_dl_send_block_req now takes explicit block_id (not peer-local index
  into dl->block_ids) — fixes OOB 'aaaaaaaa...' garbage block_id
- conn_cb sends only first unstarted block per peer (backpressure)
- BLOCK_DONE success triggers next block from same peer
- conn_mgr handles stored in peer->cm_handle, closed on assembly/cancel
topo_upd
evgeny 2 months ago
parent
commit
41e9c9b0e6
  1. 50
      src/media_delivery/media_download.c
  2. 1
      src/media_delivery/media_download.h

50
src/media_delivery/media_download.c

@ -32,6 +32,9 @@ static struct media_download* md_dl_find(struct UTUN_INSTANCE* inst, const uint8
static void md_dl_free(struct media_download* dl) {
if (!dl) return;
for (int pi = 0; pi < dl->num_peers; pi++) {
if (dl->peers[pi].cm_handle) { conn_mgr_close(dl->peers[pi].cm_handle); dl->peers[pi].cm_handle = NULL; }
}
if (dl->query_timer) { uasync_cancel_timeout(NULL, dl->query_timer); dl->query_timer = NULL; }
if (dl->timeout_timer) { uasync_cancel_timeout(NULL, dl->timeout_timer); dl->timeout_timer = NULL; }
if (dl->block_ids) { u_free(dl->block_ids); dl->block_ids = NULL; }
@ -68,19 +71,20 @@ static int md_dl_send(struct UTUN_INSTANCE* inst, uint64_t dst, const uint8_t* d
/* ── build block request ── */
static int md_dl_send_block_req(struct UTUN_INSTANCE* inst, uint64_t dst,
struct media_download* dl, int bi) {
struct media_download* dl, const uint8_t* block_id,
int chunk_idx, int peer_local_bi) {
struct media_pkt_block_req req;
memset(&req, 0, sizeof(req));
req.subcmd = MEDIA_SUBCMD_BLOCK_REQ;
req.group_id = dl->group_id;
memcpy(req.media_id, dl->media_id, 16);
memcpy(req.block_id, dl->block_ids + bi * 16, 16);
req.chunk = (uint32_t)bi;
memcpy(req.block_id, block_id, 16);
req.chunk = (uint32_t)chunk_idx;
req.offset = 0;
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: BLOCK_REQ to 0x%016llx group=%016llx bi=%d dl->num_blocks=%d block_id=%02x%02x%02x%02x...",
MDL_ID, (unsigned long long)dst, (unsigned long long)dl->group_id, bi, dl->num_blocks,
dl->block_ids[bi * 16], dl->block_ids[bi * 16 + 1], dl->block_ids[bi * 16 + 2], dl->block_ids[bi * 16 + 3]);
MDL_ID, (unsigned long long)dst, (unsigned long long)dl->group_id, peer_local_bi, dl->num_blocks,
block_id[0], block_id[1], block_id[2], block_id[3]);
return md_dl_send(inst, dst, (const uint8_t*)&req, sizeof(req));
}
@ -118,7 +122,7 @@ static void md_dl_send_have_block(struct UTUN_INSTANCE* inst, struct media_downl
static void md_dl_conn_cb(struct CONN_MGR_HANDLE* h, uint64_t node_id, uint64_t group_id,
enum conn_mgr_event event, void* arg) {
struct media_download* dl = (struct media_download*)arg; (void)h; (void)group_id;
struct media_download* dl = (struct media_download*)arg;
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: conn_cb event=%d node=0x%016llx dl=%p active=%d num_peers=%d num_blocks=%d",
MDL_ID, event, (unsigned long long)node_id, (void*)dl, dl->active, dl->num_peers, dl->num_blocks);
if (event != CONN_EVENT_UP) {
@ -126,10 +130,11 @@ static void md_dl_conn_cb(struct CONN_MGR_HANDLE* h, uint64_t node_id, uint64_t
MDL_ID, (unsigned long long)node_id, event);
return;
}
/* send block requests to this peer */
/* send first unstarted block request to this peer */
for (int pi = 0; pi < dl->num_peers; pi++) {
if (dl->peers[pi].node_id == node_id) {
dl->peers[pi].connected = 1;
dl->peers[pi].cm_handle = h; /* store for close on completion */
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: matched peer[%d] num_blocks=%d",
MDL_ID, pi, dl->peers[pi].num_blocks);
for (int bi = 0; bi < dl->peers[pi].num_blocks; bi++) {
@ -137,7 +142,8 @@ static void md_dl_conn_cb(struct CONN_MGR_HANDLE* h, uint64_t node_id, uint64_t
dl->peers[pi].blocks[bi].started = 1;
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: conn_cb BLOCK_REQ: peer-local bi=%d",
MDL_ID, bi);
md_dl_send_block_req(dl->inst, node_id, dl, bi);
md_dl_send_block_req(dl->inst, node_id, dl, dl->peers[pi].blocks[bi].block_id, bi, bi);
break;
}
}
}
@ -293,7 +299,7 @@ void media_download_handle_query_resp(struct UTUN_INSTANCE* inst,
for (int bi = 0; bi < dl->peers[pi].num_blocks; bi++) {
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: direct BLOCK_REQ: peer-local bi=%d peer_num_blocks=%d",
MDL_ID, bi, dl->peers[pi].num_blocks);
md_dl_send_block_req(inst, dl->peers[pi].node_id, dl, bi);
md_dl_send_block_req(inst, dl->peers[pi].node_id, dl, dl->peers[pi].blocks[bi].block_id, bi, bi);
}
}
}
@ -400,6 +406,25 @@ void media_download_handle_done(struct UTUN_INSTANCE* inst,
/* send HAVE_BLOCK to supernode */
md_dl_send_have_block(inst, dl, bi);
/* send next block from the same peer (backpressure: one at a time) */
for (int pi = 0; pi < dl->num_peers; pi++) {
for (int bj = 0; bj < dl->peers[pi].num_blocks; bj++) {
if (memcmp(dl->peers[pi].blocks[bj].block_id, dl->block_ids + bi * 16, 16) != 0) continue;
for (int bk = 0; bk < dl->peers[pi].num_blocks; bk++) {
if (!dl->peers[pi].blocks[bk].started) {
dl->peers[pi].blocks[bk].started = 1;
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: BLOCK_DONE → next req: peer[%d] bi=%d",
MDL_ID, pi, bk);
md_dl_send_block_req(inst, dl->peers[pi].node_id, dl,
dl->peers[pi].blocks[bk].block_id, bk, bk);
goto blockreq_sent;
}
}
goto blockreq_sent;
}
}
blockreq_sent: ;
/* check assembly */
if (dl->blocks_validated == dl->num_blocks) {
/* assemble file */
@ -428,6 +453,13 @@ void media_download_handle_done(struct UTUN_INSTANCE* inst,
} }
dl->assembled = 1;
dl->err = 0;
/* close all conn_mgr handles */
for (int pi = 0; pi < dl->num_peers; pi++) {
if (dl->peers[pi].cm_handle) {
conn_mgr_close(dl->peers[pi].cm_handle);
dl->peers[pi].cm_handle = NULL;
}
}
if (dl->done_cb) dl->done_cb(dl->done_arg, dl->err);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: file assembled: %s", MDL_ID, dl->dest_path);
}

1
src/media_delivery/media_download.h

@ -23,6 +23,7 @@ struct media_index_result;
struct media_download_peer {
uint64_t node_id;
uint8_t connected;
void* cm_handle;
int num_blocks;
struct {
uint8_t block_id[16];

Loading…
Cancel
Save