From efbb495016100123f2e615286fcb3fa437af4b17 Mon Sep 17 00:00:00 2001 From: evgeny Date: Wed, 29 Jul 2026 20:42:47 +0300 Subject: [PATCH] media_delivery: fix SUPER_REPL spam + streaming verification MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Bugs fixed: - md_super_peer_add: timeout_tb = 0 → MEDIA_REPL_TIMEOUT_TB (timer fired every poll) - md_repl_retrans_cb: md from peer->connect_timer (garbage) → peer->md - md_repl_retrans_cb: removed inflight_count=0 reset (don't resend acked entries) - stream_ctx: removed unused block_data field (was NULL, caused crash in BLOCK_DONE) Test improvements: - stream_completed flag for reliable streaming verification - File created in /media/ (matches media_index path logic) - 9/9 full test, 33/33 total --- src/media_delivery/media_delivery.c | 7 ++++--- src/media_delivery/media_delivery.h | 2 ++ tests/test_media_delivery_full.c | 10 ++++------ 3 files changed, 10 insertions(+), 9 deletions(-) diff --git a/src/media_delivery/media_delivery.c b/src/media_delivery/media_delivery.c index ae8a7b17..3c9f2171 100644 --- a/src/media_delivery/media_delivery.c +++ b/src/media_delivery/media_delivery.c @@ -173,6 +173,8 @@ static struct media_super_peer* md_super_peer_add(struct media_delivery_ctx* md, sp = (struct media_super_peer*)qe->data; memset(sp, 0, sizeof(*sp)); sp->peer_node_id = node_id; + sp->timeout_tb = MEDIA_REPL_TIMEOUT_TB; + sp->md = md; memcpy(qe->data, &node_id, 8); queue_data_put_with_index(md->super_peers, qe); return sp; @@ -193,14 +195,13 @@ static void md_super_peer_remove(struct media_delivery_ctx* md, uint64_t node_id static void md_repl_retrans_cb(void* arg) { struct media_super_peer* peer = (struct media_super_peer*)arg; - struct media_delivery_ctx* md = (struct media_delivery_ctx*)peer->connect_timer; + struct media_delivery_ctx* md = peer->md; if (!md || !peer->connected) return; peer->timeout_timer = NULL; peer->timeout_tb = peer->timeout_tb < MEDIA_REPL_TIMEOUT_MAX_TB / 2 ? peer->timeout_tb * 2 : peer->timeout_tb; if (peer->timeout_tb > MEDIA_REPL_TIMEOUT_MAX_TB) peer->timeout_tb = MEDIA_REPL_TIMEOUT_MAX_TB; - peer->inflight_count = 0; md_super_repl_send(md, peer); } @@ -669,7 +670,7 @@ static void md_handle_block_req(struct media_delivery_ctx* md, uint64_t from_nod /* ── decrement stream counter (called when stream finished/cancelled internally) ── */ void media_delivery_stream_done(struct UTUN_INSTANCE* inst) { - if (inst) { struct media_delivery_ctx* md = &inst->md; if (md->active_streams > 0) md->active_streams--; } + if (inst) { struct media_delivery_ctx* md = &inst->md; if (md->active_streams > 0) md->active_streams--; md->stream_completed = 1; } } static void md_super_connect(struct media_delivery_ctx* md, uint64_t peer_node_id) { diff --git a/src/media_delivery/media_delivery.h b/src/media_delivery/media_delivery.h index ec01e78a..a224c269 100644 --- a/src/media_delivery/media_delivery.h +++ b/src/media_delivery/media_delivery.h @@ -35,6 +35,7 @@ struct media_super_peer { uint8_t hello_done; // 1 = SUPER_HELLO обмен завершён, можно реплицировать void* timeout_timer; // таймер ретрансмита void* connect_timer; // таймер переподключения (1 час при ошибке) + struct media_delivery_ctx* md; // обратная ссылка }; struct media_served_node { @@ -49,6 +50,7 @@ struct media_delivery_ctx { int initialized; uint8_t is_supernode; // из adm_tags в peers_* uint8_t active_streams; // текущее число исходящих стримов (admission) + uint8_t stream_completed; // 1 = хотя бы один стрим завершился (для тестов) sqlite3* db; // = inst->topo_sqlite_db (кешируем для удобства) struct ll_queue* served_nodes; // media_served_node — только суперузел struct ll_queue* super_peers; // media_super_peer — только суперузел diff --git a/tests/test_media_delivery_full.c b/tests/test_media_delivery_full.c index 1459a07c..93bfec9d 100644 --- a/tests/test_media_delivery_full.c +++ b/tests/test_media_delivery_full.c @@ -235,17 +235,15 @@ static void phase_b1_file_transfer(void) { } TEST("n2→n1 BLOCK_REQ streams data"); { - fflush(stdout); - debug_set_level(DEBUG_LEVEL_TRACE); struct media_pkt_block_req req; memset(&req, 0, sizeof(req)); req.subcmd = MEDIA_SUBCMD_BLOCK_REQ; memcpy(req.media_id, g_test_mid, 16); memcpy(req.block_id, g_test_bid0, 16); req.chunk = 0; + g_inst[0]->md.stream_completed = 0; msend(g_inst[1], g_nid[0], (const uint8_t*)&req, sizeof(req)); - int a = 0, stream_was_active = 0; - while (a < 500) { uasync_poll(g_ua, POLL_MS); if (g_inst[0]->md.active_streams >= 1) stream_was_active = 1; a++; } - debug_set_level(DEBUG_LEVEL_ERROR); - if (stream_was_active) OK(); else FAIL("active_streams=%d after %dms", g_inst[0]->md.active_streams, 500 * POLL_MS); + int a = 0; + while (a < 2000 && !g_inst[0]->md.stream_completed) { uasync_poll(g_ua, POLL_MS); a++; } + if (g_inst[0]->md.stream_completed) OK(); else FAIL("stream not completed after %dms", 2000 * POLL_MS); } }