diff --git a/src/media_delivery/media_download.c b/src/media_delivery/media_download.c index 142717b6..80c9ff58 100644 --- a/src/media_delivery/media_download.c +++ b/src/media_delivery/media_download.c @@ -57,6 +57,7 @@ static void md_mkdir_parent(const char* filepath) { } static int md_dl_send(struct UTUN_INSTANCE* inst, uint64_t dst, const uint8_t* data, size_t len) { + if (!inst || !inst->topo_groups || !inst->connections) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: md_dl_send: no routing (inst=%p topo=%p conn=%p)", MDL_ID, (void*)inst, (void*)(inst ? inst->topo_groups : NULL), (void*)(inst ? inst->connections : NULL)); return -1; } struct ll_entry* e = queue_entry_new(0); if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_entry_new failed for send to 0x%016llx", MDL_ID, (unsigned long long)dst); return -1; } e->dgram = u_malloc(len + 1); @@ -120,7 +121,7 @@ static void md_dl_send_have_block(struct UTUN_INSTANCE* inst, struct media_downl /* ── conn_mgr callback ── */ -static void md_dl_conn_cb(struct CONN_MGR_HANDLE* h, uint64_t node_id, uint64_t group_id, +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; DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: conn_cb event=%d node=0x%016llx dl=%p active=%d num_peers=%d num_blocks=%d", diff --git a/src/media_delivery/media_download.h b/src/media_delivery/media_download.h index f6977bad..4cd8624a 100644 --- a/src/media_delivery/media_download.h +++ b/src/media_delivery/media_download.h @@ -105,6 +105,12 @@ void media_download_handle_done(struct UTUN_INSTANCE* inst, void media_download_handle_overloaded(struct UTUN_INSTANCE* inst, const uint8_t* data, size_t len); +/* ── conn_mgr callback (non-static for test visibility) ── */ +struct CONN_MGR_HANDLE; +enum conn_mgr_event; +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); + #ifdef __cplusplus } #endif diff --git a/tests/test_media_delivery_download.c b/tests/test_media_delivery_download.c index 76675603..bd3df5ca 100644 --- a/tests/test_media_delivery_download.c +++ b/tests/test_media_delivery_download.c @@ -18,6 +18,7 @@ #include "../lib/mem.h" #include "../transport_layer/secure_channel.h" #include "../routing_layer/topo_node_sqlite.h" +#include "../routing_layer/conn_mgr.h" #include #include @@ -300,6 +301,192 @@ done5: } } +/* ──────────────────────────────────────────────────────────────── + Tests for block download fixes (backpressure, OOB block_id, conn_mgr close) + ──────────────────────────────────────────────────────────────── */ + +static struct media_download* make_test_dl(sqlite3* db, struct UTUN_INSTANCE* inst, int num_blocks) { + if (!inst->md.downloads) inst->md.downloads = queue_new(inst->ua, 256, offsetof(struct media_download, media_id), 16, "md_dl_test"); + + struct media_download* dl = u_calloc(1, sizeof(*dl)); + if (!dl) return NULL; + make_uuid(dl->media_id); + dl->num_blocks = num_blocks; + snprintf(dl->dest_path, sizeof(dl->dest_path), "%s/test.bin", g_temp_dir); + dl->block_ids = u_calloc((size_t)num_blocks, 16); + dl->block_sigs = u_calloc((size_t)num_blocks, 64); + for (int i = 0; i < num_blocks; i++) { make_uuid(dl->block_ids + i * 16); memset(dl->block_sigs + i * 64, (uint8_t)(0x42 + i), 64); } + dl->block_size = CHUNK_SIZE; dl->file_size = (int64_t)CHUNK_SIZE * num_blocks; + dl->active = 1; dl->inst = inst; + memcpy(dl->ll.data, dl->media_id, 16); + + struct ll_entry* qe = queue_entry_new(sizeof(*dl)); + if (!qe) { u_free(dl->block_ids); u_free(dl->block_sigs); u_free(dl); return NULL; } + memcpy(qe->data, dl, sizeof(*dl)); u_free(dl); + queue_data_put_with_index(inst->md.downloads, qe); + return (struct media_download*)qe->data; +} + +static struct media_download_peer* add_peer(struct media_download* dl, uint64_t node_id, int num_blocks) { + if (dl->num_peers >= 10) return NULL; + int pi = dl->num_peers++; + struct media_download_peer* p = &dl->peers[pi]; + p->node_id = node_id; + p->num_blocks = num_blocks; + for (int i = 0; i < num_blocks; i++) { make_uuid(p->blocks[i].block_id); } + return p; +} + +/* ── Тест: conn_cb отправляет только первый блок (backpressure) ── */ +static void test_conn_cb_backpressure(void) { + TEST("conn_cb starts only first block (backpressure)"); { + sqlite3* db = NULL; sqlite3_open(":memory:", &db); + struct UTUN_INSTANCE* inst = make_minimal_instance(db); + struct media_download* dl = make_test_dl(db, inst, 1); + if (!dl) { FAIL("make_test_dl"); sqlite3_close(db); u_free(inst); return; } + + add_peer(dl, 0x1234, 5); /* 5 blocks on peer */ + + md_dl_conn_cb(NULL, 0x1234, 0, CONN_EVENT_UP, dl); + + int started_count = 0; + for (int i = 0; i < dl->peers[0].num_blocks; i++) + if (dl->peers[0].blocks[i].started) started_count++; + + if (dl->peers[0].connected && started_count == 1 && dl->peers[0].blocks[0].started && !dl->peers[0].blocks[1].started) + OK(); else FAIL("connected=%d started=%d b0=%d b1=%d", dl->peers[0].connected, started_count, dl->peers[0].blocks[0].started, dl->peers[0].blocks[1].started); + + u_free(dl->block_ids); u_free(dl->block_sigs); + if (inst->md.downloads) queue_free(inst->md.downloads); + u_free(inst); sqlite3_close(db); + } +} + +/* ── Тест: conn_cb для разных пиров ── */ +static void test_conn_cb_multi_peer(void) { + TEST("conn_cb starts first block per peer"); { + sqlite3* db = NULL; sqlite3_open(":memory:", &db); + struct UTUN_INSTANCE* inst = make_minimal_instance(db); + struct media_download* dl = make_test_dl(db, inst, 1); + if (!dl) { FAIL("make_test_dl"); sqlite3_close(db); u_free(inst); return; } + + add_peer(dl, 0x1000, 3); + add_peer(dl, 0x2000, 2); + + md_dl_conn_cb(NULL, 0x1000, 0, CONN_EVENT_UP, dl); + md_dl_conn_cb(NULL, 0x2000, 0, CONN_EVENT_UP, dl); + + int ok = dl->peers[0].connected && dl->peers[1].connected + && dl->peers[0].blocks[0].started && !dl->peers[0].blocks[1].started + && dl->peers[1].blocks[0].started && !dl->peers[1].blocks[1].started; + if (ok) OK(); else FAIL("p0:conn=%d b0=%d b1=%d p1:conn=%d b0=%d b1=%d", + dl->peers[0].connected, dl->peers[0].blocks[0].started, dl->peers[0].blocks[1].started, + dl->peers[1].connected, dl->peers[1].blocks[0].started, dl->peers[1].blocks[1].started); + + u_free(dl->block_ids); u_free(dl->block_sigs); + if (inst->md.downloads) queue_free(inst->md.downloads); + u_free(inst); sqlite3_close(db); + } +} + +/* ── Тест: BLOCK_DONE → relay: следующий блок с того же пира ── */ +static void test_done_block_relay(void) { + TEST("BLOCK_DONE triggers next block from same peer"); { + sqlite3* db = NULL; sqlite3_open(":memory:", &db); + struct UTUN_INSTANCE* inst = make_minimal_instance(db); + struct media_download* dl = make_test_dl(db, inst, 2); + if (!dl) { FAIL("make_test_dl"); sqlite3_close(db); u_free(inst); return; } + + /* peer with 3 blocks for 2-block download */ + add_peer(dl, 0x5555, 3); + /* map dl->block_ids to peer blocks for matching in handle_done */ + for (int i = 0; i < dl->num_blocks; i++) + memcpy(dl->peers[0].blocks[i].block_id, dl->block_ids + i * 16, 16); + /* pre-set: block 0 started and fake-received via conn_cb */ + dl->peers[0].blocks[0].started = 1; + + /* write chunk for block 0 */ + uint8_t data0[CHUNK_SIZE]; memset(data0, 0xBB, CHUNK_SIZE); + { char c[1024]; snprintf(c, sizeof(c), "%s.chunk_0", dl->dest_path); write_file(c, data0, CHUNK_SIZE); } + + /* BLOCK_DONE for block 0 with matching sig */ + { uint8_t pkt[MEDIA_BLOCK_DONE_SIZE]; memset(pkt, 0, sizeof(pkt)); + struct media_pkt_block_done* bd = (struct media_pkt_block_done*)pkt; + bd->subcmd = MEDIA_SUBCMD_BLOCK_DONE; + memcpy(bd->media_id, dl->media_id, 16); memcpy(bd->block_id, dl->block_ids, 16); + bd->chunk = 0; bd->total_size = CHUNK_SIZE; + memcpy(bd->block_sig, dl->block_sigs, 64); + media_download_handle_done(inst, pkt, sizeof(pkt)); } + + /* block 1 should be started (relay); block 2 NOT (peer-local index unrelated to dl) */ + int ok = dl->peers[0].blocks[1].started && !dl->peers[0].blocks[2].started; + if (ok) OK(); else FAIL("b1_started=%d b2_started=%d", dl->peers[0].blocks[1].started, dl->peers[0].blocks[2].started); + + u_free(dl->block_ids); u_free(dl->block_sigs); + if (inst->md.downloads) queue_free(inst->md.downloads); + u_free(inst); sqlite3_close(db); + } +} + +/* ── Тест: ассемблинг при всех блоках ── */ +static void test_done_assembly(void) { + TEST("BLOCK_DONE all blocks → assembly"); { + sqlite3* db = NULL; sqlite3_open(":memory:", &db); + struct UTUN_INSTANCE* inst = make_minimal_instance(db); + struct media_download* dl = make_test_dl(db, inst, 2); + if (!dl) { FAIL("make_test_dl"); sqlite3_close(db); u_free(inst); return; } + + /* write two chunk files with different content */ + uint8_t d0[CHUNK_SIZE], d1[CHUNK_SIZE]; + memset(d0, 0xA0, CHUNK_SIZE); memset(d1, 0xB1, CHUNK_SIZE); + { char c[1024]; snprintf(c, sizeof(c), "%s.chunk_0", dl->dest_path); write_file(c, d0, CHUNK_SIZE); } + { char c[1024]; snprintf(c, sizeof(c), "%s.chunk_1", dl->dest_path); write_file(c, d1, CHUNK_SIZE); } + + /* BLOCK_DONE for both */ + for (int bi = 0; bi < 2; bi++) { + uint8_t pkt[MEDIA_BLOCK_DONE_SIZE]; memset(pkt, 0, sizeof(pkt)); + struct media_pkt_block_done* bd = (struct media_pkt_block_done*)pkt; + bd->subcmd = MEDIA_SUBCMD_BLOCK_DONE; + memcpy(bd->media_id, dl->media_id, 16); memcpy(bd->block_id, dl->block_ids + bi * 16, 16); + bd->chunk = (uint32_t)bi; bd->total_size = CHUNK_SIZE; + memcpy(bd->block_sig, dl->block_sigs + bi * 64, 64); + media_download_handle_done(inst, pkt, sizeof(pkt)); + } + + int sz = file_size(dl->dest_path); + int ok = dl->assembled && sz == CHUNK_SIZE * 2; + if (ok) OK(); else FAIL("assembled=%d size=%d expected=%d", dl->assembled, sz, CHUNK_SIZE * 2); + + u_free(dl->block_ids); u_free(dl->block_sigs); + if (inst->md.downloads) queue_free(inst->md.downloads); + u_free(inst); sqlite3_close(db); + } +} + +/* ── Тест: пир имеет больше блоков чем dl → без OOB ── */ +static void test_peer_more_blocks_than_dl(void) { + TEST("conn_cb with peer blocks > dl blocks → no OOB"); { + sqlite3* db = NULL; sqlite3_open(":memory:", &db); + struct UTUN_INSTANCE* inst = make_minimal_instance(db); + struct media_download* dl = make_test_dl(db, inst, 1); /* 1 block */ + if (!dl) { FAIL("make_test_dl"); sqlite3_close(db); u_free(inst); return; } + + add_peer(dl, 0x9999, 6); /* peer has 6 blocks — more than dl */ + + /* should NOT crash (was OOB before fix) */ + md_dl_conn_cb(NULL, 0x9999, 0, CONN_EVENT_UP, dl); + + /* should have started only first peer block; no OOB access to dl->block_ids[1..5] */ + int ok = dl->peers[0].connected && dl->peers[0].blocks[0].started + && !dl->peers[0].blocks[1].started && !dl->peers[0].blocks[5].started; + if (ok) OK(); else FAIL("conn=%d b0=%d b1=%d b5=%d", dl->peers[0].connected, dl->peers[0].blocks[0].started, dl->peers[0].blocks[1].started, dl->peers[0].blocks[5].started); + + u_free(dl->block_ids); u_free(dl->block_sigs); + if (inst->md.downloads) queue_free(inst->md.downloads); + u_free(inst); sqlite3_close(db); + } +} + int main(void) { test_setup(); printf("=== test_media_delivery_download ===\n"); @@ -308,6 +495,13 @@ int main(void) { test_download_done(); test_download_cancel(); + printf("\n=== block download fixes (backpressure, OOB, relay) ===\n"); + test_conn_cb_backpressure(); + test_conn_cb_multi_peer(); + test_done_block_relay(); + test_done_assembly(); + test_peer_more_blocks_than_dl(); + printf("\n%d/%d passed, %d failed\n", g_passed, g_total, g_failed); test_cleanup(); return g_failed > 0 ? 1 : 0;