From 41e9c9b0e6c68b3feb67f876db2d1888ba5bf4dc Mon Sep 17 00:00:00 2001 From: evgeny Date: Tue, 4 Aug 2026 00:02:20 +0300 Subject: [PATCH] media_download: fix block_id OOB, add backpressure, close conn_mgr MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 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 --- src/media_delivery/media_download.c | 50 +++++++++++++++++++++++------ src/media_delivery/media_download.h | 1 + 2 files changed, 42 insertions(+), 9 deletions(-) diff --git a/src/media_delivery/media_download.c b/src/media_delivery/media_download.c index 0295a931..142717b6 100644 --- a/src/media_delivery/media_download.c +++ b/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); } diff --git a/src/media_delivery/media_download.h b/src/media_delivery/media_download.h index 1187b884..f6977bad 100644 --- a/src/media_delivery/media_download.h +++ b/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];