diff --git a/src/chat/chat_msg.c b/src/chat/chat_msg.c index 24dc09af..58e9cef9 100644 --- a/src/chat/chat_msg.c +++ b/src/chat/chat_msg.c @@ -560,6 +560,10 @@ static void md_download_done_cb(void* arg, int err) { } } else { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: media download failed ch=%s err=%d", CC_ID, ctx->channel_id, err); + /* пометить сообщение ошибкой, чтобы UI показал «повторить» */ + char attrs[64]; + snprintf(attrs, sizeof(attrs), "{\"st\":\"er\"}"); + chat_core_update_local_attrs(ctx->channel_id, ctx->ts, ctx->author_sig, attrs); } uint8_t evt[80]; uint8_t cl = (uint8_t)strlen(ctx->channel_id); evt[0] = cl; memcpy(evt + 1, ctx->channel_id, cl); diff --git a/src/media_delivery/media_delivery.c b/src/media_delivery/media_delivery.c index 2d03f371..06cda570 100644 --- a/src/media_delivery/media_delivery.c +++ b/src/media_delivery/media_delivery.c @@ -472,19 +472,26 @@ static void md_handle_query(struct media_delivery_ctx* md, uint64_t from_node, /* 1) block_availability (supernode or HAVE_BLOCK reports) */ nf = md_ba_find_nodes(md->db, block_id, nodes, 10); - /* 2) media_files: check if we own this block */ - if (nf == 0) { + /* 2) media_files: собственные блоки (автор-не-суперузел тоже отдаёт файлы). + Проверяем всегда — иначе при наличии чужих записей в block_availability + собственный блок автора не попал бы в список держателей. */ + { const char* sql = "SELECT 1 FROM media_files WHERE block_id=? AND node_id=? LIMIT 1"; sqlite3_stmt* st = NULL; if (sqlite3_prepare_v2(md->db, sql, -1, &st, NULL) == SQLITE_OK) { sqlite3_bind_blob(st, 1, block_id, 16, SQLITE_STATIC); sqlite3_bind_int64(st, 2, (sqlite3_int64)md->self_node_id); - if (sqlite3_step(st) == SQLITE_ROW) { nodes[nf++] = md->self_node_id; } + if (sqlite3_step(st) == SQLITE_ROW) { + int dup = 0; + for (int j = 0; j < nf; j++) if (nodes[j] == md->self_node_id) { dup = 1; break; } + if (!dup && nf < 10) nodes[nf++] = md->self_node_id; + } sqlite3_finalize(st); } } for (int j = 0; j < nf && off + 30 <= (int)sizeof(resp); j++) { + if (nodes[j] == from_node) continue; /* сам себе держателем не бывает */ uint16_t rtt = 0; memcpy(resp + off, &nodes[j], 8); off += 8; memcpy(resp + off, &rtt, 2); off += 2; @@ -739,10 +746,29 @@ static void stream_send_chunk_cb(struct ll_queue* q, void* arg) { ch->data_len = (uint16_t)rd; memcpy(pkt + MEDIA_BLOCK_CHUNK_HDR_SIZE, buf, rd); - md_send(md->inst, sc->group_id, sc->dst_node_id, pkt, MEDIA_BLOCK_CHUNK_HDR_SIZE + rd); + int snd = md_send(md->inst, sc->group_id, sc->dst_node_id, pkt, MEDIA_BLOCK_CHUNK_HDR_SIZE + rd); + if (snd != 0) { + DEBUG_ERROR(DEBUG_CATEGORY_MEDIA, "%s: stream chunk send FAILED rc=%d — abort stream chunk=%d to 0x%016llx", + MD_ID, snd, sc->chunk, (unsigned long long)sc->dst_node_id); + md_file_load_dec(md, sc->media_id, sc->dst_node_id); + media_delivery_stream_done(md->inst); + fclose(sc->file); + u_free(sc); + return; + } sc->offset += rd; sc->remaining -= rd; + { + uint64_t sent = sc->offset - sc->block_start; + uint64_t mb_off = (sent - rd) >> 20; + uint64_t mb_end = sent >> 20; + if (sent == rd || mb_end != mb_off || sent >= sc->block_data_len) + DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: stream chunk=%d sent=%llu/%zu (%.0f%%) to 0x%016llx", + MD_ID, sc->chunk, (unsigned long long)sent, sc->block_data_len, + sc->block_data_len ? 100.0 * (double)sent / (double)sc->block_data_len : 0.0, + (unsigned long long)sc->dst_node_id); + } if (sc->remaining == 0) break; @@ -933,8 +959,9 @@ static void md_handle_block_req(struct media_delivery_ctx* md, uint64_t from_nod int64_t offset = sqlite3_column_int64(st, 3); int64_t file_size = sqlite3_column_int64(st, 4); - /* compute chunk start and size */ - int64_t block_start = (int64_t)chunk * chunk_size; + /* compute chunk start and size: offset column хранит n * block_size + (для последнего усечённого блока chunk * chunk_size даёт неверный offset) */ + int64_t block_start = offset; int64_t block_end = block_start + chunk_size; if (block_end > file_size) block_end = file_size; int64_t block_length = block_end - block_start; @@ -1320,6 +1347,8 @@ int media_delivery_init(struct UTUN_INSTANCE* inst) { md->inst = inst; md->db = inst->topo_sqlite_db; md->max_downloads_per_file = 3; + md->dl_stall_timeout_tb = MEDIA_DL_STALL_TIMEOUT_TB; + md->dl_max_attempts = MEDIA_DL_MAX_ATTEMPTS; if (media_delivery_create_tables(inst) != 0) return -1; @@ -1389,7 +1418,20 @@ void media_delivery_destroy(struct UTUN_INSTANCE* inst) { md->super_peers = NULL; } if (md->served_nodes) { queue_free(md->served_nodes); md->served_nodes = NULL; } - if (md->downloads) { queue_free(md->downloads); md->downloads = NULL; } + if (md->downloads) { + struct ll_entry* de = md->downloads->head; + while (de) { + struct media_download* dl = (struct media_download*)de->data; + if (dl->watchdog_timer) { uasync_cancel_timeout(md->inst->ua, dl->watchdog_timer); dl->watchdog_timer = NULL; } + 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->blocks) u_free(dl->blocks); + if (dl->block_sigs) u_free(dl->block_sigs); + de = de->next; + } + queue_free(md->downloads); + md->downloads = NULL; + } if (md->relay_blocks) { queue_free(md->relay_blocks); md->relay_blocks = NULL; } if (md->file_loads) { queue_free(md->file_loads); md->file_loads = NULL; } diff --git a/src/media_delivery/media_delivery.h b/src/media_delivery/media_delivery.h index 56f967de..6fe9337c 100644 --- a/src/media_delivery/media_delivery.h +++ b/src/media_delivery/media_delivery.h @@ -22,6 +22,8 @@ struct UTUN_INSTANCE; #define MEDIA_HAVE_BLOCK_TIMEOUT_TB 20000 // таймаут HAVE_BLOCK_ACK 2s #define MEDIA_RECONNECT_COOLDOWN_TB 360000000 // 1 час (0.1ms) #define MEDIA_HELLO_TIMEOUT_TB 20000 // таймаут SUPER_HELLO 2s +#define MEDIA_DL_STALL_TIMEOUT_TB 20000 // таймаут простоя загрузки 2s (0.1ms) +#define MEDIA_DL_MAX_ATTEMPTS 3 // попыток на блок до ошибки /* ── структуры ── */ @@ -111,6 +113,8 @@ struct media_delivery_ctx { uint8_t max_downloads_per_file; // лимит параллельных загрузок с одного файла (source node), default 3 uint8_t stream_completed; // 1 = хотя бы один стрим завершился (для тестов) uint8_t streams_started; // монотонный счётчик запущенных стримов (для тестов) + int dl_stall_timeout_tb; // таймаут простоя загрузки (по умолч. MEDIA_DL_STALL_TIMEOUT_TB) + int dl_max_attempts; // попыток на блок (по умолч. MEDIA_DL_MAX_ATTEMPTS) 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/src/media_delivery/media_download.c b/src/media_delivery/media_download.c index c046a9d3..eeee831b 100644 --- a/src/media_delivery/media_download.c +++ b/src/media_delivery/media_download.c @@ -1,4 +1,10 @@ // media_download.c — скачивание блоков медиа +// +// Блочный конечный автомат с единым watchdog-таймером: +// - таймаут простоя блока (нет трафика 2с) → фейловер на следующий узел-держатель; +// - исчерпание попыток (3) → провал всей загрузки; +// - перегрузка отправителя (OVERLOADED) → переход на другой узел или ожидание; +// - RELAY_FULL → добавление новых держателей и повторный запрос. #include "media_download.h" #include "media_delivery.h" #include "media_delivery_proto.h" @@ -13,11 +19,11 @@ #include "../lib/mem.h" #include "../lib/u_async.h" #include "../lib/platform_compat.h" +#include "../media_async/media_async.h" #include #include #include #include -#include "../media_async/media_async.h" #define MDL_ID "media_download" @@ -30,18 +36,36 @@ static struct media_download* md_dl_find(struct UTUN_INSTANCE* inst, const uint8 return e ? (struct media_download*)e->data : NULL; } -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; } - if (dl->block_sigs) { u_free(dl->block_sigs); dl->block_sigs = NULL; } +static uint64_t md_dl_stall_tb(struct UTUN_INSTANCE* inst) { + return inst->md.dl_stall_timeout_tb > 0 ? (uint64_t)inst->md.dl_stall_timeout_tb : MEDIA_DL_STALL_TIMEOUT_TB; } -/* ── send helpers ── */ +static int md_dl_max_attempts(struct UTUN_INSTANCE* inst) { + return inst->md.dl_max_attempts > 0 ? inst->md.dl_max_attempts : MEDIA_DL_MAX_ATTEMPTS; +} + +static int md_dl_block_index(struct media_download* dl, const uint8_t* block_id) { + for (int i = 0; i < dl->num_blocks; i++) + if (memcmp(dl->blocks[i].block_id, block_id, 16) == 0) return i; + return -1; +} + +static void md_dl_reset_block_file(struct media_download* dl, int bi) { + char tmp[2048]; + snprintf(tmp, sizeof(tmp), "%s.chunk_%d", dl->dest_path, bi); + FILE* f = fopen(tmp, "wb"); + if (f) fclose(f); +} + +static struct media_download_peer* md_dl_get_peer(struct media_download* dl, uint64_t node_id) { + for (int pi = 0; pi < dl->num_peers; pi++) + if (dl->peers[pi].node_id == node_id) return &dl->peers[pi]; + if (dl->num_peers >= MD_MAX_PEERS) return NULL; + struct media_download_peer* p = &dl->peers[dl->num_peers++]; + memset(p, 0, sizeof(*p)); + p->node_id = node_id; + return p; +} static void md_mkdir_parent(const char* filepath) { char path[1024]; snprintf(path, sizeof(path), "%s", filepath); @@ -59,12 +83,14 @@ static void md_mkdir_parent(const char* filepath) { } } +/* ── send helpers ── */ + static int md_dl_send(struct UTUN_INSTANCE* inst, uint64_t group_id, 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 TOPO_GROUP* grp = topo_groups_find(inst->topo_groups, group_id); struct ETCP_CONN* gc = grp ? topo_group_find_conn_for_node(grp, dst) : NULL; - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: md_dl_send group=0x%016llx dst=0x%016llx grp=%p conn_in_group=%p", - MDL_ID, (unsigned long long)group_id, (unsigned long long)dst, (void*)grp, (void*)gc); + DEBUG_DEBUG(DEBUG_CATEGORY_GENERAL, "%s: md_dl_send group=0x%016llx dst=0x%016llx grp=%p conn_in_group=%p", + MDL_ID, (unsigned long long)group_id, (unsigned long long)dst, (void*)grp, (void*)gc); } 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; } @@ -77,52 +103,26 @@ static int md_dl_send(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t ds return rc; } -/* ── build block request ── */ - -static int md_dl_send_block_req(struct UTUN_INSTANCE* inst, uint64_t dst, - struct media_download* dl, const uint8_t* block_id, - int chunk_idx, int peer_local_bi) { +static int md_dl_send_block_req(struct UTUN_INSTANCE* inst, uint64_t dst, struct media_download* dl, int 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, block_id, 16); - req.chunk = (uint32_t)chunk_idx; + memcpy(req.block_id, dl->blocks[bi].block_id, 16); + req.chunk = (uint32_t)bi; req.offset = 0; - DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%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, peer_local_bi, dl->num_blocks, - block_id[0], block_id[1], block_id[2], block_id[3]); + DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: BLOCK_REQ to 0x%016llx blk=%d/%d block_id=%02x%02x%02x%02x...", + MDL_ID, (unsigned long long)dst, bi, dl->num_blocks, + dl->blocks[bi].block_id[0], dl->blocks[bi].block_id[1], dl->blocks[bi].block_id[2], dl->blocks[bi].block_id[3]); return md_dl_send(inst, dl->group_id, dst, (const uint8_t*)&req, sizeof(req)); } /* ── forward decl ── */ static void md_dl_send_block_processing(struct UTUN_INSTANCE* inst, struct media_download* dl, int bi); - -/* ── send BLOCK_REQ + create relay context + notify supernode ── */ - -static int md_dl_start_block(struct UTUN_INSTANCE* inst, uint64_t dst, - struct media_download* dl, const uint8_t* block_id, - int chunk_idx, int peer_local_bi) { - /* find global block index */ - int gbi = -1; - for (int i = 0; i < dl->num_blocks; i++) { - if (memcmp(dl->block_ids + i * 16, block_id, 16) == 0) { gbi = i; break; } - } - if (gbi >= 0) { - /* create relay context before first chunk arrives */ - char chunk_file[2048]; - snprintf(chunk_file, sizeof(chunk_file), "%s.chunk_%d", dl->dest_path, gbi); - struct relay_block_ctx* rc = md_relay_add(&inst->md, block_id, dl->media_id, (uint32_t)gbi); - if (rc) snprintf(rc->chunk_file, sizeof(rc->chunk_file), "%s", chunk_file); - /* tell supernode we're downloading this block */ - md_dl_send_block_processing(inst, dl, gbi); - } - return md_dl_send_block_req(inst, dst, dl, block_id, chunk_idx, peer_local_bi); -} - -/* ── send HAVE_BLOCK to supernode ── */ +static void md_dl_watchdog_cb(void* arg); +static void md_dl_send_query(struct UTUN_INSTANCE* inst, struct media_download* dl); static void md_dl_send_have_block(struct UTUN_INSTANCE* inst, struct media_download* dl, int bi) { uint64_t super = dl->super_count > 0 ? dl->super_nodes[dl->super_current] : 0; @@ -133,12 +133,12 @@ static void md_dl_send_have_block(struct UTUN_INSTANCE* inst, struct media_downl hb.subcmd = MEDIA_SUBCMD_HAVE_BLOCK; hb.group_id = dl->group_id; memcpy(hb.media_id, dl->media_id, 16); - memcpy(hb.block_id, dl->block_ids + bi * 16, 16); + memcpy(hb.block_id, dl->blocks[bi].block_id, 16); hb.chunk = (uint32_t)bi; hb.timestamp = (int64_t)time(NULL); uint8_t smsg[64]; size_t soff = 0; - memcpy(smsg + soff, dl->block_ids + bi * 16, 16); soff += 16; + memcpy(smsg + soff, dl->blocks[bi].block_id, 16); soff += 16; memcpy(smsg + soff, &hb.chunk, 4); soff += 4; memcpy(smsg + soff, &hb.timestamp, 8); soff += 8; uint64_t self = inst->node_id; memcpy(smsg + soff, &self, 8); soff += 8; @@ -152,8 +152,6 @@ static void md_dl_send_have_block(struct UTUN_INSTANCE* inst, struct media_downl md_dl_send(inst, dl->group_id, super, (const uint8_t*)&hb, sizeof(hb)); } -/* ── send BLOCK_PROCESSING to supernode ── */ - static void md_dl_send_block_processing(struct UTUN_INSTANCE* inst, struct media_download* dl, int bi) { uint64_t super = dl->super_count > 0 ? dl->super_nodes[dl->super_current] : 0; if (!super) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: BLOCK_PROCESSING but no supernode", MDL_ID); return; } @@ -163,12 +161,12 @@ static void md_dl_send_block_processing(struct UTUN_INSTANCE* inst, struct media hb.subcmd = MEDIA_SUBCMD_BLOCK_PROCESSING; hb.group_id = dl->group_id; memcpy(hb.media_id, dl->media_id, 16); - memcpy(hb.block_id, dl->block_ids + bi * 16, 16); + memcpy(hb.block_id, dl->blocks[bi].block_id, 16); hb.chunk = (uint32_t)bi; hb.timestamp = (int64_t)time(NULL); uint8_t smsg[64]; size_t soff = 0; - memcpy(smsg + soff, dl->block_ids + bi * 16, 16); soff += 16; + memcpy(smsg + soff, dl->blocks[bi].block_id, 16); soff += 16; memcpy(smsg + soff, &hb.chunk, 4); soff += 4; memcpy(smsg + soff, &hb.timestamp, 8); soff += 8; uint64_t self = inst->node_id; memcpy(smsg + soff, &self, 8); soff += 8; @@ -182,39 +180,218 @@ static void md_dl_send_block_processing(struct UTUN_INSTANCE* inst, struct media md_dl_send(inst, dl->group_id, super, (const uint8_t*)&hb, sizeof(hb)); } +/* ── запрос блока у текущего держателя ── */ + +static int md_dl_request_block(struct media_download* dl, int bi) { + struct media_download_block* b = &dl->blocks[bi]; + if (b->holder_idx < 0 || b->holder_idx >= b->num_holders) { + b->state = MD_BLK_WAIT; + return -1; + } + uint64_t holder = b->holders[b->holder_idx]; + struct media_download_peer* p = md_dl_get_peer(dl, holder); + if (!p) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: no peer slot for 0x%016llx", MDL_ID, (unsigned long long)holder); b->state = MD_BLK_WAIT; return -1; } + + /* перезапуск блока с начала */ + md_dl_reset_block_file(dl, bi); + b->bytes_received = 0; + b->last_progress_tb = get_time_tb(); + b->state = MD_BLK_CONN; + dl->inflight++; + + struct UTUN_INSTANCE* inst = dl->inst; + if (instance_find_conn(inst, holder)) { + p->connected = 1; + b->state = MD_BLK_REQ; + md_dl_send_block_req(inst, holder, dl, bi); + } else { + struct TOPO_GROUP* grp = inst->topo_groups ? topo_groups_find(inst->topo_groups, dl->group_id) : NULL; + if (grp && grp->conn_mgr) { + int rc = conn_mgr_open(inst, grp->group_id, holder, md_dl_conn_cb, dl, &p->cm_handle); + if (rc != 0) { + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: conn_mgr_open failed rc=%d for 0x%016llx", MDL_ID, rc, (unsigned long long)holder); + dl->inflight--; + return -1; + } + /* CONN — BLOCK_REQ отправится в md_dl_conn_cb по CONN_EVENT_UP */ + } else { + /* нет conn_mgr — шлём напрямую best-effort */ + p->connected = 1; + b->state = MD_BLK_REQ; + md_dl_send_block_req(inst, holder, dl, bi); + } + } + return 0; +} + +/* ── pump: запрашиваем IDLE/WAIT-блоки, у которых есть держатель ── */ + +static void md_dl_pump(struct media_download* dl) { + for (int bi = 0; bi < dl->num_blocks; bi++) { + struct media_download_block* b = &dl->blocks[bi]; + if (b->state == MD_BLK_IDLE || b->state == MD_BLK_WAIT) + md_dl_request_block(dl, bi); + } +} + +/* ── фейловер одного блока на следующий держатель ── */ + +int md_dl_failover_block(struct media_download* dl, int bi, uint64_t now_tb) { + struct media_download_block* b = &dl->blocks[bi]; + if (b->state == MD_BLK_DONE || b->state == MD_BLK_FAILED) return 0; + + if (b->state == MD_BLK_CONN || b->state == MD_BLK_REQ || b->state == MD_BLK_RECV) { + if (dl->inflight > 0) dl->inflight--; + } + b->state = MD_BLK_IDLE; + md_dl_reset_block_file(dl, bi); + b->bytes_received = 0; + b->last_progress_tb = now_tb; + + b->attempts++; + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: FAILOVER blk=%d attempts=%d/%d holders=%d hidx=%d", + MDL_ID, bi, b->attempts, md_dl_max_attempts(dl->inst), b->num_holders, b->holder_idx); + + if (b->attempts >= md_dl_max_attempts(dl->inst)) { + b->state = MD_BLK_FAILED; + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: blk=%d attempts exhausted → FAILED", MDL_ID, bi); + return -1; + } + + b->holder_idx++; + if (b->holder_idx >= b->num_holders) { + b->state = MD_BLK_WAIT; + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: blk=%d no more holders → WAIT (re-query later)", MDL_ID, bi); + return 0; + } + return md_dl_request_block(dl, bi); +} + +/* ── проверка простоя (watchdog): возвращает 0 (активна) / -1 (провал) ── */ + +int md_dl_check_stall(struct media_download* dl, uint64_t now_tb) { + struct UTUN_INSTANCE* inst = dl->inst; + uint64_t stall = md_dl_stall_tb(inst); + + /* зависшие блоки — фейловер */ + for (int bi = 0; bi < dl->num_blocks; bi++) { + struct media_download_block* b = &dl->blocks[bi]; + if (b->state != MD_BLK_CONN && b->state != MD_BLK_REQ && b->state != MD_BLK_RECV) continue; + if (now_tb - b->last_progress_tb >= stall) { + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: STALL blk=%d state=%d no progress %llums", + MDL_ID, bi, b->state, (unsigned long long)(now_tb - b->last_progress_tb)); + if (md_dl_failover_block(dl, bi, now_tb) < 0) return -1; + } + } + + /* WAIT-блоки — ретрай по циклу держателей */ + for (int bi = 0; bi < dl->num_blocks; bi++) { + struct media_download_block* b = &dl->blocks[bi]; + if (b->state != MD_BLK_WAIT) continue; + if (now_tb - b->last_progress_tb >= stall) { + if (b->num_holders > 0) { + b->holder_idx = (b->holder_idx + 1) % b->num_holders; + b->last_progress_tb = now_tb; + md_dl_request_block(dl, bi); + } + } + } + + /* нет вообще никакой активности и блоки не качаются → re-QUERY следующей суперноды */ + if (dl->blocks_validated < dl->num_blocks && now_tb - dl->last_activity_tb >= stall) { + int in_flight = 0; + for (int bi = 0; bi < dl->num_blocks; bi++) { + uint8_t st = dl->blocks[bi].state; + if (st == MD_BLK_CONN || st == MD_BLK_REQ || st == MD_BLK_RECV) in_flight++; + } + if (in_flight == 0) { + dl->super_current++; + if (dl->super_current >= dl->super_count) { + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: all supernodes exhausted (validated=%d/%d)", MDL_ID, dl->blocks_validated, dl->num_blocks); + return -1; + } + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: no activity → re-QUERY super_current=%d/%d", MDL_ID, dl->super_current, dl->super_count); + md_dl_send_query(dl->inst, dl); + dl->last_activity_tb = now_tb; + } + } + return 0; +} + +/* ── завершение загрузки (успех/провал/отмена) — освобождает всё ── */ + +static void md_dl_finish(struct media_download* dl, int err) { + struct UTUN_INSTANCE* inst = dl->inst; + struct media_delivery_ctx* md = &inst->md; + uint8_t* blks = (uint8_t*)dl->blocks; + uint8_t* sigs = dl->block_sigs; + void (*cb)(void*, int) = dl->done_cb; + void* cb_arg = dl->done_arg; + + dl->active = 0; dl->err = err; + if (dl->watchdog_timer) { uasync_cancel_timeout(inst->ua, dl->watchdog_timer); dl->watchdog_timer = NULL; } + 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 (md->downloads) { + struct ll_entry* e = queue_find_data_by_index(md->downloads, dl->media_id); + if (e) { queue_remove_data(md->downloads, e); queue_entry_free(e); } + } + if (blks) u_free(blks); + if (sigs) u_free(sigs); + if (cb) cb(cb_arg, err); +} + +static void md_dl_watchdog_cb(void* arg) { + struct media_download* dl = (struct media_download*)arg; + if (!dl->active) return; + struct UTUN_INSTANCE* inst = dl->inst; + if (md_dl_check_stall(dl, get_time_tb()) < 0) { md_dl_finish(dl, -1); return; } + if (!dl->active) return; + dl->watchdog_timer = uasync_set_timeout(inst->ua, (int)md_dl_stall_tb(inst), dl, md_dl_watchdog_cb, "md_dl_wd"); +} + /* ── conn_mgr callback ── */ 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_MEDIA, "%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) { - DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: connect to 0x%016llx failed rc=%d", - MDL_ID, (unsigned long long)node_id, event); - return; - } - /* 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_MEDIA, "%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++) { - if (!dl->peers[pi].blocks[bi].started) { - dl->peers[pi].blocks[bi].started = 1; - DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: conn_cb BLOCK_REQ: peer-local bi=%d", - MDL_ID, bi); - md_dl_start_block(dl->inst, node_id, dl, dl->peers[pi].blocks[bi].block_id, bi, bi); - break; - } + if (!dl || !dl->active) return; + (void)group_id; + + if (event == CONN_EVENT_UP) { + struct media_download_peer* p = md_dl_get_peer(dl, node_id); + if (p) { p->connected = 1; p->cm_handle = h; } + DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: conn UP node=0x%016llx", MDL_ID, (unsigned long long)node_id); + /* отправить BLOCK_REQ для всех блоков этого узла, ждущих соединения */ + for (int bi = 0; bi < dl->num_blocks; bi++) { + struct media_download_block* b = &dl->blocks[bi]; + if (b->state == MD_BLK_CONN && b->holder_idx >= 0 && b->holder_idx < b->num_holders + && b->holders[b->holder_idx] == node_id) { + b->state = MD_BLK_REQ; + md_dl_send_block_req(dl->inst, node_id, dl, bi); } } + } else { + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: conn FAIL node=0x%016llx event=%d", MDL_ID, (unsigned long long)node_id, event); + struct media_download_peer* p = md_dl_get_peer(dl, node_id); + if (p) { p->connected = 0; p->cm_handle = NULL; } + uint64_t now = get_time_tb(); + int failed = 0; + for (int bi = 0; bi < dl->num_blocks; bi++) { + struct media_download_block* b = &dl->blocks[bi]; + if ((b->state == MD_BLK_CONN || b->state == MD_BLK_REQ || b->state == MD_BLK_RECV) + && b->holder_idx >= 0 && b->holder_idx < b->num_holders + && b->holders[b->holder_idx] == node_id) { + if (md_dl_failover_block(dl, bi, now) < 0) failed = 1; + } + } + if (failed) { md_dl_finish(dl, -1); return; } + md_dl_pump(dl); } } -/* ── query supernode, collect supernodes ── */ +/* ── сбор суперузлов ── */ static void md_dl_collect_supernodes(struct UTUN_INSTANCE* inst, struct media_download* dl) { dl->super_count = 0; @@ -247,13 +424,13 @@ static void md_dl_collect_supernodes(struct UTUN_INSTANCE* inst, struct media_do } } -/* ── send query to current supernode ── */ +/* ── QUERY к текущей суперноде ── */ static void md_dl_send_query(struct UTUN_INSTANCE* inst, struct media_download* dl) { if (dl->super_current >= dl->super_count) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: all supernodes exhausted", MDL_ID); dl->err = -1; - if (dl->done_cb) dl->done_cb(dl->done_arg, dl->err); + md_dl_finish(dl, -1); return; } uint64_t super = dl->super_nodes[dl->super_current]; @@ -267,12 +444,14 @@ static void md_dl_send_query(struct UTUN_INSTANCE* inst, struct media_download* q->group_id = dl->group_id; memcpy(q->media_id, dl->media_id, 16); q->num_blocks = (uint16_t)dl->num_blocks; - memcpy(pkt + MEDIA_QUERY_HDR_SIZE, dl->block_ids, (size_t)dl->num_blocks * 16); + for (int i = 0; i < dl->num_blocks; i++) + memcpy(pkt + MEDIA_QUERY_HDR_SIZE + (size_t)i * 16, dl->blocks[i].block_id, 16); DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: QUERY to super 0x%016llx blocks=%d", MDL_ID, (unsigned long long)super, dl->num_blocks); md_dl_send(inst, dl->group_id, super, pkt, pkt_len); u_free(pkt); + dl->last_activity_tb = get_time_tb(); } /* ── handle QUERY_RESP ── */ @@ -284,82 +463,83 @@ void media_download_handle_query_resp(struct UTUN_INSTANCE* inst, struct media_pkt_query_resp_entry* entries = (struct media_pkt_query_resp_entry*)(data + MEDIA_QUERY_RESP_HDR_SIZE); int ne = r->num_entries; if (ne > 100) ne = 100; + if (ne < 1) { DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: QUERY_RESP empty (0 entries)", MDL_ID); return; } - /* find the download: QUERY_RESP entries contain block_ids; - iterate all downloads checking if any block_id matches dl->block_ids */ - if (ne < 1) return; + /* найти загрузку по совпадению block_id */ struct media_download* dl = NULL; if (inst->md.downloads) { struct ll_entry* qe = inst->md.downloads->head; while (qe) { struct media_download* d = (struct media_download*)qe->data; for (int i = 0; i < ne && !dl; i++) - for (int j = 0; j < d->num_blocks; j++) - if (memcmp(entries[i].block_id, d->block_ids + j * 16, 16) == 0) { dl = d; break; } + if (md_dl_block_index(d, entries[i].block_id) >= 0) { dl = d; break; } if (dl) break; qe = qe->next; } } if (!dl) { DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: QUERY_RESP — no matching download", MDL_ID); return; } + if (!dl->active) { DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: QUERY_RESP for inactive download", MDL_ID); return; } - dl->num_peers = 0; + dl->last_activity_tb = get_time_tb(); - /* group entries by node_id */ - for (int i = 0; i < ne && dl->num_peers < 10; i++) { - uint64_t nid = entries[i].node_id; - int pi = -1; - for (int j = 0; j < dl->num_peers; j++) { - if (dl->peers[j].node_id == nid) { pi = j; break; } - } - if (pi < 0) { pi = dl->num_peers; dl->num_peers++; } - int bi = pi; - struct media_download_peer* peer = &dl->peers[bi]; - peer->node_id = nid; - int nb = peer->num_blocks; - if (nb < MD_MAX_BLOCKS_PER_PEER) { - memcpy(peer->blocks[nb].block_id, entries[i].block_id, 16); - peer->blocks[nb].started = 0; - peer->blocks[nb].received = 0; - peer->blocks[nb].validated = 0; - peer->num_blocks++; + /* сбросить holders у незавершённых блоков (re-QUERY пересобирает список) */ + for (int bi = 0; bi < dl->num_blocks; bi++) { + if (dl->blocks[bi].state != MD_BLK_DONE) { + dl->blocks[bi].num_holders = 0; + dl->blocks[bi].holder_idx = -1; } } - DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: QUERY_RESP entries=%d peers=%d dl=%p num_blocks=%d", - MDL_ID, ne, dl->num_peers, (void*)dl, dl->num_blocks); - for (int pi = 0; pi < dl->num_peers; pi++) { - DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: peer[%d] node=0x%016llx num_blocks=%d", - MDL_ID, pi, (unsigned long long)dl->peers[pi].node_id, dl->peers[pi].num_blocks); + /* заполнить holders из записей */ + for (int i = 0; i < ne; i++) { + int bi = md_dl_block_index(dl, entries[i].block_id); + if (bi < 0) continue; + struct media_download_block* b = &dl->blocks[bi]; + if (b->state == MD_BLK_DONE) continue; + int dup = 0; + for (int h = 0; h < b->num_holders; h++) if (b->holders[h] == entries[i].node_id) { dup = 1; break; } + if (!dup && b->num_holders < MD_MAX_HOLDERS_PER_BLOCK) + b->holders[b->num_holders++] = entries[i].node_id; } - /* if no peers found, try querying the author directly */ - if (dl->num_peers == 0 && dl->author_node_id && dl->author_node_id != inst->node_id) { - DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: no holders from supernode, querying author 0x%016llx", - MDL_ID, (unsigned long long)dl->author_node_id); - dl->super_nodes[0] = dl->author_node_id; - dl->super_count = 1; - dl->super_current = 0; - md_dl_send_query(inst, dl); /* retry QUERY to author */ - return; + /* fallback: блок без держателя — автор */ + for (int bi = 0; bi < dl->num_blocks; bi++) { + struct media_download_block* b = &dl->blocks[bi]; + if (b->state == MD_BLK_DONE) continue; + if (b->num_holders == 0 && dl->author_node_id && dl->author_node_id != inst->node_id) + b->holders[b->num_holders++] = dl->author_node_id; } - /* find the chat group for direct connection (best-effort) */ - struct TOPO_GROUP* grp = topo_groups_find(inst->topo_groups, dl->group_id); + /* пересобрать peers (уникальные node_id всех держателей) */ + dl->num_peers = 0; + for (int bi = 0; bi < dl->num_blocks; bi++) { + struct media_download_block* b = &dl->blocks[bi]; + for (int h = 0; h < b->num_holders; h++) + md_dl_get_peer(dl, b->holders[h]); + } + DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: QUERY_RESP entries=%d peers=%d dl=%p num_blocks=%d", + MDL_ID, ne, dl->num_peers, (void*)dl, dl->num_blocks); for (int pi = 0; pi < dl->num_peers; pi++) { - dl->peers[pi].connected = 1; - DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: peer[%d] route: grp=%p conn_mgr=%p", - MDL_ID, pi, (void*)grp, grp ? (void*)grp->conn_mgr : NULL); - if (grp && grp->conn_mgr) { - int cm_rc = conn_mgr_open(inst, grp->group_id, dl->peers[pi].node_id, md_dl_conn_cb, dl, NULL); - if (cm_rc != 0) DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: conn_mgr_open failed rc=%d for 0x%016llx", MDL_ID, cm_rc, (unsigned long long)dl->peers[pi].node_id); - } else - for (int bi = 0; bi < dl->peers[pi].num_blocks; bi++) { - DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: direct BLOCK_REQ: peer-local bi=%d peer_num_blocks=%d", - MDL_ID, bi, dl->peers[pi].num_blocks); - md_dl_start_block(inst, dl->peers[pi].node_id, dl, dl->peers[pi].blocks[bi].block_id, bi, bi); - } + char blklist[256] = ""; int bl = 0; + for (int bi = 0; bi < dl->num_blocks && bl < (int)sizeof(blklist) - 8; bi++) { + struct media_download_block* b = &dl->blocks[bi]; + for (int h = 0; h < b->num_holders; h++) + if (b->holders[h] == dl->peers[pi].node_id) { bl += snprintf(blklist + bl, sizeof(blklist) - bl, "%s%d", bl ? "," : "", bi); break; } + } + DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: peer[%d] node=0x%016llx blocks=[%s]", + MDL_ID, pi, (unsigned long long)dl->peers[pi].node_id, blklist); } + + /* разброс блоков по держателям: стартовый holder = i % num_holders */ + for (int bi = 0; bi < dl->num_blocks; bi++) { + struct media_download_block* b = &dl->blocks[bi]; + if (b->state == MD_BLK_DONE || b->state == MD_BLK_CONN || b->state == MD_BLK_REQ || b->state == MD_BLK_RECV) continue; + if (b->num_holders > 0) b->holder_idx = bi % b->num_holders; + else b->holder_idx = -1; + } + + md_dl_pump(dl); } /* ── handle incoming BLOCK_CHUNK ── */ @@ -374,18 +554,19 @@ void media_download_handle_chunk(struct UTUN_INSTANCE* inst, if (!dl) { DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: CHUNK for unknown media %02x%02x...", MDL_ID, ch->media_id[0], ch->media_id[1]); return; } if (!dl->active) { DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: CHUNK for inactive download", MDL_ID); return; } - int bi = -1; - for (int i = 0; i < dl->num_blocks; i++) { - if (memcmp(dl->block_ids + i * 16, ch->block_id, 16) == 0) { bi = i; break; } - } + int bi = md_dl_block_index(dl, ch->block_id); if (bi < 0) { DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: CHUNK for unknown block %02x%02x...", MDL_ID, ch->block_id[0], ch->block_id[1]); return; } + struct media_download_block* b = &dl->blocks[bi]; + size_t wlen = ch->data_len; + char tmp[2048]; snprintf(tmp, sizeof(tmp), "%s.chunk_%d", dl->dest_path, bi); - FILE* f = fopen(tmp, "ab"); + FILE* f = fopen(tmp, "rb+"); + if (!f) f = fopen(tmp, "wb"); if (!f) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: fopen(%s) failed for chunk write: %s", MDL_ID, tmp, strerror(errno)); return; } - size_t wlen = ch->data_len; if (len >= MEDIA_BLOCK_CHUNK_HDR_SIZE + wlen) { + fseeko(f, (off_t)ch->offset, SEEK_SET); fwrite(d + MEDIA_BLOCK_CHUNK_HDR_SIZE, 1, wlen, f); } else { DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: CHUNK truncated: len=%zu need=%zu+%zu", @@ -393,6 +574,26 @@ void media_download_handle_chunk(struct UTUN_INSTANCE* inst, } fclose(f); + uint32_t off = ch->offset; + uint32_t end = off + (uint32_t)wlen; + if (end > b->bytes_received) b->bytes_received = end; + dl->last_activity_tb = get_time_tb(); + b->last_progress_tb = dl->last_activity_tb; + if (b->state == MD_BLK_REQ || b->state == MD_BLK_CONN) b->state = MD_BLK_RECV; + + /* rate-limited progress: log at block start, every 1MB, and when block data is complete */ + { + uint32_t mb_off = off >> 20; + uint32_t mb_end = end >> 20; + int64_t blk_size = b->expected_size > 0 ? b->expected_size : dl->block_size; + if (off == 0 || mb_end != mb_off || (blk_size > 0 && end >= (uint32_t)blk_size)) { + DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: CHUNK blk=%d off=%u +%zu =%u/%lld (%.0f%%)%s", + MDL_ID, bi, off, wlen, end, (long long)blk_size, + blk_size > 0 ? 100.0 * (double)end / (double)blk_size : 0.0, + (blk_size > 0 && end >= (uint32_t)blk_size) ? " [data complete]" : ""); + } + } + /* forward to relay downstreams */ struct relay_block_ctx* rc = md_relay_find(&inst->md, ch->block_id); if (rc) { @@ -400,7 +601,6 @@ void media_download_handle_chunk(struct UTUN_INSTANCE* inst, for (int i = 0; i < rc->downstream_count; i++) { struct relay_downstream* ds = &rc->downstream[i]; if (ds->sent_offset >= rc->file_offset) continue; - /* read from chunk file and send catch-up to this downstream */ FILE* rf = fopen(tmp, "rb"); if (!rf) continue; while (ds->sent_offset < rc->file_offset) { @@ -435,18 +635,18 @@ void media_download_handle_done(struct UTUN_INSTANCE* inst, const uint8_t* data, size_t len) { if (!data || len < MEDIA_BLOCK_DONE_SIZE) return; - const uint8_t* d = data; struct media_pkt_block_done* bd = (struct media_pkt_block_done*)data; struct media_download* dl = md_dl_find(inst, bd->media_id); if (!dl) { DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: BLOCK_DONE for unknown media", MDL_ID); return; } if (!dl->active) { DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: BLOCK_DONE for inactive download", MDL_ID); return; } - int bi = -1; - for (int i = 0; i < dl->num_blocks; i++) { - if (memcmp(dl->block_ids + i * 16, bd->block_id, 16) == 0) { bi = i; break; } - } + int bi = md_dl_block_index(dl, bd->block_id); if (bi < 0) { DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: BLOCK_DONE for unknown block", MDL_ID); return; } + struct media_download_block* b = &dl->blocks[bi]; + if (b->state == MD_BLK_DONE) return; + dl->last_activity_tb = get_time_tb(); + /* verify block data: Ed25519 signature from author, or size-only fallback */ int sig_ok = 0; uint8_t author_ed25519_pubkey[32] = {0}; @@ -464,7 +664,7 @@ void media_download_handle_done(struct UTUN_INSTANCE* inst, else { fseeko(f, 0, SEEK_END); off_t fsz = ftello(f); fseeko(f, 0, SEEK_SET); if (fsz == (off_t)bd->total_size && fsz > 0) { - if (author_ed25519_pubkey[0]) { /* Ed25519 verify against author's pubkey */ + if (author_ed25519_pubkey[0]) { uint8_t* buf = u_malloc((size_t)fsz); if (buf) { size_t rd = fread(buf, 1, (size_t)fsz, f); @@ -488,110 +688,141 @@ void media_download_handle_done(struct UTUN_INSTANCE* inst, fclose(f); } - if (sig_ok) { - dl->blocks_received++; - dl->blocks_validated++; - DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: BLOCK_DONE block=%d total=%d/%d", - MDL_ID, bi, dl->blocks_validated, dl->num_blocks); - - if (dl->progress_cb) dl->progress_cb(dl->progress_arg, dl->blocks_validated, dl->num_blocks); - - /* send HAVE_BLOCK to supernode */ - md_dl_send_have_block(inst, dl, bi); - - /* forward BLOCK_DONE to relay downstreams */ - { - struct relay_block_ctx* rc = md_relay_find(&inst->md, dl->block_ids + bi * 16); - if (rc) { - for (int i = 0; i < rc->downstream_count; i++) { - struct media_pkt_block_done rd; - memset(&rd, 0, sizeof(rd)); - rd.subcmd = MEDIA_SUBCMD_BLOCK_DONE; - memcpy(rd.media_id, dl->media_id, 16); - memcpy(rd.block_id, dl->block_ids + bi * 16, 16); - rd.chunk = (uint32_t)bi; - rd.total_size = (uint32_t)rc->file_offset; - md_dl_send(inst, rc->downstream[i].group_id, rc->downstream[i].node_id, (const uint8_t*)&rd, sizeof(rd)); - DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: relay BLOCK_DONE forwarded to 0x%016llx block=%d", - MDL_ID, (unsigned long long)rc->downstream[i].node_id, bi); - } - md_relay_remove(&inst->md, dl->block_ids + bi * 16); - } - } + if (!sig_ok) { + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: BLOCK_DONE invalid blk=%d — failover", MDL_ID, bi); + if (md_dl_failover_block(dl, bi, get_time_tb()) < 0) { md_dl_finish(dl, -1); return; } + md_dl_pump(dl); + return; + } - /* 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_MEDIA, "%s: BLOCK_DONE → next req: peer[%d] bi=%d", - MDL_ID, pi, bk); - md_dl_start_block(inst, dl->peers[pi].node_id, dl, - dl->peers[pi].blocks[bk].block_id, bk, bk); - goto blockreq_sent; - } - } - goto blockreq_sent; + b->state = MD_BLK_DONE; + if (dl->inflight > 0) dl->inflight--; + dl->blocks_validated++; + DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: BLOCK_DONE block=%d total=%d/%d", + MDL_ID, bi, dl->blocks_validated, dl->num_blocks); + + if (dl->progress_cb) dl->progress_cb(dl->progress_arg, dl->blocks_validated, dl->num_blocks); + + md_dl_send_have_block(inst, dl, bi); + + /* forward BLOCK_DONE to relay downstreams */ + { + struct relay_block_ctx* rc = md_relay_find(&inst->md, dl->blocks[bi].block_id); + if (rc) { + for (int i = 0; i < rc->downstream_count; i++) { + struct media_pkt_block_done rd; + memset(&rd, 0, sizeof(rd)); + rd.subcmd = MEDIA_SUBCMD_BLOCK_DONE; + memcpy(rd.media_id, dl->media_id, 16); + memcpy(rd.block_id, dl->blocks[bi].block_id, 16); + rd.chunk = (uint32_t)bi; + rd.total_size = (uint32_t)rc->file_offset; + md_dl_send(inst, rc->downstream[i].group_id, rc->downstream[i].node_id, (const uint8_t*)&rd, sizeof(rd)); } + md_relay_remove(&inst->md, dl->blocks[bi].block_id); } - blockreq_sent: ; - - /* check assembly */ - if (dl->blocks_validated == dl->num_blocks) { - /* assemble file */ - FILE* out = fopen(dl->dest_path, "wb"); - if (!out) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: fopen(%s) failed for assembly: %s", MDL_ID, dl->dest_path, strerror(errno)); } - else { - for (int n = 0; n < dl->num_blocks; n++) { - char tmp2[2048]; - snprintf(tmp2, sizeof(tmp2), "%s.chunk_%d", dl->dest_path, n); - FILE* cf = fopen(tmp2, "rb"); - if (cf) { - uint8_t buf[65536]; - size_t rd; - while ((rd = fread(buf, 1, sizeof(buf), cf)) > 0) fwrite(buf, 1, rd, out); - fclose(cf); - remove(tmp2); - } - } - fclose(out); - } - /* verify content_hash (skip if all-zero — test/backwards compat) */ - { int hash_zero = 1; for (int hi = 0; hi < 32; hi++) if (dl->content_hash[hi]) { hash_zero = 0; break; } - if (!hash_zero) { - uint8_t fhash[32]; ma_sha256_file(dl->dest_path, fhash); - if (memcmp(fhash, dl->content_hash, 32) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: content_hash mismatch for %s", MDL_ID, dl->dest_path); remove(dl->dest_path); dl->blocks_validated--; return; } - } } - 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->blocks_validated == dl->num_blocks) { + FILE* out = fopen(dl->dest_path, "wb"); + if (!out) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: fopen(%s) failed for assembly: %s", MDL_ID, dl->dest_path, strerror(errno)); } + else { + for (int n = 0; n < dl->num_blocks; n++) { + char tmp2[2048]; + snprintf(tmp2, sizeof(tmp2), "%s.chunk_%d", dl->dest_path, n); + FILE* cf = fopen(tmp2, "rb"); + if (cf) { + uint8_t buf[65536]; + size_t rd; + while ((rd = fread(buf, 1, sizeof(buf), cf)) > 0) fwrite(buf, 1, rd, out); + fclose(cf); + remove(tmp2); } } - 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); + fclose(out); } + { int hash_zero = 1; for (int hi = 0; hi < 32; hi++) if (dl->content_hash[hi]) { hash_zero = 0; break; } + if (!hash_zero) { + uint8_t fhash[32]; ma_sha256_file(dl->dest_path, fhash); + if (memcmp(fhash, dl->content_hash, 32) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: content_hash mismatch for %s", MDL_ID, dl->dest_path); remove(dl->dest_path); dl->blocks_validated--; return; } + } } + dl->assembled = 1; + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: file assembled: %s", MDL_ID, dl->dest_path); + md_dl_finish(dl, 0); + return; + } + + /* освободился слот — добрать ожидающие блоки */ + md_dl_pump(dl); +} + +/* ── handle OVERLOADED ── */ + +void media_download_handle_overloaded(struct UTUN_INSTANCE* inst, + const uint8_t* data, size_t len) { + if (!data || len < sizeof(struct media_pkt_block_overloaded)) return; + struct media_pkt_block_overloaded* ov = (struct media_pkt_block_overloaded*)data; + + struct media_download* dl = md_dl_find(inst, ov->media_id); + if (!dl || !dl->active) { DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: OVERLOADED — no active download", MDL_ID); return; } + + int bi = md_dl_block_index(dl, ov->block_id); + if (bi < 0) { DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: OVERLOADED for unknown block", MDL_ID); return; } + + struct media_download_block* b = &dl->blocks[bi]; + if (b->state == MD_BLK_DONE || b->state == MD_BLK_FAILED) return; + + uint32_t delay = ov->retry_after_ms ? ov->retry_after_ms : 2000; + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: OVERLOADED blk=%d retry_after=%ums holders=%d hidx=%d", + MDL_ID, bi, delay, b->num_holders, b->holder_idx); + + if (b->state == MD_BLK_CONN || b->state == MD_BLK_REQ || b->state == MD_BLK_RECV) { + if (dl->inflight > 0) dl->inflight--; + } + b->state = MD_BLK_WAIT; + b->last_progress_tb = get_time_tb(); + b->holder_idx++; + if (b->holder_idx < b->num_holders) { + /* пробуем следующий держатель сразу (перегруз — не провал блока) */ + md_dl_request_block(dl, bi); + } + /* иначе остаёмся WAIT — watchdog/pump доберёт позже */ +} + +/* ── handle RELAY_FULL: добавляем присланные узлы и повторяем запрос ── */ + +void media_download_handle_relay_full(struct UTUN_INSTANCE* inst, + const uint8_t* media_id, const uint8_t* block_id, + uint32_t chunk, const uint64_t* node_ids, int num_nodes) { + struct media_download* dl = md_dl_find(inst, media_id); + if (!dl || !dl->active) { DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: RELAY_FULL — no active download", MDL_ID); return; } + + int bi = md_dl_block_index(dl, block_id); + if (bi < 0) { DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: RELAY_FULL — unknown block", MDL_ID); return; } + (void)chunk; + + struct media_download_block* b = &dl->blocks[bi]; + if (b->state == MD_BLK_DONE || b->state == MD_BLK_FAILED) return; + + if (b->state == MD_BLK_CONN || b->state == MD_BLK_REQ || b->state == MD_BLK_RECV) { + if (dl->inflight > 0) dl->inflight--; + } + b->state = MD_BLK_IDLE; + + /* добавить новых держателей (уникальные) */ + for (int ni = 0; ni < num_nodes && ni < MD_MAX_HOLDERS_PER_BLOCK; ni++) { + int dup = 0; + for (int h = 0; h < b->num_holders; h++) if (b->holders[h] == node_ids[ni]) { dup = 1; break; } + if (!dup && b->num_holders < MD_MAX_HOLDERS_PER_BLOCK) b->holders[b->num_holders++] = node_ids[ni]; + } + /* стартуем с первого нового держателя */ + b->holder_idx = 0; + if (b->num_holders > 0) { + DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: RELAY_FULL blk=%d → %d holders, retry", MDL_ID, bi, b->num_holders); + md_dl_request_block(dl, bi); } else { - /* invalid signature — mark block for retry from another peer */ - DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: BLOCK_DONE sig fail block=%d — retry", MDL_ID, bi); - char tmp[2048]; - snprintf(tmp, sizeof(tmp), "%s.chunk_%d", dl->dest_path, bi); - remove(tmp); - /* reassign to another peer */ - 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) { - dl->peers[pi].blocks[bj].started = 0; - dl->peers[pi].blocks[bj].received = 0; - dl->peers[pi].blocks[bj].validated = 0; - } - } - } + b->state = MD_BLK_WAIT; } } @@ -611,6 +842,12 @@ int media_download_start(struct UTUN_INSTANCE* inst, uint64_t group_id, struct media_delivery_ctx* md = &inst->md; if (!md->initialized) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: start before media_delivery init", MDL_ID); return -1; } + /* дедуп: та же медиа уже качается */ + if (md_dl_find(inst, result->media_id)) { + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: DUPLICATE download media=%02x%02x... ignored", MDL_ID, result->media_id[0], result->media_id[1]); + return 0; + } + struct media_download* dl = u_calloc(1, sizeof(*dl)); if (!dl) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: u_calloc for download failed", MDL_ID); return -1; } @@ -631,20 +868,30 @@ int media_download_start(struct UTUN_INSTANCE* inst, uint64_t group_id, dl->inst = inst; dl->author_node_id = author_node_id; - dl->block_ids = u_malloc((size_t)dl->num_blocks * 16); + dl->blocks = u_calloc((size_t)dl->num_blocks, sizeof(struct media_download_block)); dl->block_sigs = u_malloc((size_t)dl->num_blocks * 64); - if (!dl->block_ids || !dl->block_sigs) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: u_malloc for block_ids/sigs failed", MDL_ID); u_free(dl->block_ids); u_free(dl); return -1; } - memcpy(dl->block_ids, result->block_ids, (size_t)dl->num_blocks * 16); + if (!dl->blocks || !dl->block_sigs) { + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: alloc blocks/sigs failed", MDL_ID); + u_free(dl->blocks); u_free(dl->block_sigs); u_free(dl); return -1; + } + for (int i = 0; i < dl->num_blocks; i++) { + memcpy(dl->blocks[i].block_id, result->block_ids + (size_t)i * 16, 16); + dl->blocks[i].state = MD_BLK_IDLE; + dl->blocks[i].holder_idx = -1; + int64_t start = (int64_t)i * dl->block_size; + int64_t rem = dl->file_size > 0 ? dl->file_size - start : 0; + dl->blocks[i].expected_size = (uint32_t)((rem > 0 && rem < dl->block_size) ? rem : dl->block_size); + } memcpy(dl->block_sigs, result->block_sigs, (size_t)dl->num_blocks * 64); - /* register in downloads queue */ + /* регистрация в очереди */ if (!md->downloads) { md->downloads = queue_new(inst->ua, 256, offsetof(struct media_download, media_id), 16, "md_dl"); - if (!md->downloads) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_new(downloads) failed", MDL_ID); md_dl_free(dl); u_free(dl); return -1; } + if (!md->downloads) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_new(downloads) failed", MDL_ID); u_free(dl->blocks); u_free(dl->block_sigs); u_free(dl); return -1; } } memcpy(dl->ll.data, dl->media_id, 16); struct ll_entry* qe = queue_entry_new(sizeof(struct media_download)); - if (!qe) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_entry_new(download) failed", MDL_ID); md_dl_free(dl); u_free(dl); return -1; } + if (!qe) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_entry_new(download) failed", MDL_ID); u_free(dl->blocks); u_free(dl->block_sigs); u_free(dl); return -1; } memcpy(qe->data, dl, sizeof(*dl)); u_free(dl); queue_data_put_with_index(md->downloads, qe); @@ -652,6 +899,8 @@ int media_download_start(struct UTUN_INSTANCE* inst, uint64_t group_id, dl = (struct media_download*)qe->data; md_dl_collect_supernodes(inst, dl); + dl->last_activity_tb = get_time_tb(); + dl->watchdog_timer = uasync_set_timeout(inst->ua, (int)md_dl_stall_tb(inst), dl, md_dl_watchdog_cb, "md_dl_wd"); md_dl_send_query(inst, dl); DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: download started media=%02x%02x... blocks=%d supers=%d", @@ -659,31 +908,12 @@ int media_download_start(struct UTUN_INSTANCE* inst, uint64_t group_id, return 0; } -/* ── handle OVERLOADED from author ── */ - -void media_download_handle_overloaded(struct UTUN_INSTANCE* inst, - const uint8_t* data, size_t len) { - if (!data || len < sizeof(struct media_pkt_block_overloaded)) return; - struct media_pkt_block_overloaded* ov = (struct media_pkt_block_overloaded*)data; - - struct media_download* dl = md_dl_find(inst, ov->media_id); - if (!dl || !dl->active) return; - - uint32_t delay = ov->retry_after_ms ? ov->retry_after_ms : 2000; - DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: OVERLOADED for block %02x%02x..., retry in %ums", - MDL_ID, ov->block_id[0], ov->block_id[1], delay); - (void)dl; - /* retry: re-enqueue QUERY to supernode after delay */ - /* for now, just log — full retry comes with timer integration */ -} - int media_download_cancel(struct UTUN_INSTANCE* inst, const uint8_t* media_id, const uint8_t* block_id, uint32_t chunk) { (void)chunk; struct media_download* dl = md_dl_find(inst, media_id); if (!dl) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: cancel — download not found", MDL_ID); return -1; } - /* send CANCEL to all connected peers */ struct media_pkt_cancel can; memset(&can, 0, sizeof(can)); can.subcmd = MEDIA_SUBCMD_CANCEL; @@ -691,79 +921,25 @@ int media_download_cancel(struct UTUN_INSTANCE* inst, if (block_id) memcpy(can.block_id, block_id, 16); can.chunk = chunk; - for (int pi = 0; pi < dl->num_peers; pi++) { - if (dl->peers[pi].connected) { + for (int pi = 0; pi < dl->num_peers; pi++) + if (dl->peers[pi].connected) md_dl_send(inst, dl->group_id, dl->peers[pi].node_id, (const uint8_t*)&can, sizeof(can)); - } - } - - dl->active = 0; - dl->err = -2; - if (dl->done_cb) dl->done_cb(dl->done_arg, dl->err); DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: download cancelled", MDL_ID); + md_dl_finish(dl, -2); return 0; } -/* ── handle RELAY_FULL from relay node: retry block request from one of the listed nodes ── */ +/* ── дамп состояния (для отладки/тестов) ── */ -void media_download_handle_relay_full(struct UTUN_INSTANCE* inst, - const uint8_t* media_id, const uint8_t* block_id, - uint32_t chunk, const uint64_t* node_ids, int num_nodes) { - struct media_download* dl = md_dl_find(inst, media_id); - if (!dl || !dl->active) { DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: RELAY_FULL — no active download", MDL_ID); return; } - - int bi = -1; - for (int i = 0; i < dl->num_blocks; i++) { - if (memcmp(dl->block_ids + i * 16, block_id, 16) == 0) { bi = i; break; } - } - if (bi < 0) { DEBUG_WARN(DEBUG_CATEGORY_MEDIA, "%s: RELAY_FULL — unknown block", MDL_ID); return; } - - /* try each listed node — if already a peer, retry block; otherwise add new */ - for (int ni = 0; ni < num_nodes && ni < MD_MAX_RELAY_DOWNSTREAM; ni++) { - int found_pi = -1; - for (int j = 0; j < dl->num_peers; j++) { - if (dl->peers[j].node_id == node_ids[ni]) { found_pi = j; break; } - } - if (found_pi >= 0) { - /* already a peer — reset block started flag and retry */ - for (int bj = 0; bj < dl->peers[found_pi].num_blocks; bj++) { - if (memcmp(dl->peers[found_pi].blocks[bj].block_id, block_id, 16) == 0) { - dl->peers[found_pi].blocks[bj].started = 0; - dl->peers[found_pi].blocks[bj].received = 0; - DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: RELAY_FULL → retry existing peer[%d] node=0x%016llx bi=%d", - MDL_ID, found_pi, (unsigned long long)node_ids[ni], bj); - md_dl_start_block(inst, node_ids[ni], dl, block_id, (int)chunk, bj); - return; - } - } - /* peer doesn't have this block — add it */ - int nb = dl->peers[found_pi].num_blocks; - if (nb >= MD_MAX_BLOCKS_PER_PEER) continue; - memcpy(dl->peers[found_pi].blocks[nb].block_id, block_id, 16); - dl->peers[found_pi].blocks[nb].started = 0; - dl->peers[found_pi].blocks[nb].received = 0; - dl->peers[found_pi].blocks[nb].validated = 0; - dl->peers[found_pi].num_blocks++; - DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: RELAY_FULL → add block to existing peer[%d] bi=%d", - MDL_ID, found_pi, nb); - md_dl_start_block(inst, node_ids[ni], dl, block_id, (int)chunk, nb); - return; - } - /* new peer */ - int pi = dl->num_peers; - if (pi >= 10) break; - dl->peers[pi].node_id = node_ids[ni]; - dl->peers[pi].connected = 1; - dl->peers[pi].num_blocks = 1; - memcpy(dl->peers[pi].blocks[0].block_id, block_id, 16); - dl->peers[pi].blocks[0].started = 0; - dl->peers[pi].blocks[0].received = 0; - dl->peers[pi].blocks[0].validated = 0; - dl->num_peers++; - DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: RELAY_FULL → new peer[%d] node=0x%016llx", - MDL_ID, pi, (unsigned long long)node_ids[ni]); - md_dl_start_block(inst, node_ids[ni], dl, block_id, (int)chunk, 0); - break; +void md_dl_dump(struct media_download* dl) { + if (!dl) return; + DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: dump media=%02x%02x... blocks=%d validated=%d inflight=%d active=%d attempts_budget=%d", + MDL_ID, dl->media_id[0], dl->media_id[1], dl->num_blocks, dl->blocks_validated, dl->inflight, dl->active, + dl->inst ? md_dl_max_attempts(dl->inst) : 0); + for (int bi = 0; bi < dl->num_blocks; bi++) { + struct media_download_block* b = &dl->blocks[bi]; + DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: blk=%d state=%d holders=%d hidx=%d attempts=%d bytes=%u/%u", + MDL_ID, bi, b->state, b->num_holders, b->holder_idx, b->attempts, b->bytes_received, b->expected_size); } } diff --git a/src/media_delivery/media_download.h b/src/media_delivery/media_download.h index b06614a4..31b6b7f7 100644 --- a/src/media_delivery/media_download.h +++ b/src/media_delivery/media_download.h @@ -13,27 +13,45 @@ extern "C" { struct UTUN_INSTANCE; struct media_index_result; +struct CONN_MGR_HANDLE; -struct UTUN_INSTANCE; -struct media_index_result; +/* максимальное число узлов-держателей одного блока (кандидатов на фейловер) */ +#define MD_MAX_HOLDERS_PER_BLOCK 4 +/* максимальное число уникальных узлов-держателей в загрузке */ +#define MD_MAX_PEERS 10 + +/* ── состояние блока в конечном автомате загрузки ── */ +enum md_block_state { + MD_BLK_IDLE = 0, /* ещё не запрошен */ + MD_BLK_CONN, /* conn_mgr_open в процессе, BLOCK_REQ ещё не отправлен */ + MD_BLK_REQ, /* BLOCK_REQ отправлен, ждём первый чанк */ + MD_BLK_RECV, /* чанки идут */ + MD_BLK_WAIT, /* перегружен/нет свободного держателя — ретрай позже */ + MD_BLK_DONE, /* BLOCK_DONE + сигнатура валидна */ + MD_BLK_FAILED /* попытки исчерпаны */ +}; -/* максимальное число параллельных блоков на пира (один стрим на блок) */ -#define MD_MAX_BLOCKS_PER_PEER 64 +/* ── состояние одного блока ── */ +struct media_download_block { + uint8_t block_id[16]; + uint64_t holders[MD_MAX_HOLDERS_PER_BLOCK]; /* узлы, у кого есть блок */ + int num_holders; + int holder_idx; /* текущий держатель (-1 = нет) */ + int attempts; /* 0..max (фейловеры/провалы блока) */ + uint32_t bytes_received; /* max(offset+len) полученных байт */ + uint32_t expected_size; /* реальный размер блока (последний короче) */ + uint8_t state; /* enum md_block_state */ + uint64_t last_progress_tb; /* get_time_tb() последнего чанка/прогресса */ +}; +/* ── соединение с узлом-держателем ── */ struct media_download_peer { uint64_t node_id; uint8_t connected; - void* cm_handle; - int num_blocks; - struct { - uint8_t block_id[16]; - uint8_t started; - uint8_t received; - uint8_t validated; - } blocks[MD_MAX_BLOCKS_PER_PEER]; + struct CONN_MGR_HANDLE* cm_handle; }; -/* внутреннее состояние загрузки (доступно для тестов) */ +/* ── внутреннее состояние загрузки (доступно для тестов) ── */ struct media_download { struct ll_entry ll; uint8_t media_id[16]; @@ -41,24 +59,24 @@ struct media_download { char dest_path[1024]; char media_base[512]; int num_blocks; - uint8_t* block_ids; - uint8_t* block_sigs; + uint8_t* block_sigs; /* num_blocks * 64 (Ed25519 подписи блоков) */ uint8_t content_hash[32]; int64_t file_size; int64_t block_size; - int blocks_received; + struct media_download_block* blocks; /* num_blocks */ int blocks_validated; - int num_peers; - struct media_download_peer peers[10]; + int inflight; /* блоков в CONN/REQ/RECV */ uint8_t active; uint8_t assembled; int err; uint64_t super_nodes[10]; int super_count; int super_current; - uint64_t author_node_id; /* src_node_id from chat message, for fallback when supernode has no info */ - void* query_timer; - void* timeout_timer; + uint64_t author_node_id; /* fallback, когда у суперноды нет инфо */ + void* watchdog_timer; /* таймер простоя (2с) */ + uint64_t last_activity_tb; /* любая входящая активность (chunk/done/query_resp) */ + int num_peers; + struct media_download_peer peers[MD_MAX_PEERS]; void (*done_cb)(void* arg, int err); void* done_arg; void (*progress_cb)(void* arg, int blocks_done, int num_blocks); @@ -66,21 +84,6 @@ struct media_download { struct UTUN_INSTANCE* inst; }; -/* состояние одного стрима на стороне блок-холдера */ -struct media_block_stream { - uint8_t media_id[16]; - uint8_t block_id[16]; - uint32_t chunk; - uint64_t offset; // текущее смещение в блоке - uint32_t bytes_sent; // всего отправлено байт - uint32_t total_size; // ожидаемый размер блока - uint8_t done; // 1 = стрим завершён - void* waiter; // handle etcp_router_on_send_ready - FILE* file; // fd исходного файла - uint64_t dst_node_id; - uint64_t group_id; -}; - int media_download_start(struct UTUN_INSTANCE* inst, uint64_t group_id, const struct media_index_result* result, const char* dest_path, const char* media_base, @@ -92,22 +95,17 @@ int media_download_cancel(struct UTUN_INSTANCE* inst, const uint8_t* media_id, const uint8_t* block_id, uint32_t chunk); /* ── handlers for incoming data, called from md_etcp_recv dispatch ── */ - -struct ll_entry; -struct UTUN_INSTANCE; - void media_download_handle_query_resp(struct UTUN_INSTANCE* inst, const uint8_t* data, size_t len); void media_download_handle_chunk(struct UTUN_INSTANCE* inst, - const uint8_t* data, size_t len); -void media_download_handle_done(struct UTUN_INSTANCE* inst, const uint8_t* data, size_t len); +void media_download_handle_done(struct UTUN_INSTANCE* inst, + const uint8_t* data, size_t len); void media_download_handle_overloaded(struct UTUN_INSTANCE* inst, - const uint8_t* data, size_t len); - + const uint8_t* data, size_t len); void media_download_handle_relay_full(struct UTUN_INSTANCE* inst, - const uint8_t* media_id, const uint8_t* block_id, - uint32_t chunk, const uint64_t* node_ids, int num_nodes); + const uint8_t* media_id, const uint8_t* block_id, + uint32_t chunk, const uint64_t* node_ids, int num_nodes); /* ── conn_mgr callback (non-static for test visibility) ── */ struct CONN_MGR_HANDLE; @@ -115,6 +113,14 @@ 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); +/* ── тестовые хуки (для детерминированных юнит-тестов с «фейковым» временем) ── + * md_dl_check_stall — проверка простоя: возвращает 0 (активна) / -1 (провал, нужно finish). + * md_dl_failover_block — фейловер одного блока на следующий держатель: 0 / -1 (попытки исчерпаны). + * md_dl_dump — дамп состояния загрузки в лог (DEBUG_CATEGORY_MEDIA). */ +int md_dl_check_stall(struct media_download* dl, uint64_t now_tb); +int md_dl_failover_block(struct media_download* dl, int bi, uint64_t now_tb); +void md_dl_dump(struct media_download* dl); + #ifdef __cplusplus } #endif diff --git a/tests/Makefile.am b/tests/Makefile.am index 7c3ff59b..a1a4f70c 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -71,6 +71,7 @@ check_PROGRAMS = \ test_media_index \ test_media_delivery_sql \ test_media_delivery_download \ + test_media_download_timeout \ test_media_delivery_integration \ test_media_delivery_full \ test_etcp_link_stress \ @@ -379,6 +380,10 @@ test_media_delivery_download_SOURCES = test_media_delivery_download.c test_media_delivery_download_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/src/media_delivery -I$(top_srcdir)/src/media_async -I$(top_srcdir)/lib -I$(top_srcdir)/src/transport_layer -I$(top_srcdir)/src/routing_layer -I$(top_srcdir)/src/chat test_media_delivery_download_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_media_download_timeout_SOURCES = test_media_download_timeout.c +test_media_download_timeout_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/src/media_delivery -I$(top_srcdir)/src/media_async -I$(top_srcdir)/lib -I$(top_srcdir)/src/transport_layer -I$(top_srcdir)/src/routing_layer -I$(top_srcdir)/src/chat +test_media_download_timeout_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) + test_media_delivery_integration_SOURCES = test_media_delivery_integration.c test_media_delivery_integration_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/src/media_delivery -I$(top_srcdir)/src/media_async -I$(top_srcdir)/lib -I$(top_srcdir)/src/transport_layer -I$(top_srcdir)/src/routing_layer -I$(top_srcdir)/src/chat test_media_delivery_integration_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) diff --git a/tests/test_media_delivery_download.c b/tests/test_media_delivery_download.c index 4402556f..cba8e996 100644 --- a/tests/test_media_delivery_download.c +++ b/tests/test_media_delivery_download.c @@ -1,13 +1,16 @@ -// test_media_delivery_download.c — тест скачивания блоков (chunk → done → assembly) +// test_media_delivery_download.c — юнит-тесты скачивания блоков (блочный конечный автомат) // // Покрытие: -// 1. media_download_start — регистрация в очереди, сбор суперузлов -// 2. media_download_handle_chunk — запись чанка в .chunk_N -// 3. media_download_handle_chunk — неизвестный media_id → игнор -// 4. media_download_handle_done — проверка подписи, отметка validated -// 5. media_download_handle_done — неверная подпись → сброс started/received -// 6. media_download_handle_done — все блоки получены → сборка файла -// 7. media_download_cancel — отправка CANCEL, done_cb с err=-2 +// - запись чанков по offset (идемпотентность при дублях) +// - BLOCK_DONE: сигнатура валидна → сборка; неверная → фейловер +// - фейловер блока по простою (фейковое время) +// - исчерпание попыток → FAILED +// - проверка простоя md_dl_check_stall (только зависшие блоки) +// - OVERLOADED → следующий держатель / WAIT +// - RELAY_FULL → добавление держателей +// - пустой QUERY_RESP +// - назначение держателей + разброс по узлам +// - дедуп повторного media_download_start #include "media_delivery.h" #include "media_delivery_proto.h" @@ -36,7 +39,7 @@ static int g_passed = 0, g_failed = 0, g_total = 0; static char g_temp_dir[256]; static struct UASYNC* g_ua = NULL; -#define TEST(name) do { g_total++; printf(" %-55s", name); } while(0) +#define TEST(name) do { g_total++; printf(" %-58s", name); fflush(stdout); } while(0) #define OK() do { g_passed++; printf("OK\n"); } while(0) #define FAIL(fmt, ...) do { g_failed++; printf("FAIL: " fmt "\n", ##__VA_ARGS__); } while(0) @@ -57,14 +60,9 @@ static int file_size(const char* path) { static void test_setup(void) { snprintf(g_temp_dir, sizeof(g_temp_dir), "%s", TEMP_DIR); -#ifdef _WIN32 - { char tmp_path[512]; GetTempPathA(sizeof(tmp_path), tmp_path); - snprintf(g_temp_dir, sizeof(g_temp_dir), "%s\\utun_test_%08x", tmp_path, (unsigned)rand()); - _mkdir(g_temp_dir); } -#else if (!mkdtemp(g_temp_dir)) { fprintf(stderr, "mkdtemp failed\n"); exit(1); } -#endif debug_config_init(); debug_set_level(DEBUG_LEVEL_ERROR); + if (getenv("UTUN_TEST_DEBUG")) { debug_set_category_level_by_name("media", "trace"); debug_set_category_level_by_name("general", "info"); } g_ua = uasync_create(); } @@ -73,7 +71,7 @@ static void test_cleanup(void) { if (g_ua) { uasync_destroy(g_ua, 0); g_ua = NULL; } } -/* ── create minimal UTUN_INSTANCE with media_delivery ── */ +/* ── минимальный UTUN_INSTANCE с media_delivery ── */ #include "../src/utun_instance.h" @@ -83,413 +81,398 @@ static struct UTUN_INSTANCE* make_minimal_instance(sqlite3* db) { inst->ua = g_ua; inst->node_id = 0xDEADBEEFDEADBEEFULL; inst->topo_sqlite_db = db; - - /* generate Ed25519 key pair — just random bytes for test */ for (int i = 0; i < 32; i++) inst->my_ed25519_privkey[i] = (uint8_t)(rand() & 0xFF); memset(inst->my_ed25519_pubkey, 0xDD, 32); - /* init media_delivery */ inst->md.inst = inst; inst->md.db = db; inst->md.self_node_id = inst->node_id; inst->md.initialized = 1; - + inst->md.dl_stall_timeout_tb = 20000; /* 2s */ + inst->md.dl_max_attempts = 3; return inst; } -/* ── test cases ── */ - -static void test_download_chunk(void) { - TEST("download handle_chunk writes to temp file"); { - sqlite3* db = NULL; sqlite3_open(":memory:", &db); - struct UTUN_INSTANCE* inst = make_minimal_instance(db); - - inst->md.downloads = queue_new(inst->ua, 256, offsetof(struct media_download, media_id), 16, "md_dl_test"); - struct media_download dl_buf; memset(&dl_buf, 0, sizeof(dl_buf)); - make_uuid(dl_buf.media_id); - dl_buf.num_blocks = NUM_BLOCKS; - snprintf(dl_buf.dest_path, sizeof(dl_buf.dest_path), "%s/test_output.bin", g_temp_dir); - dl_buf.block_ids = u_malloc(NUM_BLOCKS * 16); - dl_buf.block_sigs = u_malloc(NUM_BLOCKS * 64); - for (int i = 0; i < NUM_BLOCKS; i++) make_uuid(dl_buf.block_ids + i * 16); - dl_buf.active = 1; - memcpy(dl_buf.ll.data, dl_buf.media_id, 16); - - struct ll_entry* qe = queue_entry_new(sizeof(struct media_download)); - if (qe) { memcpy(qe->data, &dl_buf, sizeof(dl_buf)); queue_data_put_with_index(inst->md.downloads, qe); } - - uint8_t chunk_data[256]; - for (int i = 0; i < 256; i++) chunk_data[i] = (uint8_t)(i & 0xFF); - - uint8_t pkt[MEDIA_BLOCK_CHUNK_HDR_SIZE + sizeof(chunk_data)]; - struct media_pkt_block_chunk* ch = (struct media_pkt_block_chunk*)pkt; - ch->subcmd = MEDIA_SUBCMD_BLOCK_CHUNK; - memcpy(ch->media_id, dl_buf.media_id, 16); - memcpy(ch->block_id, dl_buf.block_ids, 16); - ch->chunk = 0; ch->offset = 0; ch->data_len = (uint16_t)sizeof(chunk_data); - memcpy(pkt + MEDIA_BLOCK_CHUNK_HDR_SIZE, chunk_data, sizeof(chunk_data)); +/* создаёт загрузку и регистрирует в очереди; возвращает указатель на данные в очереди */ +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"); - media_download_handle_chunk(inst, pkt, sizeof(pkt)); + 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->blocks = u_calloc((size_t)num_blocks, sizeof(struct media_download_block)); + dl->block_sigs = u_calloc((size_t)num_blocks, 64); + for (int i = 0; i < num_blocks; i++) { + make_uuid(dl->blocks[i].block_id); + memset(dl->block_sigs + i * 64, (uint8_t)(0x42 + i), 64); + dl->blocks[i].state = MD_BLK_IDLE; + dl->blocks[i].holder_idx = -1; + dl->blocks[i].expected_size = CHUNK_SIZE; + } + dl->block_size = CHUNK_SIZE; + dl->file_size = (int64_t)CHUNK_SIZE * num_blocks; + dl->active = 1; + dl->inst = inst; + dl->author_node_id = 0xAAAA000000000001ULL; + memcpy(dl->ll.data, dl->media_id, 16); - char tmp[1024]; snprintf(tmp, sizeof(tmp), "%s.chunk_0", dl_buf.dest_path); - int ok1 = file_exists(tmp) && file_size(tmp) == (int)sizeof(chunk_data); + struct ll_entry* qe = queue_entry_new(sizeof(*dl)); + if (!qe) { u_free(dl->blocks); 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; +} - /* send second chunk with same block_id → appends */ - ch->offset = (uint32_t)sizeof(chunk_data); - media_download_handle_chunk(inst, pkt, sizeof(pkt)); - int ok2 = file_size(tmp) == 2 * (int)sizeof(chunk_data); +static void free_test_dl(struct UTUN_INSTANCE* inst, struct media_download* dl) { + if (!dl) return; + if (dl->watchdog_timer) { uasync_cancel_timeout(inst->ua, dl->watchdog_timer); dl->watchdog_timer = NULL; } + u_free(dl->blocks); + u_free(dl->block_sigs); + if (inst->md.downloads) { + struct ll_entry* e = queue_find_data_by_index(inst->md.downloads, dl->media_id); + if (e) { queue_remove_data(inst->md.downloads, e); queue_entry_free(e); } + } +} - if (ok1 && ok2) OK(); else FAIL("chunk: ok1=%d ok2=%d size=%d", ok1, ok2, file_size(tmp)); +static void add_holder(struct media_download* dl, int bi, uint64_t node_id) { + struct media_download_block* b = &dl->blocks[bi]; + if (b->num_holders < MD_MAX_HOLDERS_PER_BLOCK) + b->holders[b->num_holders++] = node_id; +} - u_free(dl_buf.block_ids); u_free(dl_buf.block_sigs); - if (inst->md.downloads) queue_free(inst->md.downloads); - u_free(inst); sqlite3_close(db); - } +/* ── 1. запись чанков по offset ── */ - TEST("download handle_chunk unknown media_id"); { +static void test_chunk_offset_write(void) { + TEST("chunk write by offset (dup/idempotent)"); { sqlite3* db = NULL; sqlite3_open(":memory:", &db); struct UTUN_INSTANCE* inst = make_minimal_instance(db); - inst->md.downloads = queue_new(inst->ua, 256, offsetof(struct media_download, media_id), 16, "md_dl_test2"); + struct media_download* dl = make_test_dl(db, inst, 1); + if (!dl) { FAIL("make_test_dl"); sqlite3_close(db); u_free(inst); return; } - uint8_t pkt[MEDIA_BLOCK_CHUNK_HDR_SIZE]; + uint8_t cdata[256]; for (int i = 0; i < 256; i++) cdata[i] = (uint8_t)(i & 0xFF); + + uint8_t pkt[MEDIA_BLOCK_CHUNK_HDR_SIZE + sizeof(cdata)]; struct media_pkt_block_chunk* ch = (struct media_pkt_block_chunk*)pkt; + memset(ch, 0, MEDIA_BLOCK_CHUNK_HDR_SIZE); ch->subcmd = MEDIA_SUBCMD_BLOCK_CHUNK; - make_uuid(ch->media_id); make_uuid(ch->block_id); - ch->chunk = 0; ch->offset = 0; ch->data_len = 0; + memcpy(ch->media_id, dl->media_id, 16); + memcpy(ch->block_id, dl->blocks[0].block_id, 16); + ch->chunk = 0; ch->offset = 0; ch->data_len = (uint16_t)sizeof(cdata); + memcpy(pkt + MEDIA_BLOCK_CHUNK_HDR_SIZE, cdata, sizeof(cdata)); + + media_download_handle_chunk(inst, pkt, sizeof(pkt)); + char tmp[1024]; snprintf(tmp, sizeof(tmp), "%s.chunk_0", dl->dest_path); + int s1 = file_size(tmp); + + /* дубль с тем же offset — не должен добавить данные */ + media_download_handle_chunk(inst, pkt, sizeof(pkt)); + int s2 = file_size(tmp); + /* чанк со смещением 200 — пишется по offset, размер растёт до 200+256 */ + ch->offset = 200; media_download_handle_chunk(inst, pkt, sizeof(pkt)); - OK(); /* should not crash or create files */ + int s3 = file_size(tmp); + + int ok = s1 == 256 && s2 == 256 && s3 == 456 && dl->blocks[0].bytes_received == 456; + if (ok) OK(); else FAIL("s1=%d s2=%d s3=%d bytes=%u", s1, s2, s3, dl->blocks[0].bytes_received); - if (inst->md.downloads) queue_free(inst->md.downloads); - u_free(inst); sqlite3_close(db); + free_test_dl(inst, dl); + sqlite3_close(db); u_free(inst); } } -static void test_download_done(void) { - TEST("download handle_done sig ok + assembly"); { +/* ── 2. BLOCK_DONE: сборка ── */ + +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); - inst->md.downloads = queue_new(inst->ua, 256, offsetof(struct media_download, media_id), 16, "md_dl_test3"); - - struct media_download dl; memset(&dl, 0, sizeof(dl)); - make_uuid(dl.media_id); - dl.num_blocks = 2; - snprintf(dl.dest_path, sizeof(dl.dest_path), "%s/out.bin", g_temp_dir); - dl.block_ids = u_malloc(16 * 2); dl.block_sigs = u_malloc(64 * 2); - make_uuid(dl.block_ids); make_uuid(dl.block_ids + 16); - memset(dl.block_sigs, 0x42, 64); memset(dl.block_sigs + 64, 0x43, 64); - dl.block_size = CHUNK_SIZE; dl.file_size = CHUNK_SIZE * 2; - dl.active = 1; - memcpy(dl.ll.data, dl.media_id, 16); - - struct ll_entry* qe = queue_entry_new(sizeof(struct media_download)); - if (!qe) { FAIL("queue_entry_new"); goto done3; } - memcpy(qe->data, &dl, sizeof(dl)); queue_data_put_with_index(inst->md.downloads, qe); - - /* write chunk files */ - uint8_t data0[CHUNK_SIZE]; memset(data0, 0xA0, CHUNK_SIZE); - uint8_t data1[CHUNK_SIZE]; memset(data1, 0xB1, CHUNK_SIZE); - { - char c0[1024]; snprintf(c0, sizeof(c0), "%s.chunk_0", dl.dest_path); - char c1[1024]; snprintf(c1, sizeof(c1), "%s.chunk_1", dl.dest_path); - write_file(c0, data0, CHUNK_SIZE); write_file(c1, data1, CHUNK_SIZE); - } + struct media_download* dl = make_test_dl(db, inst, 2); + if (!dl) { FAIL("make_test_dl"); sqlite3_close(db); u_free(inst); return; } + 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); } + char dest[1024]; snprintf(dest, sizeof(dest), "%s", dl->dest_path); - /* BLOCK_DONE for block 0 */ - { - uint8_t pkt[MEDIA_BLOCK_DONE_SIZE]; - struct media_pkt_block_done* bd = (struct media_pkt_block_done*)pkt; - memset(pkt, 0, sizeof(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); /* match expected */ - media_download_handle_done(inst, pkt, sizeof(pkt)); - } - /* BLOCK_DONE for block 1 */ - { - uint8_t pkt[MEDIA_BLOCK_DONE_SIZE]; + 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; - memset(pkt, 0, sizeof(pkt)); bd->subcmd = MEDIA_SUBCMD_BLOCK_DONE; - memcpy(bd->media_id, dl.media_id, 16); - memcpy(bd->block_id, dl.block_ids + 16, 16); - bd->chunk = 1; bd->total_size = CHUNK_SIZE; - memcpy(bd->block_sig, dl.block_sigs + 64, 64); + memcpy(bd->media_id, dl->media_id, 16); memcpy(bd->block_id, dl->blocks[bi].block_id, 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)); } - /* verify assembly */ - if (file_exists(dl.dest_path)) { - int sz = file_size(dl.dest_path); - if (sz == CHUNK_SIZE * 2) OK(); else FAIL("assembled size %d != %d", sz, CHUNK_SIZE * 2); - } else { FAIL("assembled file not created"); } - -done3: - 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 sz = file_size(dest); + if (sz == CHUNK_SIZE * 2) OK(); else FAIL("assembled size %d != %d", sz, CHUNK_SIZE * 2); + /* dl уже освобождён через md_dl_finish */ + sqlite3_close(db); u_free(inst); } +} - TEST("download handle_done bad sig → retry"); { +/* ── 3. BLOCK_DONE: неверная сигнатура → фейловер ── */ + +static void test_done_bad_sig(void) { + TEST("BLOCK_DONE bad sig → failover (no assembly)"); { sqlite3* db = NULL; sqlite3_open(":memory:", &db); struct UTUN_INSTANCE* inst = make_minimal_instance(db); - inst->md.downloads = queue_new(inst->ua, 256, offsetof(struct media_download, media_id), 16, "md_dl_test4"); - - struct media_download dl; memset(&dl, 0, sizeof(dl)); - make_uuid(dl.media_id); - dl.num_blocks = 1; - snprintf(dl.dest_path, sizeof(dl.dest_path), "%s/bad_sig.bin", g_temp_dir); - dl.block_ids = u_malloc(16); dl.block_sigs = u_malloc(64); - make_uuid(dl.block_ids); - memset(dl.block_sigs, 0x55, 64); /* expected sig */ - dl.block_size = CHUNK_SIZE; dl.file_size = CHUNK_SIZE; - dl.active = 1; - memcpy(dl.ll.data, dl.media_id, 16); - - struct ll_entry* qe = queue_entry_new(sizeof(struct media_download)); - if (!qe) { FAIL("queue_entry_new"); goto done4; } - memcpy(qe->data, &dl, sizeof(dl)); queue_data_put_with_index(inst->md.downloads, qe); - - /* write chunk file */ + struct media_download* dl = make_test_dl(db, inst, 1); + if (!dl) { FAIL("make_test_dl"); sqlite3_close(db); u_free(inst); return; } + add_holder(dl, 0, 0x1111); + dl->blocks[0].holder_idx = 0; + uint8_t data[CHUNK_SIZE]; memset(data, 0xFF, CHUNK_SIZE); - char c0[1024]; snprintf(c0, sizeof(c0), "%s.chunk_0", dl.dest_path); - write_file(c0, data, CHUNK_SIZE); + { char c[1024]; snprintf(c, sizeof(c), "%s.chunk_0", dl->dest_path); write_file(c, data, CHUNK_SIZE); } - /* send BLOCK_DONE with WRONG signature */ - uint8_t pkt[MEDIA_BLOCK_DONE_SIZE]; + uint8_t pkt[MEDIA_BLOCK_DONE_SIZE]; memset(pkt, 0, sizeof(pkt)); struct media_pkt_block_done* bd = (struct media_pkt_block_done*)pkt; - memset(pkt, 0, sizeof(pkt)); bd->subcmd = MEDIA_SUBCMD_BLOCK_DONE; - memcpy(bd->media_id, dl.media_id, 16); - memcpy(bd->block_id, dl.block_ids, 16); + memcpy(bd->media_id, dl->media_id, 16); memcpy(bd->block_id, dl->blocks[0].block_id, 16); bd->chunk = 0; bd->total_size = CHUNK_SIZE; memset(bd->block_sig, 0xAA, 64); /* WRONG */ media_download_handle_done(inst, pkt, sizeof(pkt)); - /* file should NOT be assembled (sig mismatch) */ - if (!file_exists(dl.dest_path)) OK(); else FAIL("assembled despite bad sig"); + /* файл не собран; блок в WAIT (1 держатель, попытка ушла) */ + int ok = dl->blocks_validated == 0 && !dl->assembled && dl->blocks[0].attempts == 1 && dl->blocks[0].state == MD_BLK_WAIT; + if (ok) OK(); else FAIL("validated=%d assembled=%d attempts=%d state=%d", dl->blocks_validated, dl->assembled, dl->blocks[0].attempts, dl->blocks[0].state); -done4: - 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); + free_test_dl(inst, dl); + sqlite3_close(db); u_free(inst); } } -static void test_download_cancel(void) { - TEST("download cancel marks inactive and calls done_cb"); { +/* ── 4. фейловер блока по простою (фейковое время) ── */ + +static void test_failover_stall(void) { + TEST("failover on stall → next holder"); { sqlite3* db = NULL; sqlite3_open(":memory:", &db); struct UTUN_INSTANCE* inst = make_minimal_instance(db); - inst->md.downloads = queue_new(inst->ua, 256, offsetof(struct media_download, media_id), 16, "md_dl_test5"); - - struct media_download dl; memset(&dl, 0, sizeof(dl)); - make_uuid(dl.media_id); - dl.num_blocks = 1; - snprintf(dl.dest_path, sizeof(dl.dest_path), "%s/cancel.bin", g_temp_dir); - dl.block_ids = u_malloc(16); make_uuid(dl.block_ids); - dl.block_sigs = u_malloc(64); memset(dl.block_sigs, 0, 64); - dl.active = 1; - memcpy(dl.ll.data, dl.media_id, 16); - - struct ll_entry* qe = queue_entry_new(sizeof(struct media_download)); - if (!qe) { FAIL("queue_entry_new"); goto done5; } - memcpy(qe->data, &dl, sizeof(dl)); queue_data_put_with_index(inst->md.downloads, qe); - - int rc = media_download_cancel(inst, dl.media_id, NULL, 0); - if (rc == 0) OK(); else FAIL("cancel returned %d", rc); - -done5: - 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); + struct media_download* dl = make_test_dl(db, inst, 1); + if (!dl) { FAIL("make_test_dl"); sqlite3_close(db); u_free(inst); return; } + add_holder(dl, 0, 0x1111); + add_holder(dl, 0, 0x2222); + dl->blocks[0].holder_idx = 0; + dl->blocks[0].state = MD_BLK_RECV; + dl->blocks[0].bytes_received = 12345; + dl->inflight = 1; + uint64_t now = 1000000000ULL; + dl->blocks[0].last_progress_tb = now - 30000; /* 3s назад */ + + int rc = md_dl_failover_block(dl, 0, now); + int ok = rc == 0 && dl->blocks[0].attempts == 1 && dl->blocks[0].holder_idx == 1 + && dl->blocks[0].bytes_received == 0 && dl->blocks[0].state == MD_BLK_REQ; + if (ok) OK(); else FAIL("rc=%d attempts=%d hidx=%d bytes=%u state=%d", rc, dl->blocks[0].attempts, dl->blocks[0].holder_idx, dl->blocks[0].bytes_received, dl->blocks[0].state); + + free_test_dl(inst, dl); + sqlite3_close(db); u_free(inst); } } -/* ──────────────────────────────────────────────────────────────── - Tests for block download fixes (backpressure, OOB block_id, conn_mgr close) - ──────────────────────────────────────────────────────────────── */ +/* ── 5. исчерпание попыток → FAILED ── */ -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"); +static void test_attempts_exhausted(void) { + TEST("attempts exhausted → FAILED"); { + 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_holder(dl, 0, 0x1111); + dl->blocks[0].holder_idx = 0; + dl->blocks[0].state = MD_BLK_RECV; + dl->blocks[0].attempts = 2; /* уже 2 провала */ + dl->inflight = 1; + + int rc = md_dl_failover_block(dl, 0, 1000000000ULL); + int ok = rc == -1 && dl->blocks[0].state == MD_BLK_FAILED && dl->blocks[0].attempts == 3; + if (ok) OK(); else FAIL("rc=%d state=%d attempts=%d", rc, dl->blocks[0].state, dl->blocks[0].attempts); + + free_test_dl(inst, dl); + sqlite3_close(db); u_free(inst); + } +} - 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); +/* ── 6. md_dl_check_stall: только зависшие блоки ── */ - 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 void test_check_stall(void) { + TEST("check_stall fails over only stalled block"); { + 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; } + for (int bi = 0; bi < 2; bi++) { add_holder(dl, bi, 0x1111); add_holder(dl, bi, 0x2222); dl->blocks[bi].holder_idx = 0; dl->blocks[bi].state = MD_BLK_RECV; } + dl->inflight = 2; + uint64_t now = 1000000000ULL; + dl->blocks[0].last_progress_tb = now - 30000; /* завис */ + dl->blocks[1].last_progress_tb = now - 100; /* живой */ + + int rc = md_dl_check_stall(dl, now); + int ok = rc == 0 && dl->blocks[0].attempts == 1 && dl->blocks[1].attempts == 0; + if (ok) OK(); else FAIL("rc=%d b0.attempts=%d b1.attempts=%d", rc, dl->blocks[0].attempts, dl->blocks[1].attempts); + + free_test_dl(inst, dl); + sqlite3_close(db); u_free(inst); + } } -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; -} +/* ── 7. OVERLOADED → следующий держатель ── */ -/* ── Тест: conn_cb отправляет только первый блок (backpressure) ── */ -static void test_conn_cb_backpressure(void) { - TEST("conn_cb starts only first block (backpressure)"); { +static void test_overloaded_next_holder(void) { + TEST("OVERLOADED → next holder (no attempt)"); { 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_holder(dl, 0, 0x1111); + add_holder(dl, 0, 0x2222); + dl->blocks[0].holder_idx = 0; + dl->blocks[0].state = MD_BLK_RECV; + dl->inflight = 1; + + struct media_pkt_block_overloaded ov; memset(&ov, 0, sizeof(ov)); + ov.subcmd = MEDIA_SUBCMD_BLOCK_OVERLOADED; + memcpy(ov.media_id, dl->media_id, 16); memcpy(ov.block_id, dl->blocks[0].block_id, 16); + ov.retry_after_ms = 2000; + media_download_handle_overloaded(inst, (const uint8_t*)&ov, sizeof(ov)); + + int ok = dl->blocks[0].attempts == 0 && dl->blocks[0].holder_idx == 1 && dl->blocks[0].state == MD_BLK_REQ; + if (ok) OK(); else FAIL("attempts=%d hidx=%d state=%d", dl->blocks[0].attempts, dl->blocks[0].holder_idx, dl->blocks[0].state); + + free_test_dl(inst, dl); + sqlite3_close(db); u_free(inst); + } +} - add_peer(dl, 0x1234, 5); /* 5 blocks on peer */ +/* ── 8. OVERLOADED без кандидатов → WAIT ── */ - md_dl_conn_cb(NULL, 0x1234, 0, CONN_EVENT_UP, dl); +static void test_overloaded_wait(void) { + TEST("OVERLOADED no more holders → WAIT"); { + 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_holder(dl, 0, 0x1111); /* один держатель */ + dl->blocks[0].holder_idx = 0; + dl->blocks[0].state = MD_BLK_RECV; + dl->inflight = 1; - int started_count = 0; - for (int i = 0; i < dl->peers[0].num_blocks; i++) - if (dl->peers[0].blocks[i].started) started_count++; + struct media_pkt_block_overloaded ov; memset(&ov, 0, sizeof(ov)); + ov.subcmd = MEDIA_SUBCMD_BLOCK_OVERLOADED; + memcpy(ov.media_id, dl->media_id, 16); memcpy(ov.block_id, dl->blocks[0].block_id, 16); + media_download_handle_overloaded(inst, (const uint8_t*)&ov, sizeof(ov)); - 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); + int ok = dl->blocks[0].state == MD_BLK_WAIT && dl->blocks[0].attempts == 0; + if (ok) OK(); else FAIL("state=%d attempts=%d", dl->blocks[0].state, dl->blocks[0].attempts); - 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); + free_test_dl(inst, dl); + sqlite3_close(db); u_free(inst); } } -/* ── Тест: conn_cb для разных пиров ── */ -static void test_conn_cb_multi_peer(void) { - TEST("conn_cb starts first block per peer"); { +/* ── 9. RELAY_FULL → добавление держателей ── */ + +static void test_relay_full(void) { + TEST("RELAY_FULL → add holders + retry"); { 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_holder(dl, 0, 0x1111); /* 1 держатель */ + dl->blocks[0].holder_idx = 0; + dl->blocks[0].state = MD_BLK_RECV; + dl->inflight = 1; - add_peer(dl, 0x1000, 3); - add_peer(dl, 0x2000, 2); + uint64_t nodes[2] = { 0x3333, 0x4444 }; + media_download_handle_relay_full(inst, dl->media_id, dl->blocks[0].block_id, 0, nodes, 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->blocks[0].num_holders == 3 && dl->blocks[0].state == MD_BLK_REQ; + if (ok) OK(); else FAIL("holders=%d state=%d", dl->blocks[0].num_holders, dl->blocks[0].state); - 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); + free_test_dl(inst, dl); + sqlite3_close(db); u_free(inst); } } -/* ── Тест: BLOCK_DONE → relay: следующий блок с того же пира ── */ -static void test_done_block_relay(void) { - TEST("BLOCK_DONE triggers next block from same peer"); { +/* ── 10. пустой QUERY_RESP ── */ + +static void test_query_resp_empty(void) { + TEST("QUERY_RESP empty → no crash"); { 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); + struct media_download* dl = make_test_dl(db, inst, 1); 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; + uint8_t pkt[MEDIA_QUERY_RESP_HDR_SIZE]; memset(pkt, 0, sizeof(pkt)); + pkt[0] = MEDIA_SUBCMD_QUERY_RESP; /* num_entries = 0 */ + media_download_handle_query_resp(inst, pkt, sizeof(pkt)); + OK(); - /* 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); + free_test_dl(inst, dl); + sqlite3_close(db); u_free(inst); } } -/* ── Тест: ассемблинг при всех блоках ── */ -static void test_done_assembly(void) { - TEST("BLOCK_DONE all blocks → assembly"); { +/* ── 11. назначение держателей + разброс ── */ + +static void test_query_resp_assign(void) { + TEST("QUERY_RESP assigns holders (spread)"); { 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); + struct media_download* dl = make_test_dl(db, inst, 3); 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)); + /* 2 держателя на каждый блок */ + uint64_t n1 = 0x1111, n2 = 0x2222; + struct media_pkt_query_resp_entry entries[6]; + memset(entries, 0, sizeof(entries)); + for (int i = 0; i < 3; i++) { + entries[i * 2 + 0].node_id = n1; memcpy(entries[i * 2 + 0].block_id, dl->blocks[i].block_id, 16); + entries[i * 2 + 1].node_id = n2; memcpy(entries[i * 2 + 1].block_id, dl->blocks[i].block_id, 16); } - - 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); + uint8_t pkt[3 + sizeof(entries)]; + pkt[0] = MEDIA_SUBCMD_QUERY_RESP; + uint16_t ne = 6; memcpy(pkt + 1, &ne, 2); + memcpy(pkt + 3, entries, sizeof(entries)); + media_download_handle_query_resp(inst, pkt, sizeof(pkt)); + + /* блоки разложены: holder_idx = i % num_holders (стартовый разброс) */ + int ok = 1; + for (int i = 0; i < 3; i++) + if (dl->blocks[i].num_holders != 2) ok = 0; + if (dl->num_peers != 2) ok = 0; + if (ok) OK(); else FAIL("h0=%d h1=%d h2=%d peers=%d", dl->blocks[0].num_holders, dl->blocks[1].num_holders, dl->blocks[2].num_holders, dl->num_peers); + + free_test_dl(inst, dl); + sqlite3_close(db); u_free(inst); } } -/* ── Тест: пир имеет больше блоков чем dl → без OOB ── */ -static void test_peer_more_blocks_than_dl(void) { - TEST("conn_cb with peer blocks > dl blocks → no OOB"); { +/* ── 12. дедуп ── */ + +static void dedup_done_cb(void* arg, int err) { (void)arg; (void)err; } + +static void test_dedup(void) { + TEST("duplicate start ignored"); { 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 */ + 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, 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); + struct media_index_result r; memset(&r, 0, sizeof(r)); + memcpy(r.media_id, dl->media_id, 16); + r.num_blocks = 1; r.file_size = CHUNK_SIZE; r.block_size = CHUNK_SIZE; + int rc = media_download_start(inst, 0x1234, &r, "/tmp/x.bin", "/tmp", 0xAAA, dedup_done_cb, NULL, NULL, NULL); + /* rc == 0 (не ошибка), новая запись не создана (дедуп) */ + int count = 0; + if (inst->md.downloads) { struct ll_entry* e = inst->md.downloads->head; while (e) { count++; e = e->next; } } + if (rc == 0 && count == 1) OK(); else FAIL("rc=%d count=%d", rc, count); + + free_test_dl(inst, dl); + sqlite3_close(db); u_free(inst); } } @@ -497,16 +480,18 @@ int main(void) { test_setup(); printf("=== test_media_delivery_download ===\n"); - test_download_chunk(); - 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_chunk_offset_write(); test_done_assembly(); - test_peer_more_blocks_than_dl(); + test_done_bad_sig(); + test_failover_stall(); + test_attempts_exhausted(); + test_check_stall(); + test_overloaded_next_holder(); + test_overloaded_wait(); + test_relay_full(); + test_query_resp_empty(); + test_query_resp_assign(); + test_dedup(); printf("\n%d/%d passed, %d failed\n", g_passed, g_total, g_failed); test_cleanup(); diff --git a/tests/test_media_delivery_full.c b/tests/test_media_delivery_full.c index cf3f001c..e2c8d472 100644 --- a/tests/test_media_delivery_full.c +++ b/tests/test_media_delivery_full.c @@ -30,6 +30,7 @@ #include #include #include +#include static int G_PASSED = 0, G_FAILED = 0, G_TOTAL = 0; #define TEST(n) do { G_TOTAL++; printf(" %-60s", n); fflush(stdout); } while(0) @@ -46,6 +47,7 @@ static struct UTUN_INSTANCE* g_inst[N_NODES]; static uint64_t g_nid[N_NODES]; static uint64_t g_group_id = 0; static uint8_t g_test_mid[16], g_test_bid0[16], g_test_bid1[16]; +static uint8_t g_t2_mid[16], g_t2_bids[3][16]; static char g_tdir[256] = "/tmp/utun_mdf_XXXXXX"; static char g_cfg[N_NODES][256]; static char g_db_dir[N_NODES][320]; @@ -311,9 +313,110 @@ static void phase_b1_stream_test(void) { } } +/* ══════════════════════════════════════════════════════════ + Phase B2: усечённый последний блок — проверка offset на отправителе + (баг: block_start = chunk * chunk_size давал неверный offset для + последнего блока, где chunk_size усечён) + ══════════════════════════════════════════════════════════ */ + +static int create_truncated_file(void) { + char src_tmp[512]; snprintf(src_tmp, sizeof(src_tmp), "/tmp/mdl_trunc_%d.bin", getpid()); + char media_dir[512]; snprintf(media_dir, sizeof(media_dir), "%s/media", g_db_dir[0]); + char dst_path[512]; snprintf(dst_path, sizeof(dst_path), "%s/trunc_src.bin", media_dir); + char media_base[512]; snprintf(media_base, sizeof(media_base), "%s", g_db_dir[0]); + + /* 12288 байт, контент = (i >> 8), чтобы первый байт последнего блока отличал offset */ + const int FS = 12288, BS = 5120, NB = 3; + uint8_t file_data[FS]; + for (int i = 0; i < FS; i++) file_data[i] = (uint8_t)(i >> 8); + FILE* f = fopen(src_tmp, "wb"); + if (!f) return -1; + fwrite(file_data, 1, FS, f); fclose(f); + + { FILE* fin = fopen(src_tmp, "rb"), *fout = fopen(dst_path, "wb"); + if (fin && fout) { uint8_t buf[4096]; size_t rd; while ((rd = fread(buf, 1, sizeof(buf), fin)) > 0) fwrite(buf, 1, rd, fout); } + if (fin) fclose(fin); if (fout) fclose(fout); } + + uint8_t hash[32]; + { EVP_MD_CTX* ctx = EVP_MD_CTX_new(); EVP_DigestInit_ex(ctx, EVP_sha256(), NULL); + EVP_DigestUpdate(ctx, file_data, FS); EVP_DigestFinal_ex(ctx, hash, NULL); EVP_MD_CTX_free(ctx); } + + struct media_index_result result; memset(&result, 0, sizeof(result)); + media_index_generate_uuid(result.media_id); + memcpy(result.content_hash, hash, 32); + result.file_size = FS; result.block_size = BS; result.num_blocks = NB; + result.block_ids = u_malloc(NB * 16); result.block_sigs = u_malloc(NB * 64); + for (int i = 0; i < NB; i++) { + media_index_generate_uuid(result.block_ids + i * 16); + int sz = (i == NB - 1) ? FS - i * BS : BS; + uint8_t smsg[5200]; size_t soff = 0; + memcpy(smsg + soff, file_data + i * BS, sz); soff += (size_t)sz; + uint64_t nid = g_nid[0]; memcpy(smsg + soff, &nid, 8); soff += 8; + sc_ed25519_sign(g_inst[0]->my_ed25519_privkey, smsg, soff, result.block_sigs + i * 64); + } + int rc = media_index_commit(g_inst[0]->topo_sqlite_db, &result, g_nid[0], + g_inst[0]->my_ed25519_privkey, "test_ch", dst_path, media_base); + memcpy(g_t2_mid, result.media_id, 16); + for (int i = 0; i < NB; i++) memcpy(g_t2_bids[i], result.block_ids + i * 16, 16); + media_index_result_free(&result); + unlink(src_tmp); + return rc; +} + +static void phase_b2_truncated_last_block(void) { + TEST("truncated last block served from correct offset"); { + if (create_truncated_file() != 0) { FAIL("create_truncated_file"); return; } + + /* регистрируем вручную загрузку одного блока (последнего, усечённого) на n1 */ + struct UTUN_INSTANCE* inst = g_inst[1]; + if (!inst->md.downloads) inst->md.downloads = queue_new(inst->ua, 256, offsetof(struct media_download, media_id), 16, "md_dl_b2"); + struct media_download* dl = u_calloc(1, sizeof(*dl)); + if (!dl) { FAIL("u_calloc dl"); return; } + memcpy(dl->media_id, g_t2_mid, 16); + dl->num_blocks = 1; + snprintf(dl->dest_path, sizeof(dl->dest_path), "/tmp/utun_mdf_trunc.bin"); + dl->blocks = u_calloc(1, sizeof(struct media_download_block)); + dl->block_sigs = u_calloc(1, 64); + memcpy(dl->blocks[0].block_id, g_t2_bids[2], 16); /* блок 2 — усечённый */ + dl->blocks[0].state = MD_BLK_RECV; + dl->blocks[0].expected_size = 2048; + dl->block_size = 5120; dl->file_size = 12288; + dl->active = 1; dl->inst = inst; + dl->author_node_id = 0; /* чтобы BLOCK_DONE прошёл по memcmp-пути (sigs нулевые) */ + memcpy(dl->ll.data, dl->media_id, 16); + struct ll_entry* qe = queue_entry_new(sizeof(*dl)); + if (!qe) { u_free(dl->blocks); u_free(dl->block_sigs); u_free(dl); FAIL("queue_entry_new"); return; } + memcpy(qe->data, dl, sizeof(*dl)); u_free(dl); + queue_data_put_with_index(inst->md.downloads, qe); + + /* запрашиваем блок 2 (маршрут по UTUN-группе — тест не поднимает chat-групп routing) */ + struct media_pkt_block_req req; memset(&req, 0, sizeof(req)); + req.subcmd = MEDIA_SUBCMD_BLOCK_REQ; req.group_id = 0; /* → TOPO_GROUP_UTUN */ + memcpy(req.media_id, g_t2_mid, 16); memcpy(req.block_id, g_t2_bids[2], 16); + req.chunk = 2; + msend(g_inst[1], g_nid[0], (const uint8_t*)&req, sizeof(req)); + + /* ждём сборки (блок валиден → num_blocks=1 → сразу ассемблируется) */ + const char* dest = "/tmp/utun_mdf_trunc.bin"; + int a = 0; + while (a < 3000 && access(dest, F_OK) != 0) { uasync_poll(g_ua, POLL_MS); a++; } + + FILE* df = fopen(dest, "rb"); + if (!df) { FAIL("assembled file not created"); return; } + fseeko(df, 0, SEEK_END); off_t sz = ftello(df); fseeko(df, 0, SEEK_SET); + uint8_t first = 0; fread(&first, 1, 1, df); fclose(df); + + /* верный offset 10240 → первый байт = 40; баг (offset 4096) → 16 */ + if (sz == 2048 && first == 40) OK(); + else FAIL("size=%lld first=%d (expected 2048/40; bug would give 2048/16)", (long long)sz, first); + unlink(dest); + } +} + /* ── main ── */ int main(void) { debug_config_init(); debug_set_level(DEBUG_LEVEL_ERROR); + if (getenv("UTUN_TEST_DEBUG")) { debug_set_category_level_by_name("media", "trace"); debug_set_category_level_by_name("general", "info"); } utun_instance_set_tun_init_enabled(0); srand((unsigned)time(NULL)); printf("=== test_media_delivery_full ===\n"); @@ -381,6 +484,7 @@ int main(void) { phase_a3_admission(); phase_b1_stream_test(); + phase_b2_truncated_last_block(); fflush(stdout); fflush(stderr); printf("\n%d/%d passed, %d failed\n", G_PASSED, G_TOTAL, G_FAILED); fflush(stdout); diff --git a/tests/test_media_download_timeout.c b/tests/test_media_download_timeout.c new file mode 100644 index 00000000..4b72107f --- /dev/null +++ b/tests/test_media_download_timeout.c @@ -0,0 +1,123 @@ +// test_media_download_timeout.c — интеграционный тест watchdog-таймера скачивания +// +// Проверяет, что реальный uasync-таймер простоя срабатывает: +// - загрузка с недостижимым автором → watchdog (100ms) → re-QUERY → суперноды +// исчерпаны → done_cb(-1). +// - загрузка с держателем, который «завис» (нет чанков) → фейловер → 3 попытки → +// done_cb(-1). +// +// Таймауты укорочены (dl_stall_timeout_tb=100ms), фейковое время не используется — +// гоняем реальный цикл uasync_poll. + +#include "media_delivery.h" +#include "media_delivery_proto.h" +#include "media_download.h" +#include "media_index.h" +#include "../lib/u_async.h" +#include "../lib/debug_config.h" +#include "../lib/mem.h" +#include "../src/utun_instance.h" +#include +#include +#include +#include +#include + +static int G_PASSED = 0, G_FAILED = 0, G_TOTAL = 0; +#define TEST(n) do { G_TOTAL++; printf(" %-58s", n); fflush(stdout); } while(0) +#define OK() do { G_PASSED++; printf("OK\n"); } while(0) +#define FAIL(f,...) do { G_FAILED++; printf("FAIL: " f "\n", ##__VA_ARGS__); } while(0) + +static struct UASYNC* g_ua = NULL; +static int g_done = 0, g_err = 0; + +static void make_uuid(uint8_t b[16]) { for (int i = 0; i < 16; i++) b[i] = (uint8_t)(rand() & 0xFF); } + +static void done_cb(void* a, int e) { (void)a; g_done = 1; g_err = e; printf(" [done_cb err=%d]\n", e); } + +static struct UTUN_INSTANCE* make_minimal(void) { + struct UTUN_INSTANCE* inst = u_calloc(1, sizeof(*inst)); + if (!inst) return NULL; + inst->ua = g_ua; + inst->node_id = 0xDEADBEEFDEADBEEFULL; + sqlite3* db = NULL; sqlite3_open(":memory:", &db); + inst->topo_sqlite_db = db; + for (int i = 0; i < 32; i++) inst->my_ed25519_privkey[i] = (uint8_t)(rand() & 0xFF); + memset(inst->my_ed25519_pubkey, 0xDD, 32); + inst->md.inst = inst; inst->md.db = db; inst->md.self_node_id = inst->node_id; inst->md.initialized = 1; + inst->md.dl_stall_timeout_tb = 1000; /* 100ms */ + inst->md.dl_max_attempts = 3; + return inst; +} + +static void make_result(struct media_index_result* r, int num_blocks) { + memset(r, 0, sizeof(*r)); + make_uuid(r->media_id); + r->num_blocks = num_blocks; r->file_size = 100 * num_blocks; r->block_size = 100; + r->block_ids = u_calloc((size_t)num_blocks, 16); + r->block_sigs = u_calloc((size_t)num_blocks, 64); + for (int i = 0; i < num_blocks; i++) make_uuid(r->block_ids + i * 16); +} + +static int poll_until_done(int max_ms) { + int el = 0; + while (!g_done && el < max_ms) { uasync_poll(g_ua, 5); el += 5; } + return g_done; +} + +int main(void) { + debug_config_init(); debug_set_level(DEBUG_LEVEL_ERROR); + if (getenv("UTUN_TEST_DEBUG")) { debug_set_category_level_by_name("media", "trace"); debug_set_category_level_by_name("general", "info"); } + srand(1); + printf("=== test_media_download_timeout ===\n"); + g_ua = uasync_create(); + + /* ── 1. недостижимый автор → watchdog → done_cb(-1) ── */ + TEST("unreachable author → watchdog → fail"); { + struct UTUN_INSTANCE* inst = make_minimal(); + struct media_index_result r; make_result(&r, 1); + g_done = 0; g_err = 0; + int rc = media_download_start(inst, 0x1234, &r, "/tmp/utun_mdt_1.bin", "/tmp", 0xBADF00D00000001ULL, done_cb, NULL, NULL, NULL); + if (rc != 0) { FAIL("start rc=%d", rc); } + else if (poll_until_done(3000) && g_err != 0) OK(); + else FAIL("done=%d err=%d", g_done, g_err); + + media_index_result_free(&r); + /* нет активной загрузки после finish */ + if (inst->md.downloads) { struct ll_entry* e = inst->md.downloads->head; int n = 0; while (e) { n++; e = e->next; } if (n != 0) FAIL("downloads leftover=%d", n); } + sqlite3_close(inst->topo_sqlite_db); u_free(inst); + } + + /* ── 2. зависший держатель → фейловер → re-QUERY → fail ── */ + TEST("stalled holder → failover → fail"); { + struct UTUN_INSTANCE* inst = make_minimal(); + struct media_index_result r; make_result(&r, 1); + g_done = 0; g_err = 0; + int rc = media_download_start(inst, 0x1234, &r, "/tmp/utun_mdt_2.bin", "/tmp", 0xBADF00D00000002ULL, done_cb, NULL, NULL, NULL); + if (rc != 0) { + FAIL("start rc=%d", rc); + } else { + /* имитируем QUERY_RESP: автор сообщает, что держит блок */ + struct media_pkt_query_resp_entry e; + memset(&e, 0, sizeof(e)); + e.node_id = 0xBADF00D00000002ULL; + memcpy(e.block_id, r.block_ids, 16); + uint8_t pkt[3 + sizeof(e)]; + pkt[0] = MEDIA_SUBCMD_QUERY_RESP; + uint16_t ne = 1; memcpy(pkt + 1, &ne, 2); + memcpy(pkt + 3, &e, sizeof(e)); + media_download_handle_query_resp(inst, pkt, sizeof(pkt)); + + /* держатель «завис» — чанков нет, watchdog сделает фейловер → попытки → fail */ + if (poll_until_done(5000) && g_err != 0) OK(); + else FAIL("done=%d err=%d (expected error after 3 attempts)", g_done, g_err); + } + + media_index_result_free(&r); + sqlite3_close(inst->topo_sqlite_db); u_free(inst); + } + + printf("\n%d/%d passed, %d failed\n", G_PASSED, G_TOTAL, G_FAILED); + if (g_ua) { uasync_destroy(g_ua, 0); g_ua = NULL; } + return G_FAILED > 0 ? 1 : 0; +} diff --git a/tools/chatgui-android/app/src/main/java/com/utun/chat/data/ChatRepository.kt b/tools/chatgui-android/app/src/main/java/com/utun/chat/data/ChatRepository.kt index 565afde7..f96871d7 100644 --- a/tools/chatgui-android/app/src/main/java/com/utun/chat/data/ChatRepository.kt +++ b/tools/chatgui-android/app/src/main/java/com/utun/chat/data/ChatRepository.kt @@ -274,6 +274,7 @@ class ChatRepository { val localAttrs = obj.optString("localAttrs", "") if (filePath.isNotEmpty()) downloadState = "downloaded" else if (localAttrs.contains("\"st\":\"fl\"")) downloadState = "downloaded" + else if (localAttrs.contains("\"st\":\"er\"")) downloadState = "error" var waveform = emptyList() var voiceDurationMs = 0 diff --git a/tools/chatgui-android/app/src/main/java/com/utun/chat/ui/components/MessageBubble.kt b/tools/chatgui-android/app/src/main/java/com/utun/chat/ui/components/MessageBubble.kt index 6a9f0310..7c3cac8f 100644 --- a/tools/chatgui-android/app/src/main/java/com/utun/chat/ui/components/MessageBubble.kt +++ b/tools/chatgui-android/app/src/main/java/com/utun/chat/ui/components/MessageBubble.kt @@ -208,11 +208,13 @@ private fun FileBubble(message: Message, textColor: Color, onDownload: (Message) val stateIcon = when { state == "downloaded" -> "\u2705" state == "downloading" -> "\u23F3" + state == "error" -> "\u26A0\uFE0F" else -> "\u2B07\uFE0F" } val stateText = when { state == "downloaded" -> "Tap to save" state == "downloading" -> "Downloading..." + state == "error" -> "Download failed, tap to retry" else -> "Tap to download" } @@ -220,7 +222,7 @@ private fun FileBubble(message: Message, textColor: Color, onDownload: (Message) modifier = Modifier.clickable { when { state == "downloaded" -> onOpen(message) - state.isEmpty() -> onDownload(message) + else -> onDownload(message) } } ) { @@ -272,6 +274,7 @@ private fun VideoBubble(message: Message, textColor: Color, onPlay: (Message) -> Unit = {}) { val downloading = message.downloadState == "downloading" val downloaded = message.filePath.isNotEmpty() + val isError = message.downloadState == "error" var thumb by remember(message.filePath) { mutableStateOf(null) } var durMs by remember(message.filePath) { mutableStateOf(0) } @@ -302,7 +305,7 @@ private fun VideoBubble(message: Message, textColor: Color, Box(modifier = Modifier.fillMaxSize(), contentAlignment = Alignment.Center) { Box(modifier = Modifier.size(44.dp).background(Color(0x88000000), CircleShape), contentAlignment = Alignment.Center) { Text( - when { downloading -> "\u23F3"; !downloaded -> "\u2B07"; else -> "\u25B6" }, + when { downloading -> "\u23F3"; isError -> "\u267B"; !downloaded -> "\u2B07"; else -> "\u25B6" }, fontSize = 18.sp, color = Color.White ) } @@ -323,6 +326,7 @@ private fun VideoBubble(message: Message, textColor: Color, Text( "${formatFileSize(message.fileSize)} · " + when { downloading -> "Downloading..." + isError -> "Failed, tap to retry" !downloaded -> "Tap to download" else -> "Tap to play" },