diff --git a/src/media_delivery/media_delivery.c b/src/media_delivery/media_delivery.c index 9dba170a..769e07ff 100644 --- a/src/media_delivery/media_delivery.c +++ b/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)); diff --git a/src/media_delivery/media_download.c b/src/media_delivery/media_download.c index eb96a9e3..84b8b04a 100644 --- a/src/media_delivery/media_download.c +++ b/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; diff --git a/src/media_delivery/media_download.h b/src/media_delivery/media_download.h index 9677f51e..188b12aa 100644 --- a/src/media_delivery/media_download.h +++ b/src/media_delivery/media_download.h @@ -9,6 +9,10 @@ extern "C" { #include #include #include +#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]; diff --git a/tests/Makefile.am b/tests/Makefile.am index d9e44799..9bb7c84f 100644 --- a/tests/Makefile.am +++ b/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 diff --git a/tests/test_media_delivery_download.c b/tests/test_media_delivery_download.c new file mode 100644 index 00000000..76675603 --- /dev/null +++ b/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 +#include +#include +#include +#include +#include +#include + +#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; +} diff --git a/tests/test_media_delivery_sql.c b/tests/test_media_delivery_sql.c new file mode 100644 index 00000000..7fbac041 --- /dev/null +++ b/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 +#include +#include +#include +#include + +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; +}