Browse Source

media_delivery: fix SUPER_REPL spam + streaming verification

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 <db_dir>/media/ (matches media_index path logic)
- 9/9 full test, 33/33 total
topo_upd
evgeny 2 months ago
parent
commit
efbb495016
  1. 7
      src/media_delivery/media_delivery.c
  2. 2
      src/media_delivery/media_delivery.h
  3. 10
      tests/test_media_delivery_full.c

7
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) {

2
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 — только суперузел

10
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);
}
}

Loading…
Cancel
Save