Browse Source

media_delivery: add tests + fix queue_new NULL ua and index_offset bugs

- test_media_delivery_sql: 15 tests (SQL tables, CRUD, packed structs)
- test_media_delivery_download: 5 tests (chunk write/append, sig verify, assembly, cancel)
- Fix: queue_new(NULL) returns NULL — pass inst->ua everywhere
- Fix: queue index offset must be offsetof(struct, key_field) not 0
- Move struct media_download to header for test visibility
topo_upd
evgeny 2 months ago
parent
commit
a863da7737
  1. 84
      src/media_delivery/media_delivery.c
  2. 145
      src/media_delivery/media_download.c
  3. 43
      src/media_delivery/media_download.h
  4. 10
      tests/Makefile.am
  5. 314
      tests/test_media_delivery_download.c
  6. 362
      tests/test_media_delivery_sql.c

84
src/media_delivery/media_delivery.c

@ -62,7 +62,10 @@ static int md_ba_insert(sqlite3* db, const uint8_t* block_uuid, uint64_t group_i
"INSERT OR REPLACE INTO block_availability(block_uuid,group_id,node_id,chunk,timestamp)"
" VALUES(?,?,?,?,?)";
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) return -1;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: ba_insert prepare failed: %s", MD_ID, sqlite3_errmsg(db));
return -1;
}
sqlite3_bind_blob(st, 1, block_uuid, 16, SQLITE_STATIC);
sqlite3_bind_int64(st, 2, (sqlite3_int64)group_id);
sqlite3_bind_int64(st, 3, (sqlite3_int64)node_id);
@ -70,7 +73,11 @@ static int md_ba_insert(sqlite3* db, const uint8_t* block_uuid, uint64_t group_i
sqlite3_bind_int64(st, 5, (sqlite3_int64)ts);
int rc = sqlite3_step(st);
sqlite3_finalize(st);
return (rc == SQLITE_DONE) ? 0 : -1;
if (rc != SQLITE_DONE) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: ba_insert step failed rc=%d: %s", MD_ID, rc, sqlite3_errmsg(db));
return -1;
}
return 0;
}
static int md_ba_find_nodes(sqlite3* db, const uint8_t* block_uuid,
@ -78,7 +85,10 @@ static int md_ba_find_nodes(sqlite3* db, const uint8_t* block_uuid,
const char* sql =
"SELECT node_id, timestamp FROM block_availability WHERE block_uuid=? ORDER BY timestamp DESC LIMIT ?";
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) return 0;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: ba_find_nodes prepare failed: %s", MD_ID, sqlite3_errmsg(db));
return 0;
}
sqlite3_bind_blob(st, 1, block_uuid, 16, SQLITE_STATIC);
sqlite3_bind_int(st, 2, max);
int n = 0;
@ -96,7 +106,10 @@ static int md_ba_get_since(sqlite3* db, uint64_t since_id,
"SELECT id,block_uuid,group_id,node_id,chunk,timestamp"
" FROM block_availability WHERE id > ? ORDER BY id ASC LIMIT 4";
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) return 0;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: ba_get_since prepare failed: %s", MD_ID, sqlite3_errmsg(db));
return 0;
}
sqlite3_bind_int64(st, 1, (sqlite3_int64)since_id);
int off = 0;
while (sqlite3_step(st) == SQLITE_ROW) {
@ -116,7 +129,10 @@ static int md_ba_get_since(sqlite3* db, uint64_t since_id,
static uint64_t md_ss_get(sqlite3* db, uint64_t peer_node_id) {
const char* sql = "SELECT last_recv_id FROM super_sync WHERE peer_node_id=?";
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) return 0;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: ss_get prepare failed: %s", MD_ID, sqlite3_errmsg(db));
return 0;
}
sqlite3_bind_int64(st, 1, (sqlite3_int64)peer_node_id);
uint64_t v = 0;
if (sqlite3_step(st) == SQLITE_ROW) v = (uint64_t)sqlite3_column_int64(st, 0);
@ -127,7 +143,10 @@ static uint64_t md_ss_get(sqlite3* db, uint64_t peer_node_id) {
static void md_ss_set(sqlite3* db, uint64_t peer_node_id, uint64_t last_recv_id) {
const char* sql = "INSERT OR REPLACE INTO super_sync(peer_node_id,last_recv_id) VALUES(?,?)";
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) return;
if (sqlite3_prepare_v2(db, sql, -1, &st, NULL) != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: ss_set prepare failed: %s", MD_ID, sqlite3_errmsg(db));
return;
}
sqlite3_bind_int64(st, 1, (sqlite3_int64)peer_node_id);
sqlite3_bind_int64(st, 2, (sqlite3_int64)last_recv_id);
sqlite3_step(st);
@ -146,11 +165,11 @@ static struct media_super_peer* md_super_peer_add(struct media_delivery_ctx* md,
struct media_super_peer* sp = md_super_peer_find(md, node_id);
if (sp) return sp;
if (!md->super_peers) {
md->super_peers = queue_new(NULL, 64, 0, 8, "md_super");
if (!md->super_peers) return NULL;
md->super_peers = queue_new(md->inst->ua, 64, offsetof(struct media_super_peer, peer_node_id), 8, "md_super");
if (!md->super_peers) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_new(super_peers) failed", MD_ID); return NULL; }
}
struct ll_entry* qe = queue_entry_new(sizeof(struct media_super_peer));
if (!qe) return NULL;
if (!qe) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_entry_new(super_peer) failed", MD_ID); return NULL; }
sp = (struct media_super_peer*)qe->data;
memset(sp, 0, sizeof(*sp));
sp->peer_node_id = node_id;
@ -274,13 +293,13 @@ static void md_handle_serve_reg(struct media_delivery_ctx* md, uint64_t from_nod
struct media_pkt_serve_reg* pkt = (struct media_pkt_serve_reg*)data;
if (!md->served_nodes) {
md->served_nodes = queue_new(NULL, 64, 0, 8, "md_served");
if (!md->served_nodes) return;
md->served_nodes = queue_new(md->inst->ua, 64, offsetof(struct media_served_node, node_id), 8, "md_served");
if (!md->served_nodes) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_new(served_nodes) failed", MD_ID); return; }
}
struct ll_entry* e = queue_find_data_by_index(md->served_nodes, &from_node);
if (!e) {
struct ll_entry* qe = queue_entry_new(sizeof(struct media_served_node));
if (!qe) return;
if (!qe) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_entry_new(served_node) failed", MD_ID); return; }
struct media_served_node* sn = (struct media_served_node*)qe->data;
memset(sn, 0, sizeof(*sn));
sn->node_id = from_node; sn->group_id = pkt->group_id;
@ -298,7 +317,8 @@ static void md_handle_serve_reg(struct media_delivery_ctx* md, uint64_t from_nod
static void md_handle_query(struct media_delivery_ctx* md, uint64_t from_node,
const uint8_t* data, size_t len) {
if (len < MEDIA_QUERY_HDR_SIZE || !md->db) return;
if (len < MEDIA_QUERY_HDR_SIZE) return;
if (!md->db) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: QUERY but no DB", MD_ID); return; }
struct media_pkt_query* q = (struct media_pkt_query*)data;
uint16_t nb = q->num_blocks;
if (nb > 64) nb = 64;
@ -332,7 +352,8 @@ static void md_handle_query(struct media_delivery_ctx* md, uint64_t from_node,
static void md_handle_have_block(struct media_delivery_ctx* md, uint64_t from_node,
const uint8_t* data, size_t len) {
if (len < MEDIA_HAVE_BLOCK_SIZE || !md->db) return;
if (len < MEDIA_HAVE_BLOCK_SIZE) return;
if (!md->db) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: HAVE_BLOCK but no DB", MD_ID); return; }
struct media_pkt_have_block* hb = (struct media_pkt_have_block*)data;
int rc = md_ba_insert(md->db, hb->block_id, hb->group_id, from_node, hb->chunk, hb->timestamp);
@ -364,7 +385,8 @@ static void md_handle_have_block(struct media_delivery_ctx* md, uint64_t from_no
static void md_handle_super_repl(struct media_delivery_ctx* md, uint64_t from_node,
const uint8_t* data, size_t len) {
if (len < MEDIA_SUPER_REPL_HDR_SIZE || !md->db) return;
if (len < MEDIA_SUPER_REPL_HDR_SIZE) return;
if (!md->db) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: SUPER_REPL but no DB", MD_ID); return; }
struct media_pkt_super_repl* rp = (struct media_pkt_super_repl*)data;
uint16_t num = rp->num_entries;
int64_t max_id = 0;
@ -399,7 +421,7 @@ static void md_handle_super_ack(struct media_delivery_ctx* md, uint64_t from_nod
struct media_pkt_super_ack* sa = (struct media_pkt_super_ack*)data;
struct media_super_peer* peer = md_super_peer_find(md, from_node);
if (!peer) return;
if (!peer) { DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: SUPER_ACK from unknown peer 0x%016llx", MD_ID, (unsigned long long)from_node); return; }
peer->peer_last_recv_id = sa->ack_seq;
peer->inflight_count--;
@ -413,11 +435,12 @@ static void md_handle_super_ack(struct media_delivery_ctx* md, uint64_t from_nod
static void md_handle_super_hello(struct media_delivery_ctx* md, uint64_t from_node,
const uint8_t* data, size_t len) {
if (len < sizeof(struct media_pkt_super_hello) || !md->db) return;
if (len < sizeof(struct media_pkt_super_hello)) return;
if (!md->db) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: SUPER_HELLO but no DB", MD_ID); return; }
struct media_pkt_super_hello* sh = (struct media_pkt_super_hello*)data;
struct media_super_peer* peer = md_super_peer_add(md, from_node);
if (!peer) return;
if (!peer) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: super_peer_add failed for 0x%016llx", MD_ID, (unsigned long long)from_node); return; }
peer->peer_last_recv_id = sh->last_recv_id;
peer->connected = 1;
@ -443,7 +466,7 @@ static void md_handle_super_hello(struct media_delivery_ctx* md, uint64_t from_n
static void md_super_conn_cb(int result, uint64_t node_id, void* arg) {
struct media_delivery_ctx* md = (struct media_delivery_ctx*)arg;
struct media_super_peer* peer = md_super_peer_find(md, node_id);
if (!peer) return;
if (!peer) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: conn_cb for unknown peer 0x%016llx", MD_ID, (unsigned long long)node_id); return; }
if (result != CONN_MGR_OK) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: connect to supernode 0x%016llx failed rc=%d",
@ -466,13 +489,16 @@ static void md_super_conn_cb(int result, uint64_t node_id, void* arg) {
if (md_send(md->inst, node_id, (const uint8_t*)&hello, sizeof(hello)) == 0) {
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: SUPER_HELLO sent to 0x%016llx last_recv=%lld",
MD_ID, (unsigned long long)node_id, (long long)my_last);
} else {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: SUPER_HELLO send failed to 0x%016llx",
MD_ID, (unsigned long long)node_id);
}
/* wait for HELLO response — md_handle_super_hello will start replication */
}
static void md_super_connect(struct media_delivery_ctx* md, uint64_t peer_node_id) {
struct media_super_peer* peer = md_super_peer_add(md, peer_node_id);
if (!peer) return;
if (!peer) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: super_peer_add failed for 0x%016llx", MD_ID, (unsigned long long)peer_node_id); return; }
struct TOPO_GROUP* grp = NULL;
/* find group containing this peer */
@ -495,7 +521,10 @@ static void md_super_connect(struct media_delivery_ctx* md, uint64_t peer_node_i
/* ── supernode start/stop ── */
static int md_super_start(struct media_delivery_ctx* md) {
if (!md->inst || !md->inst->topo_sqlite_db || !md->inst->topo_groups) return -1;
if (!md->inst || !md->inst->topo_sqlite_db || !md->inst->topo_groups) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: super_start missing inst/db/groups", MD_ID);
return -1;
}
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: supernode START node_id=0x%016llx",
MD_ID, (unsigned long long)md->self_node_id);
@ -524,7 +553,7 @@ static int md_super_start(struct media_delivery_ctx* md) {
}
static void md_super_stop(struct media_delivery_ctx* md) {
if (!md->db) return;
if (!md->db) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: super_stop but no DB", MD_ID); return; }
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: supernode STOP, cleaning foreign blocks", MD_ID);
@ -631,7 +660,7 @@ static void md_etcp_recv(struct ETCP_CONN* conn, struct ll_entry* entry) {
struct UTUN_INSTANCE* inst = conn->instance;
if (!inst) return;
struct media_delivery_ctx* md = &inst->md;
if (!md->initialized) return;
if (!md->initialized) { DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: recv before init", MD_ID); return; }
const uint8_t* data = entry->dgram;
size_t len = entry->len;
@ -646,6 +675,8 @@ static void md_etcp_recv(struct ETCP_CONN* conn, struct ll_entry* entry) {
md_handle_serve_reg(md, from_node, data, len);
break;
case MEDIA_SUBCMD_SERVE_LEAVE:
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: SERVE_LEAVE from 0x%016llx",
MD_ID, (unsigned long long)from_node);
if (md->served_nodes) {
struct ll_entry* e = queue_find_data_by_index(md->served_nodes, &from_node);
if (e) { queue_remove_data(md->served_nodes, e); queue_entry_free(e); }
@ -653,18 +684,23 @@ static void md_etcp_recv(struct ETCP_CONN* conn, struct ll_entry* entry) {
break;
case MEDIA_SUBCMD_QUERY:
if (md->is_supernode) md_handle_query(md, from_node, data, len);
else DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: QUERY ignored — not supernode", MD_ID);
break;
case MEDIA_SUBCMD_HAVE_BLOCK:
if (md->is_supernode) md_handle_have_block(md, from_node, data, len);
else DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: HAVE_BLOCK ignored — not supernode", MD_ID);
break;
case MEDIA_SUBCMD_SUPER_REPL:
if (md->is_supernode) md_handle_super_repl(md, from_node, data, len);
else DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: SUPER_REPL ignored — not supernode", MD_ID);
break;
case MEDIA_SUBCMD_SUPER_ACK:
if (md->is_supernode) md_handle_super_ack(md, from_node, data, len);
else DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: SUPER_ACK ignored — not supernode", MD_ID);
break;
case MEDIA_SUBCMD_SUPER_HELLO:
if (md->is_supernode) md_handle_super_hello(md, from_node, data, len);
else DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: SUPER_HELLO ignored — not supernode", MD_ID);
break;
case MEDIA_SUBCMD_QUERY_RESP:
media_download_handle_query_resp(inst, entry->dgram, entry->len);
@ -727,7 +763,7 @@ static int media_delivery_create_tables(struct UTUN_INSTANCE* inst) {
/* ── public API ── */
int media_delivery_init(struct UTUN_INSTANCE* inst) {
if (!inst) return -1;
if (!inst) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: init with NULL inst", MD_ID); return -1; }
struct media_delivery_ctx* md = &inst->md;
memset(md, 0, sizeof(*md));

145
src/media_delivery/media_download.c

@ -19,39 +19,6 @@
#define MDL_ID "media_download"
/* ── internal download state ── */
struct media_download {
struct ll_entry ll; // индекс по media_id (16 байт)
uint8_t media_id[16];
uint64_t group_id;
char dest_path[1024];
char media_base[512];
int num_blocks;
uint8_t* block_ids; // num_blocks * 16
uint8_t* block_sigs; // num_blocks * 64
uint8_t content_hash[32];
int64_t file_size;
int64_t block_size;
int blocks_received;
int blocks_validated;
int num_peers;
struct media_download_peer peers[10];
uint8_t active;
uint8_t assembled;
int err;
/* supernode list */
uint64_t super_nodes[10];
int super_count;
int super_current;
void* query_timer;
void* timeout_timer;
void (*done_cb)(void* arg, int err);
void* done_arg;
};
/* ── helpers ── */
static struct media_download* md_dl_find(struct UTUN_INSTANCE* inst, const uint8_t* media_id) {
@ -73,12 +40,12 @@ static void md_dl_free(struct media_download* dl) {
static int md_dl_send(struct UTUN_INSTANCE* inst, uint64_t dst, const uint8_t* data, size_t len) {
struct ll_entry* e = queue_entry_new(0);
if (!e) return -1;
if (!e) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_entry_new failed for send to 0x%016llx", MDL_ID, (unsigned long long)dst); return -1; }
e->dgram = u_malloc(len);
if (!e->dgram) { queue_entry_free(e); return -1; }
if (!e->dgram) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: u_malloc(%zu) failed for send to 0x%016llx", MDL_ID, len, (unsigned long long)dst); queue_entry_free(e); return -1; }
memcpy(e->dgram, data, len); e->len = (uint16_t)len;
int rc = etcp_route_send(inst, 0, dst, e, 1);
if (rc != 0) { u_free(e->dgram); queue_entry_free(e); }
if (rc != 0) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: etcp_route_send failed rc=%d to 0x%016llx", MDL_ID, rc, (unsigned long long)dst); u_free(e->dgram); queue_entry_free(e); }
return rc;
}
@ -103,7 +70,7 @@ static int md_dl_send_block_req(struct UTUN_INSTANCE* inst, uint64_t dst,
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;
if (!super) return;
if (!super) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: HAVE_BLOCK but no supernode available", MDL_ID); return; }
struct media_pkt_have_block hb;
memset(&hb, 0, sizeof(hb));
@ -114,13 +81,15 @@ static void md_dl_send_have_block(struct UTUN_INSTANCE* inst, struct media_downl
hb.chunk = (uint32_t)bi;
hb.timestamp = (int64_t)time(NULL);
/* sign: block_id + chunk + timestamp + my_node_id */
uint8_t smsg[64]; size_t soff = 0;
memcpy(smsg + soff, dl->block_ids + bi * 16, 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;
sc_ed25519_sign(inst->my_ed25519_privkey, smsg, soff, hb.node_sign);
if (sc_ed25519_sign(inst->my_ed25519_privkey, smsg, soff, hb.node_sign) != SC_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: Ed25519 sign failed for HAVE_BLOCK", MDL_ID);
return;
}
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: HAVE_BLOCK to super 0x%016llx block=%d",
MDL_ID, (unsigned long long)super, bi);
@ -191,7 +160,7 @@ static void md_dl_send_query(struct UTUN_INSTANCE* inst, struct media_download*
size_t pkt_len = MEDIA_QUERY_HDR_SIZE + (size_t)dl->num_blocks * 16;
uint8_t* pkt = u_malloc(pkt_len);
if (!pkt) return;
if (!pkt) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: u_malloc(%zu) failed for QUERY pkt", MDL_ID, pkt_len); return; }
struct media_pkt_query* q = (struct media_pkt_query*)pkt;
memset(q, 0, sizeof(*q));
q->subcmd = MEDIA_SUBCMD_QUERY;
@ -221,9 +190,8 @@ void media_download_handle_query_resp(struct UTUN_INSTANCE* inst,
const uint8_t* bid = entries[0].block_id;
struct media_download* dl = md_dl_find(inst, bid);
if (!dl) {
/* search all blocks */
for (int i = 0; i < ne; i++) { dl = md_dl_find(inst, entries[i].block_id); if (dl) break; }
if (!dl) return;
if (!dl) { DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: QUERY_RESP — no matching download", MDL_ID); return; }
}
dl->num_peers = 0;
@ -260,7 +228,7 @@ void media_download_handle_query_resp(struct UTUN_INSTANCE* inst,
if (g->conn_mgr) { grp = g; break; }
gle = gle->next;
}
if (!grp || !grp->conn_mgr) return;
if (!grp || !grp->conn_mgr) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: QUERY_RESP — no conn_mgr to connect peers", MDL_ID); return; }
for (int pi = 0; pi < dl->num_peers; pi++) {
conn_mgr_connect_node(grp->conn_mgr, dl->peers[pi].node_id, 0, md_dl_conn_cb, dl);
@ -276,26 +244,27 @@ void media_download_handle_chunk(struct UTUN_INSTANCE* inst,
const uint8_t* d = data;
struct media_pkt_block_chunk* ch = (struct media_pkt_block_chunk*)data;
struct media_download* dl = md_dl_find(inst, ch->media_id);
if (!dl || !dl->active) return;
if (!dl) { DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%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_DEBUG, "%s: CHUNK for inactive download", MDL_ID); return; }
/* find block index */
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; }
}
if (bi < 0) return;
if (bi < 0) { DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: CHUNK for unknown block %02x%02x...", MDL_ID, ch->block_id[0], ch->block_id[1]); return; }
/* write chunk to temp file */
char tmp[2048];
snprintf(tmp, sizeof(tmp), "%s.chunk_%d", dl->dest_path, bi);
FILE* f = fopen(tmp, "ab");
if (f) {
size_t wlen = ch->data_len;
if (len >= MEDIA_BLOCK_CHUNK_HDR_SIZE + wlen) {
fwrite(d + MEDIA_BLOCK_CHUNK_HDR_SIZE, 1, wlen, f);
}
fclose(f);
if (!f) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: fopen(%s) failed for chunk write", MDL_ID, tmp); return; }
size_t wlen = ch->data_len;
if (len >= MEDIA_BLOCK_CHUNK_HDR_SIZE + wlen) {
fwrite(d + MEDIA_BLOCK_CHUNK_HDR_SIZE, 1, wlen, f);
} else {
DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: CHUNK truncated: len=%zu need=%zu+%zu",
MDL_ID, len, (size_t)MEDIA_BLOCK_CHUNK_HDR_SIZE, (size_t)wlen);
}
fclose(f);
}
/* ── handle incoming BLOCK_DONE ── */
@ -307,13 +276,14 @@ void media_download_handle_done(struct UTUN_INSTANCE* inst,
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 || !dl->active) return;
if (!dl) { DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: BLOCK_DONE for unknown media", MDL_ID); return; }
if (!dl->active) { DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%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; }
}
if (bi < 0) return;
if (bi < 0) { DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: BLOCK_DONE for unknown block", MDL_ID); return; }
/* verify signature */
EVP_PKEY* pkey = NULL;
@ -324,31 +294,27 @@ void media_download_handle_done(struct UTUN_INSTANCE* inst,
char tmp[2048];
snprintf(tmp, sizeof(tmp), "%s.chunk_%d", dl->dest_path, bi);
FILE* f = fopen(tmp, "rb");
if (f) {
fseeko(f, 0, SEEK_END);
off_t fsz = ftello(f);
fseeko(f, 0, SEEK_SET);
uint8_t* buf = u_malloc((size_t)fsz);
if (buf) {
size_t rd = fread(buf, 1, (size_t)fsz, f);
uint64_t node_id = inst->node_id;
uint8_t smsg[128]; size_t soff = 0;
memcpy(smsg + soff, buf, rd); soff += rd;
memcpy(smsg + soff, &node_id, 8); soff += 8;
EVP_DigestVerifyInit(vctx, NULL, EVP_sha256(), NULL, pkey);
/* fallback: compare raw block_sig */
sig_ok = (rd == (size_t)bd->total_size && rd > 0) ? 1 : 0;
if (sig_ok && rd == (size_t)bd->total_size) {
/* basic check: block_sig should match pre-computed */
if (memcmp(bd->block_sig, dl->block_sigs + bi * 64, 64) == 0) sig_ok = 1;
else sig_ok = 0;
}
u_free(buf);
}
fclose(f);
if (!f) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: fopen(%s) failed for sig check", MDL_ID, tmp); EVP_MD_CTX_free(vctx); return; }
fseeko(f, 0, SEEK_END);
off_t fsz = ftello(f);
fseeko(f, 0, SEEK_SET);
uint8_t* buf = u_malloc((size_t)fsz);
if (!buf) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: u_malloc(%lld) failed for sig check", MDL_ID, (long long)fsz); fclose(f); EVP_MD_CTX_free(vctx); return; }
size_t rd = fread(buf, 1, (size_t)fsz, f);
uint64_t node_id = inst->node_id;
uint8_t smsg[128]; size_t soff = 0;
memcpy(smsg + soff, buf, rd); soff += rd;
memcpy(smsg + soff, &node_id, 8); soff += 8;
EVP_DigestVerifyInit(vctx, NULL, EVP_sha256(), NULL, pkey);
sig_ok = (rd == (size_t)bd->total_size && rd > 0) ? 1 : 0;
if (sig_ok && rd == (size_t)bd->total_size) {
if (memcmp(bd->block_sig, dl->block_sigs + bi * 64, 64) == 0) sig_ok = 1;
else sig_ok = 0;
}
EVP_MD_CTX_free(vctx);
u_free(buf);
fclose(f);
}
EVP_MD_CTX_free(vctx);
if (sig_ok) {
dl->blocks_received++;
@ -363,7 +329,8 @@ void media_download_handle_done(struct UTUN_INSTANCE* inst,
if (dl->blocks_validated == dl->num_blocks) {
/* assemble file */
FILE* out = fopen(dl->dest_path, "wb");
if (out) {
if (!out) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: fopen(%s) failed for assembly", MDL_ID, dl->dest_path); }
else {
for (int n = 0; n < dl->num_blocks; n++) {
char tmp2[2048];
snprintf(tmp2, sizeof(tmp2), "%s.chunk_%d", dl->dest_path, n);
@ -408,12 +375,16 @@ 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,
void (*done_cb)(void* arg, int err), void* done_arg) {
if (!inst || !result || !dest_path || !media_base || !done_cb) return -1;
if (!inst || !result || !dest_path || !media_base || !done_cb) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: start with NULL args inst=%p result=%p dest=%p base=%p cb=%p",
MDL_ID, (void*)inst, (void*)result, (void*)dest_path, (void*)media_base, (void*)done_cb);
return -1;
}
struct media_delivery_ctx* md = &inst->md;
if (!md->initialized) return -1;
if (!md->initialized) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: start before media_delivery init", MDL_ID); return -1; }
struct media_download* dl = u_calloc(1, sizeof(*dl));
if (!dl) return -1;
if (!dl) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: u_calloc for download failed", MDL_ID); return -1; }
memcpy(dl->media_id, result->media_id, 16);
dl->group_id = group_id;
@ -429,18 +400,18 @@ int media_download_start(struct UTUN_INSTANCE* inst, uint64_t group_id,
dl->block_ids = u_malloc((size_t)dl->num_blocks * 16);
dl->block_sigs = u_malloc((size_t)dl->num_blocks * 64);
if (!dl->block_ids || !dl->block_sigs) { u_free(dl->block_ids); u_free(dl); return -1; }
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);
memcpy(dl->block_sigs, result->block_sigs, (size_t)dl->num_blocks * 64);
/* register in downloads queue */
if (!md->downloads) {
md->downloads = queue_new(NULL, 256, 0, 16, "md_dl");
if (!md->downloads) { md_dl_free(dl); u_free(dl); return -1; }
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; }
}
memcpy(dl->ll.data, dl->media_id, 16);
struct ll_entry* qe = queue_entry_new(sizeof(struct media_download));
if (!qe) { md_dl_free(dl); u_free(dl); return -1; }
if (!qe) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: queue_entry_new(download) failed", MDL_ID); md_dl_free(dl); u_free(dl); return -1; }
memcpy(qe->data, dl, sizeof(*dl));
u_free(dl);
queue_data_put_with_index(md->downloads, qe);
@ -459,7 +430,7 @@ 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) return -1;
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;

43
src/media_delivery/media_download.h

@ -9,6 +9,10 @@ extern "C" {
#include <stdint.h>
#include <stddef.h>
#include <stdio.h>
#include "../lib/ll_queue.h"
struct UTUN_INSTANCE;
struct media_index_result;
struct UTUN_INSTANCE;
struct media_index_result;
@ -18,16 +22,45 @@ struct media_index_result;
struct media_download_peer {
uint64_t node_id;
uint8_t connected; // 1 = conn_mgr подключён
int num_blocks; // сколько блоков назначено этому узлу
uint8_t connected;
int num_blocks;
struct {
uint8_t block_id[16];
uint8_t started; // 1 = BLOCK_REQ отправлен
uint8_t received; // 1 = BLOCK_DONE получен
uint8_t validated; // 1 = подпись проверена
uint8_t started;
uint8_t received;
uint8_t validated;
} blocks[MD_MAX_BLOCKS_PER_PEER];
};
/* внутреннее состояние загрузки (доступно для тестов) */
struct media_download {
struct ll_entry ll;
uint8_t media_id[16];
uint64_t group_id;
char dest_path[1024];
char media_base[512];
int num_blocks;
uint8_t* block_ids;
uint8_t* block_sigs;
uint8_t content_hash[32];
int64_t file_size;
int64_t block_size;
int blocks_received;
int blocks_validated;
int num_peers;
struct media_download_peer peers[10];
uint8_t active;
uint8_t assembled;
int err;
uint64_t super_nodes[10];
int super_count;
int super_current;
void* query_timer;
void* timeout_timer;
void (*done_cb)(void* arg, int err);
void* done_arg;
};
/* состояние одного стрима на стороне блок-холдера */
struct media_block_stream {
uint8_t media_id[16];

10
tests/Makefile.am

@ -62,6 +62,8 @@ check_PROGRAMS = \
test_opus_codec \
test_media_async \
test_media_index \
test_media_delivery_sql \
test_media_delivery_download \
bench_timeout_heap \
bench_uasync_timeouts
@ -325,6 +327,14 @@ test_media_index_SOURCES = test_media_index.c
test_media_index_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/src/media_delivery -I$(top_srcdir)/src/media_async -I$(top_srcdir)/lib
test_media_index_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS)
test_media_delivery_sql_SOURCES = test_media_delivery_sql.c
test_media_delivery_sql_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/src/media_delivery -I$(top_srcdir)/lib
test_media_delivery_sql_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS)
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)
# Copy test configs to build directory (tests run from build/tests/)
all-local: copy-test-configs

314
tests/test_media_delivery_download.c

@ -0,0 +1,314 @@
// test_media_delivery_download.c — тест скачивания блоков (chunk → done → assembly)
//
// Покрытие:
// 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
#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 "../transport_layer/secure_channel.h"
#include "../routing_layer/topo_node_sqlite.h"
#include <sqlite3.h>
#include <openssl/evp.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
#include <sys/stat.h>
#define TEMP_DIR "/tmp/test_mdl_dl_XXXXXX"
#define CHUNK_SIZE 100
#define NUM_BLOCKS 3
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 OK() do { g_passed++; printf("OK\n"); } while(0)
#define FAIL(fmt, ...) do { g_failed++; printf("FAIL: " fmt "\n", ##__VA_ARGS__); } while(0)
/* ── helpers ── */
static void make_uuid(uint8_t buf[16]) { for (int i = 0; i < 16; i++) buf[i] = (uint8_t)(rand() & 0xFF); }
static void write_file(const char* path, const uint8_t* data, size_t len) {
FILE* f = fopen(path, "wb");
if (f) { fwrite(data, 1, len, f); fclose(f); }
}
static int file_exists(const char* path) { return access(path, F_OK) == 0; }
static int file_size(const char* path) {
struct stat st; if (stat(path, &st) != 0) return -1; return (int)st.st_size;
}
static void test_setup(void) {
snprintf(g_temp_dir, sizeof(g_temp_dir), "%s", TEMP_DIR);
if (!mkdtemp(g_temp_dir)) { fprintf(stderr, "mkdtemp failed\n"); exit(1); }
debug_config_init(); debug_set_level(DEBUG_LEVEL_ERROR);
g_ua = uasync_create();
}
static void test_cleanup(void) {
char cmd[512]; snprintf(cmd, sizeof(cmd), "rm -rf %s", g_temp_dir); system(cmd);
if (g_ua) { uasync_destroy(g_ua, 0); g_ua = NULL; }
}
/* ── create minimal UTUN_INSTANCE with media_delivery ── */
#include "../src/utun_instance.h"
static struct UTUN_INSTANCE* make_minimal_instance(sqlite3* db) {
struct UTUN_INSTANCE* inst = u_calloc(1, sizeof(*inst));
if (!inst) return NULL;
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;
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));
media_download_handle_chunk(inst, pkt, sizeof(pkt));
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);
/* 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);
if (ok1 && ok2) OK(); else FAIL("chunk: ok1=%d ok2=%d size=%d", ok1, ok2, file_size(tmp));
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);
}
TEST("download handle_chunk unknown media_id"); {
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");
uint8_t pkt[MEDIA_BLOCK_CHUNK_HDR_SIZE];
struct media_pkt_block_chunk* ch = (struct media_pkt_block_chunk*)pkt;
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;
media_download_handle_chunk(inst, pkt, sizeof(pkt));
OK(); /* should not crash or create files */
if (inst->md.downloads) queue_free(inst->md.downloads);
u_free(inst); sqlite3_close(db);
}
}
static void test_download_done(void) {
TEST("download handle_done sig ok + 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);
}
/* 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];
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);
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);
}
TEST("download handle_done bad sig → retry"); {
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 */
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);
/* send BLOCK_DONE with WRONG signature */
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;
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");
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);
}
}
static void test_download_cancel(void) {
TEST("download cancel marks inactive and calls done_cb"); {
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);
}
}
int main(void) {
test_setup();
printf("=== test_media_delivery_download ===\n");
test_download_chunk();
test_download_done();
test_download_cancel();
printf("\n%d/%d passed, %d failed\n", g_passed, g_total, g_failed);
test_cleanup();
return g_failed > 0 ? 1 : 0;
}

362
tests/test_media_delivery_sql.c

@ -0,0 +1,362 @@
// test_media_delivery_sql.c — SQL-таблицы + протокольные структуры media_delivery
//
// Покрытие:
// 1. Создание таблиц block_availability + super_sync + индексов
// 2. INSERT OR REPLACE (upsert) — дубликат не создаёт вторую строку
// 3. SELECT с фильтром по block_uuid, сортировка timestamp DESC
// 4. SELECT с id > N, сортировка ASC, LIMIT
// 5. super_sync: INSERT OR REPLACE + SELECT (дефолт 0)
// 6. DELETE WHERE node_id=?
// 7. DELETE WHERE node_id!=? (чистка чужих при выключении суперузла)
// 8. Размеры packed-структур
// 9. Проверка полей packed-структур (нет padding-сюрпризов)
// 10. Построение/разбор MEDIA_QUERY
// 11. Построение/разбор MEDIA_HAVE_BLOCK
// 12. Построение/разбор MEDIA_SUPER_HELLO
#include "../src/media_delivery/media_delivery_proto.h"
#include "../lib/debug_config.h"
#include <sqlite3.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <time.h>
static int g_passed = 0, g_failed = 0, g_total = 0;
#define TEST(name) do { g_total++; printf(" %-50s", name); } 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)
static void make_uuid(uint8_t buf[16]) {
for (int i = 0; i < 16; i++) buf[i] = (uint8_t)(rand() & 0xFF);
}
static sqlite3* make_db(void) {
sqlite3* db = NULL;
sqlite3_open(":memory:", &db);
sqlite3_exec(db, "PRAGMA journal_mode=memory", NULL, NULL, NULL);
return db;
}
static int create_tables(sqlite3* db) {
const char* sql_ba =
"CREATE TABLE IF NOT EXISTS block_availability ("
" id INTEGER PRIMARY KEY AUTOINCREMENT,"
" block_uuid BLOB NOT NULL,"
" group_id INTEGER NOT NULL,"
" node_id INTEGER NOT NULL,"
" chunk INTEGER NOT NULL,"
" timestamp INTEGER NOT NULL)";
int rc = sqlite3_exec(db, sql_ba, NULL, NULL, NULL);
if (rc != SQLITE_OK) return -1;
sqlite3_exec(db, "CREATE UNIQUE INDEX IF NOT EXISTS idx_ba_uuid_node ON block_availability(block_uuid, node_id)", NULL, NULL, NULL);
sqlite3_exec(db, "CREATE INDEX IF NOT EXISTS idx_ba_id ON block_availability(id)", NULL, NULL, NULL);
sqlite3_exec(db, "CREATE INDEX IF NOT EXISTS idx_ba_node_id ON block_availability(node_id)", NULL, NULL, NULL);
const char* sql_ss =
"CREATE TABLE IF NOT EXISTS super_sync ("
" peer_node_id INTEGER PRIMARY KEY,"
" last_recv_id INTEGER NOT NULL DEFAULT 0)";
rc = sqlite3_exec(db, sql_ss, NULL, NULL, NULL);
return (rc == SQLITE_OK) ? 0 : -1;
}
static int table_exists(sqlite3* db, const char* name) {
sqlite3_stmt* s = NULL;
int ok = 0;
if (sqlite3_prepare_v2(db, "SELECT name FROM sqlite_master WHERE type='table' AND name=?", -1, &s, NULL) == SQLITE_OK) {
sqlite3_bind_text(s, 1, name, -1, SQLITE_STATIC);
ok = (sqlite3_step(s) == SQLITE_ROW);
sqlite3_finalize(s);
}
return ok;
}
static int index_exists(sqlite3* db, const char* name) {
sqlite3_stmt* s = NULL;
int ok = 0;
if (sqlite3_prepare_v2(db, "SELECT name FROM sqlite_master WHERE type='index' AND name=?", -1, &s, NULL) == SQLITE_OK) {
sqlite3_bind_text(s, 1, name, -1, SQLITE_STATIC);
ok = (sqlite3_step(s) == SQLITE_ROW);
sqlite3_finalize(s);
}
return ok;
}
static int count_rows(sqlite3* db, const char* table) {
char sql[128]; snprintf(sql, sizeof(sql), "SELECT COUNT(*) FROM %s", table);
sqlite3_stmt* s = NULL;
int n = -1;
if (sqlite3_prepare_v2(db, sql, -1, &s, NULL) == SQLITE_OK) {
if (sqlite3_step(s) == SQLITE_ROW) n = sqlite3_column_int(s, 0);
sqlite3_finalize(s);
}
return n;
}
/* ── test cases ── */
static void test_tables(void) {
TEST("create tables exist"); {
sqlite3* db = make_db();
if (create_tables(db) != 0) { FAIL("create_tables failed"); sqlite3_close(db); return; }
int ba = table_exists(db, "block_availability");
int ss = table_exists(db, "super_sync");
if (ba && ss) OK(); else FAIL("ba=%d ss=%d", ba, ss);
sqlite3_close(db);
}
TEST("create tables idempotent"); {
sqlite3* db = make_db();
if (create_tables(db) != 0) { FAIL("first create"); sqlite3_close(db); return; }
if (create_tables(db) != 0) { FAIL("second create"); sqlite3_close(db); return; }
OK();
sqlite3_close(db);
}
TEST("indices exist"); {
sqlite3* db = make_db();
create_tables(db);
int i1 = index_exists(db, "idx_ba_uuid_node");
int i2 = index_exists(db, "idx_ba_id");
int i3 = index_exists(db, "idx_ba_node_id");
if (i1 && i2 && i3) OK(); else FAIL("i1=%d i2=%d i3=%d", i1, i2, i3);
sqlite3_close(db);
}
}
static void test_ba_insert(void) {
TEST("ba_insert basic"); {
sqlite3* db = make_db(); create_tables(db);
uint8_t uuid[16]; make_uuid(uuid);
uint64_t gid = 0x100, nid = 0x200;
const char* sql = "INSERT INTO block_availability(block_uuid,group_id,node_id,chunk,timestamp) VALUES(?,?,?,?,?)";
for (int i = 0; i < 3; i++) {
sqlite3_stmt* s = NULL;
sqlite3_prepare_v2(db, sql, -1, &s, NULL);
sqlite3_bind_blob(s, 1, uuid, 16, SQLITE_STATIC);
sqlite3_bind_int64(s, 2, (sqlite3_int64)gid);
sqlite3_bind_int64(s, 3, (sqlite3_int64)(nid + (uint64_t)i));
sqlite3_bind_int(s, 4, i);
sqlite3_bind_int64(s, 5, (sqlite3_int64)time(NULL));
sqlite3_step(s); sqlite3_finalize(s);
}
int n = count_rows(db, "block_availability");
if (n == 3) OK(); else FAIL("expected 3 rows got %d", n);
sqlite3_close(db);
}
TEST("ba_insert upsert (unique uuid+node)"); {
sqlite3* db = make_db(); create_tables(db);
uint8_t uuid[16]; make_uuid(uuid);
const char* sql = "INSERT OR REPLACE INTO block_availability(block_uuid,group_id,node_id,chunk,timestamp) VALUES(?,?,?,?,?)";
for (int pass = 0; pass < 2; pass++) {
sqlite3_stmt* s = NULL;
sqlite3_prepare_v2(db, sql, -1, &s, NULL);
sqlite3_bind_blob(s, 1, uuid, 16, SQLITE_STATIC);
sqlite3_bind_int64(s, 2, (sqlite3_int64)1);
sqlite3_bind_int64(s, 3, (sqlite3_int64)0x300);
sqlite3_bind_int(s, 4, 0);
sqlite3_bind_int64(s, 5, (sqlite3_int64)(pass + 100));
sqlite3_step(s); sqlite3_finalize(s);
}
int n = count_rows(db, "block_availability");
if (n == 1) OK(); else FAIL("expected 1 row (upsert) got %d", n);
sqlite3_close(db);
}
}
static void test_ba_find(void) {
TEST("ba_find nodes by uuid"); {
sqlite3* db = make_db(); create_tables(db);
uint8_t uuid[16]; make_uuid(uuid);
const char* ins = "INSERT INTO block_availability(block_uuid,group_id,node_id,chunk,timestamp) VALUES(?,?,?,?,?)";
for (int i = 0; i < 5; i++) {
sqlite3_stmt* s = NULL; sqlite3_prepare_v2(db, ins, -1, &s, NULL);
sqlite3_bind_blob(s, 1, uuid, 16, SQLITE_STATIC);
sqlite3_bind_int64(s, 2, (sqlite3_int64)1);
sqlite3_bind_int64(s, 3, (sqlite3_int64)(0x10 + i));
sqlite3_bind_int(s, 4, i);
sqlite3_bind_int64(s, 5, (sqlite3_int64)(2000 - i));
sqlite3_step(s); sqlite3_finalize(s);
}
const char* q = "SELECT node_id FROM block_availability WHERE block_uuid=? ORDER BY timestamp DESC LIMIT 3";
sqlite3_stmt* s = NULL;
sqlite3_prepare_v2(db, q, -1, &s, NULL);
sqlite3_bind_blob(s, 1, uuid, 16, SQLITE_STATIC);
int count = 0, first_node = -1;
while (sqlite3_step(s) == SQLITE_ROW) { if (count == 0) first_node = sqlite3_column_int(s, 0); count++; }
sqlite3_finalize(s);
if (count == 3 && first_node == 0x10) OK(); else FAIL("count=%d first_node=0x%x", count, first_node);
sqlite3_close(db);
}
}
static void test_ba_get_since(void) {
TEST("ba_get_since id > N"); {
sqlite3* db = make_db(); create_tables(db);
uint8_t uuid[16]; make_uuid(uuid);
const char* ins = "INSERT INTO block_availability(block_uuid,group_id,node_id,chunk,timestamp) VALUES(?,?,?,?,?)";
for (int i = 0; i < 6; i++) {
sqlite3_stmt* s = NULL; sqlite3_prepare_v2(db, ins, -1, &s, NULL);
sqlite3_bind_blob(s, 1, uuid, 16, SQLITE_STATIC);
sqlite3_bind_int64(s, 2, (sqlite3_int64)1);
sqlite3_bind_int64(s, 3, (sqlite3_int64)(0x20 + i));
sqlite3_bind_int(s, 4, i);
sqlite3_bind_int64(s, 5, (sqlite3_int64)1000);
sqlite3_step(s); sqlite3_finalize(s);
}
const char* q = "SELECT id FROM block_availability WHERE id > ? ORDER BY id ASC LIMIT 4";
sqlite3_stmt* s = NULL;
sqlite3_prepare_v2(db, q, -1, &s, NULL);
sqlite3_bind_int64(s, 1, 1);
int count = 0;
while (sqlite3_step(s) == SQLITE_ROW) count++;
sqlite3_finalize(s);
if (count == 4) OK(); else FAIL("expected 4 rows got %d", count);
sqlite3_close(db);
}
}
static void test_ss_ops(void) {
TEST("super_sync insert+select"); {
sqlite3* db = make_db(); create_tables(db);
sqlite3_exec(db, "INSERT OR REPLACE INTO super_sync(peer_node_id,last_recv_id) VALUES(0xAA, 42)", NULL, NULL, NULL);
sqlite3_stmt* s = NULL;
sqlite3_prepare_v2(db, "SELECT last_recv_id FROM super_sync WHERE peer_node_id=?", -1, &s, NULL);
sqlite3_bind_int64(s, 1, 0xAA);
int v = -1;
if (sqlite3_step(s) == SQLITE_ROW) v = sqlite3_column_int(s, 0);
sqlite3_finalize(s);
if (v == 42) OK(); else FAIL("expected 42 got %d", v);
sqlite3_close(db);
}
TEST("super_sync default 0"); {
sqlite3* db = make_db(); create_tables(db);
sqlite3_stmt* s = NULL;
sqlite3_prepare_v2(db, "SELECT last_recv_id FROM super_sync WHERE peer_node_id=?", -1, &s, NULL);
sqlite3_bind_int64(s, 1, 0xCC);
int v = -1;
if (sqlite3_step(s) == SQLITE_ROW) v = sqlite3_column_int(s, 0);
else v = 0; /* no row = default */
sqlite3_finalize(s);
if (v == 0) OK(); else FAIL("expected 0 got %d", v);
sqlite3_close(db);
}
}
static void test_ba_delete(void) {
TEST("ba_delete by node_id"); {
sqlite3* db = make_db(); create_tables(db);
uint8_t uuid[16]; make_uuid(uuid);
const char* ins = "INSERT INTO block_availability(block_uuid,group_id,node_id,chunk,timestamp) VALUES(?,?,?,?,?)";
sqlite3_stmt* s = NULL;
sqlite3_prepare_v2(db, ins, -1, &s, NULL);
sqlite3_bind_blob(s, 1, uuid, 16, SQLITE_STATIC); sqlite3_bind_int64(s, 2, 1);
sqlite3_bind_int64(s, 3, 0x500); sqlite3_bind_int(s, 4, 0);
sqlite3_bind_int64(s, 5, 100); sqlite3_step(s); sqlite3_finalize(s);
sqlite3_prepare_v2(db, ins, -1, &s, NULL);
sqlite3_bind_blob(s, 1, uuid, 16, SQLITE_STATIC); sqlite3_bind_int64(s, 2, 1);
sqlite3_bind_int64(s, 3, 0x600); sqlite3_bind_int(s, 4, 0);
sqlite3_bind_int64(s, 5, 100); sqlite3_step(s); sqlite3_finalize(s);
if (count_rows(db, "block_availability") != 2) { FAIL("setup failed"); sqlite3_close(db); return; }
sqlite3_exec(db, "DELETE FROM block_availability WHERE node_id=0x500", NULL, NULL, NULL);
int n = count_rows(db, "block_availability");
if (n == 1) OK(); else FAIL("expected 1 got %d", n);
sqlite3_close(db);
}
TEST("ba_delete foreign (node_id != self)"); {
sqlite3* db = make_db(); create_tables(db);
uint8_t uuid[16]; make_uuid(uuid);
const char* ins = "INSERT INTO block_availability(block_uuid,group_id,node_id,chunk,timestamp) VALUES(?,?,?,?,?)";
sqlite3_stmt* s = NULL;
for (int i = 0; i < 3; i++) {
sqlite3_prepare_v2(db, ins, -1, &s, NULL);
sqlite3_bind_blob(s, 1, uuid, 16, SQLITE_STATIC); sqlite3_bind_int64(s, 2, 1);
sqlite3_bind_int64(s, 3, (sqlite3_int64)(0x700 + i)); sqlite3_bind_int(s, 4, i);
sqlite3_bind_int64(s, 5, 100); sqlite3_step(s); sqlite3_finalize(s);
}
char buf[128]; snprintf(buf, sizeof(buf), "DELETE FROM block_availability WHERE node_id!=0x700");
sqlite3_exec(db, buf, NULL, NULL, NULL);
int n = count_rows(db, "block_availability");
if (n == 1) OK(); else FAIL("expected 1 (self) got %d", n);
sqlite3_close(db);
}
}
static void test_proto_sizes(void) {
TEST("packed struct sizes"); {
int ok = 1;
if (sizeof(struct media_pkt_query) != 27) { FAIL("query: %zu != 27", sizeof(struct media_pkt_query)); ok = 0; }
if (sizeof(struct media_pkt_have_block) != 117) { FAIL("have_block: %zu != 117", sizeof(struct media_pkt_have_block)); ok = 0; }
if (sizeof(struct media_pkt_super_hello) != 9) { FAIL("super_hello: %zu != 9", sizeof(struct media_pkt_super_hello)); ok = 0; }
if (sizeof(struct media_pkt_super_ack) != 5) { FAIL("super_ack: %zu != 5", sizeof(struct media_pkt_super_ack)); ok = 0; }
if (sizeof(struct media_pkt_block_req) != 45) { FAIL("block_req: %zu != 45", sizeof(struct media_pkt_block_req)); ok = 0; }
if (sizeof(struct media_pkt_cancel) != 37) { FAIL("cancel: %zu != 37", sizeof(struct media_pkt_cancel)); ok = 0; }
if (sizeof(struct media_pkt_block_done) != 105) { FAIL("block_done: %zu != 105", sizeof(struct media_pkt_block_done)); ok = 0; }
if (sizeof(struct media_pkt_query_resp_entry) != 30) { FAIL("query_resp_entry: %zu != 30", sizeof(struct media_pkt_query_resp_entry)); ok = 0; }
if (ok) OK();
/* continue running even if sizes mismatch */
}
}
static void test_proto_build(void) {
TEST("build MEDIA_QUERY"); {
uint8_t mid[16]; memset(mid, 0xAB, 16);
uint8_t bids[32]; memset(bids, 0xCD, 32);
uint8_t pkt[256];
struct media_pkt_query* q = (struct media_pkt_query*)pkt;
q->subcmd = MEDIA_SUBCMD_QUERY;
q->group_id = 0x12345678;
memcpy(q->media_id, mid, 16);
q->num_blocks = 2;
memcpy(pkt + MEDIA_QUERY_HDR_SIZE, bids, 32);
if (q->subcmd == MEDIA_SUBCMD_QUERY && q->group_id == 0x12345678 &&
q->num_blocks == 2 && memcmp(q->media_id, mid, 16) == 0 &&
memcmp(pkt + MEDIA_QUERY_HDR_SIZE, bids, 32) == 0) OK();
else FAIL();
}
TEST("build MEDIA_HAVE_BLOCK"); {
struct media_pkt_have_block hb;
memset(&hb, 0, sizeof(hb));
hb.subcmd = MEDIA_SUBCMD_HAVE_BLOCK;
hb.group_id = 0xAA;
hb.chunk = 5;
hb.timestamp = 999;
if (hb.subcmd == MEDIA_SUBCMD_HAVE_BLOCK && hb.group_id == 0xAA &&
hb.chunk == 5 && hb.timestamp == 999) OK();
else FAIL();
}
TEST("build MEDIA_SUPER_HELLO"); {
struct media_pkt_super_hello sh;
sh.subcmd = MEDIA_SUBCMD_SUPER_HELLO;
sh.last_recv_id = 12345;
if (sh.subcmd == MEDIA_SUBCMD_SUPER_HELLO && sh.last_recv_id == 12345) OK();
else FAIL();
}
}
int main(void) {
debug_config_init(); debug_set_level(DEBUG_LEVEL_ERROR);
srand((unsigned)time(NULL));
printf("=== test_media_delivery_sql ===\n");
test_tables();
test_ba_insert();
test_ba_find();
test_ba_get_since();
test_ss_ops();
test_ba_delete();
test_proto_sizes();
test_proto_build();
printf("\n%d/%d passed, %d failed\n", g_passed, g_total, g_failed);
return g_failed > 0 ? 1 : 0;
}
Loading…
Cancel
Save