Browse Source

media: рефакторинг загрузки блоков — конечный автомат, фейловер, watchdog простоя

- media_download: состояние блока (IDLE/CONN/REQ/RECV/WAIT/DONE/FAILED),
  кандидаты-держатели (MD_MAX_HOLDERS_PER_BLOCK), фейловер + ретраи (MEDIA_DL_MAX_ATTEMPTS),
  watchdog простоя (MEDIA_DL_STALL_TIMEOUT_TB) вместо фейкового времени
- media_delivery: настраиваемые dl_stall_timeout_tb / dl_max_attempts
- chat_msg: пометка ошибки скачивания в local_attrs ("st":"er") для UI
- tests: новый test_media_download_timeout (watchdog/фейловер), рефакторинг
  test_media_delivery_download/full
- chatgui-android: состояние 'error' у вложений/видео (иконка + tap to retry)
v2
evgeny 3 weeks ago
parent
commit
f28d3e91f7
  1. 4
      src/chat/chat_msg.c
  2. 56
      src/media_delivery/media_delivery.c
  3. 4
      src/media_delivery/media_delivery.h
  4. 748
      src/media_delivery/media_download.c
  5. 88
      src/media_delivery/media_download.h
  6. 5
      tests/Makefile.am
  7. 639
      tests/test_media_delivery_download.c
  8. 104
      tests/test_media_delivery_full.c
  9. 123
      tests/test_media_download_timeout.c
  10. 1
      tools/chatgui-android/app/src/main/java/com/utun/chat/data/ChatRepository.kt
  11. 8
      tools/chatgui-android/app/src/main/java/com/utun/chat/ui/components/MessageBubble.kt

4
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);

56
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; }

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

748
src/media_delivery/media_download.c

File diff suppressed because it is too large Load Diff

88
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,10 +95,6 @@ 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,
@ -104,7 +103,6 @@ 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);
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);
@ -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

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

639
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 cdata[256]; for (int i = 0; i < 256; i++) cdata[i] = (uint8_t)(i & 0xFF);
uint8_t pkt[MEDIA_BLOCK_CHUNK_HDR_SIZE];
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));
OK(); /* should not crash or create files */
char tmp[1024]; snprintf(tmp, sizeof(tmp), "%s.chunk_0", dl->dest_path);
int s1 = file_size(tmp);
if (inst->md.downloads) queue_free(inst->md.downloads);
u_free(inst); sqlite3_close(db);
/* дубль с тем же 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));
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);
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);
int sz = file_size(dest);
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);
/* dl уже освобождён через md_dl_finish */
sqlite3_close(db); u_free(inst);
}
}
/* ── 3. BLOCK_DONE: неверная сигнатура → фейловер ── */
TEST("download handle_done bad sig → retry"); {
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;
/* 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); }
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();
/* 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();

104
tests/test_media_delivery_full.c

@ -30,6 +30,7 @@
#include <stdlib.h>
#include <string.h>
#include <time.h>
#include <unistd.h>
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);

123
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 <sqlite3.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
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;
}

1
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<Float>()
var voiceDurationMs = 0

8
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<Bitmap?>(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"
},

Loading…
Cancel
Save