diff --git a/src/media_delivery/media_delivery.c b/src/media_delivery/media_delivery.c index 136a6c35..09a07eb5 100644 --- a/src/media_delivery/media_delivery.c +++ b/src/media_delivery/media_delivery.c @@ -631,14 +631,14 @@ static void md_handle_block_req(struct media_delivery_ctx* md, uint64_t from_nod /* find file in media_files DB */ sqlite3* db = md->db; - if (!db) return; + if (!db) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: BLOCK_REQ no DB", MD_ID); md->active_streams--; return; } const char* sql = "SELECT location,chunk_size,chunk,offset,file_size FROM media_files WHERE block_id=? AND node_id=? LIMIT 1"; 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: BLOCK_REQ prepare failed: %s", MD_ID, sqlite3_errmsg(db)); md->active_streams--; return; } sqlite3_bind_blob(st, 1, req->block_id, 16, SQLITE_STATIC); sqlite3_bind_int64(st, 2, (sqlite3_int64)md->inst->node_id); - if (sqlite3_step(st) != SQLITE_ROW) { sqlite3_finalize(st); return; } + if (sqlite3_step(st) != SQLITE_ROW) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: BLOCK_REQ block_id=%02x%02x... not found in media_files", MD_ID, req->block_id[0], req->block_id[1]); sqlite3_finalize(st); md->active_streams--; return; } const char* loc = (const char*)sqlite3_column_text(st, 0); int64_t chunk_size = sqlite3_column_int64(st, 1); @@ -691,6 +691,7 @@ static void md_handle_block_req(struct media_delivery_ctx* md, uint64_t from_nod DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: streaming file=%s chunk=%d start=%lld len=%lld group=%016llx to 0x%016llx", MD_ID, path, chunk, (long long)block_start, (long long)block_length, (unsigned long long)group_id, (unsigned long long)from_node); + md->streams_started++; /* monotonic — for test admission checks */ /* start streaming — first chunk immediately, then backpressure */ memset(&sc->waiter, 0, sizeof(sc->waiter)); @@ -1037,8 +1038,18 @@ void media_delivery_destroy(struct UTUN_INSTANCE* inst) { gle = gle->next; } + if (md->super_peers) { + struct ll_entry* se = md->super_peers->head; + while (se) { + struct media_super_peer* sp = (struct media_super_peer*)se->data; + if (sp->timeout_timer) { uasync_cancel_timeout(md->inst->ua, sp->timeout_timer); sp->timeout_timer = NULL; } + if (sp->connect_timer) { uasync_cancel_timeout(md->inst->ua, sp->connect_timer); sp->connect_timer = NULL; } + se = se->next; + } + queue_free(md->super_peers); + md->super_peers = NULL; + } if (md->served_nodes) { queue_free(md->served_nodes); md->served_nodes = NULL; } - if (md->super_peers) { queue_free(md->super_peers); md->super_peers = NULL; } if (md->downloads) { queue_free(md->downloads); md->downloads = NULL; } memset(md, 0, sizeof(*md)); diff --git a/src/media_delivery/media_delivery.h b/src/media_delivery/media_delivery.h index a224c269..dd924187 100644 --- a/src/media_delivery/media_delivery.h +++ b/src/media_delivery/media_delivery.h @@ -51,6 +51,7 @@ struct media_delivery_ctx { uint8_t is_supernode; // из adm_tags в peers_* uint8_t active_streams; // текущее число исходящих стримов (admission) uint8_t stream_completed; // 1 = хотя бы один стрим завершился (для тестов) + uint8_t streams_started; // монотонный счётчик запущенных стримов (для тестов) sqlite3* db; // = inst->topo_sqlite_db (кешируем для удобства) struct ll_queue* served_nodes; // media_served_node — только суперузел struct ll_queue* super_peers; // media_super_peer — только суперузел diff --git a/tests/test_conn_mgr_handles.c b/tests/test_conn_mgr_handles.c new file mode 100644 index 00000000..1c71d9fe --- /dev/null +++ b/tests/test_conn_mgr_handles.c @@ -0,0 +1,154 @@ +/** + * @file test_conn_mgr_handles.c + * @brief conn_mgr handle lifecycle: multiple handles sharing one connection + * + * Test 1: Two conn_mgr_connect_node → both OK, same entry, one conn + * Test 2: Release first handle → conn stays (second handle still active) + * Test 3: Release last handle → conn torn down + */ +#include +#include +#include +#include +#include "../lib/platform_compat.h" +#include "test_utils.h" + +#include "etcp.h" +#include "etcp_connections.h" +#include "../src/config_parser.h" +#include "../src/config_updater.h" +#include "../src/utun_instance.h" +#include "topo_group.h" +#include "conn_mgr.h" +#include "../lib/u_async.h" +#include "../lib/debug_config.h" +#include "../lib/mem.h" + +#define TIMEOUT_TB 300000 +#define POLL_MS 5 +static struct UTUN_INSTANCE *g_a = NULL, *g_b = NULL; +static struct UASYNC *ua = NULL; +static volatile int result = 0; +static void* ttimer = NULL; +static char tdir[] = "/tmp/utun_cmh_XXXXXX"; +static char ca[256], cb[256]; +static int pa = 0, pb = 0; +static uint64_t g_nid_a = 0, g_nid_b = 0; + +static volatile int cb1_ok = 0, cb2_ok = 0; +static struct CONN_MGR_HANDLE *gh1 = NULL, *gh2 = NULL; + +static int wf(const char* p, const char* f, ...) { + va_list ap; FILE* fp = fopen(p, "w"); if (!fp) return -1; + va_start(ap, f); vfprintf(fp, f, ap); va_end(ap); fclose(fp); return 0; +} +static char* gv(const char* p, const char* k) { + struct utun_config* c = parse_config(p); if (!c) return NULL; + char* r = (strcmp(k, "pub") == 0) ? u_strdup(c->global.my_public_key_hex) : u_strdup(c->global.my_private_key_hex); + free_config(c); return r; +} +static void fail(const char* msg) { fprintf(stderr, "FAIL: %s\n", msg); fflush(stderr); result = 2; } +static void done(void) { if (result == 0) result = 1; } +static void to_cb(void* arg) { (void)arg; fail("timeout"); done(); } + +static void ccb1(int r, uint64_t id, void* arg) { + (void)arg; cb1_ok = (r == CONN_MGR_OK); + fprintf(stderr, "ccb1: result=%d node=0x%llx ok=%d\n", r, (unsigned long long)id, cb1_ok); fflush(stderr); +} +static void ccb2(int r, uint64_t id, void* arg) { + (void)arg; cb2_ok = (r == CONN_MGR_OK); + fprintf(stderr, "ccb2: result=%d node=0x%llx ok=%d\n", r, (unsigned long long)id, cb2_ok); fflush(stderr); +} + +static void test1_check(void* arg); +static void test2_check(void* arg); +static void test2_verify(void* arg); +static void test3_check(void* arg); +static void test3_verify(void* arg); + +static void test1(void* arg) { + (void)arg; if (result) return; + if (!topo_node_find_by_id(topo_groups_get_default(g_a->topo_groups), g_nid_b)) { uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test1, "t1w"); return; } + struct ETCP_CONN* c = topo_group_find_conn_for_node(g_a->conn_mgr->group, g_nid_b); + if (!c) c = instance_find_conn(g_a, g_nid_b); + if (!c || !c->links) { uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test1, "t1w2"); return; } + int has = 0; struct ETCP_LINK* l = c->links; + while (l) { if (l->link_state == 3 && l->initialized) { has = 1; break; } l = l->next; } + if (!has) { uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test1, "t1w3"); return; } + fprintf(stderr, "Test 1: two conn_mgr_connect_node to same node\n"); fflush(stderr); + cb1_ok = 0; cb2_ok = 0; + gh1 = conn_mgr_connect_node(g_a->conn_mgr, g_nid_b, 0, ccb1, NULL); + gh2 = conn_mgr_connect_node(g_a->conn_mgr, g_nid_b, 0, ccb2, NULL); + if (!gh1 || !gh2) { fail("test1: handle returned NULL"); done(); return; } + uasync_set_timeout(ua, 200, NULL, (timeout_callback_t)test1_check, "t1c"); +} +static void test1_check(void* arg) { + (void)arg; + if (!cb1_ok || !cb2_ok) { fail("test1: both callbacks should be OK"); done(); return; } + uint8_t st, ty; conn_mgr_get_status(g_a->conn_mgr, g_nid_b, &st, &ty); + if (st != CONN_MGR_STATE_CONNECTED || ty != CONN_TYPE_DIRECT) { fail("test1: status mismatch"); done(); return; } + fprintf(stderr, " OK: both handles returned OK, state=%d type=%d\n", st, ty); fflush(stderr); + uasync_call_soon(ua, NULL, (timeout_callback_t)test2_check); +} + +static void test2_check(void* arg) { + (void)arg; + fprintf(stderr, "Test 2: release first handle — conn stays alive\n"); fflush(stderr); + conn_mgr_release_handle(gh1); gh1 = NULL; + uasync_set_timeout(ua, 100, NULL, (timeout_callback_t)test2_verify, "t2v"); +} +static void test2_verify(void* arg) { + (void)arg; + uint8_t st, ty; conn_mgr_get_status(g_a->conn_mgr, g_nid_b, &st, &ty); + if (st != CONN_MGR_STATE_CONNECTED) { fail("test2: conn died after first handle close"); done(); return; } + fprintf(stderr, " OK: conn alive after h1 close, state=%d\n", st); fflush(stderr); + uasync_call_soon(ua, NULL, (timeout_callback_t)test3_check); +} + +static void test3_check(void* arg) { + (void)arg; + fprintf(stderr, "Test 3: release last handle — conn torn down\n"); fflush(stderr); + conn_mgr_release_handle(gh2); gh2 = NULL; + uasync_set_timeout(ua, 200, NULL, (timeout_callback_t)test3_verify, "t3v"); +} +static void test3_verify(void* arg) { + (void)arg; + uint8_t st, ty; conn_mgr_get_status(g_a->conn_mgr, g_nid_b, &st, &ty); + if (st != CONN_MGR_STATE_DISCONNECTED) { fail("test3: conn should be disconnected"); done(); return; } + fprintf(stderr, " OK: conn DISCONNECTED after last handle close\n"); fflush(stderr); + fprintf(stderr, "=== ALL DONE ===\n"); fflush(stderr); + done(); +} + +int main(void) { + debug_config_init(); debug_set_level(DEBUG_LEVEL_WARN); + utun_instance_set_tun_init_enabled(0); + + test_mkdtemp(tdir); + int base = 50000 + (getpid() % 10000); pa = base; pb = base + 1; + snprintf(ca, sizeof(ca), "%s/a.conf", tdir); snprintf(cb, sizeof(cb), "%s/b.conf", tdir); + wf(ca, "[global]\ntun_ip=10.96.0.1/24\ntun_ifname=tun93\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", pa); + wf(cb, "[global]\ntun_ip=10.96.0.2/24\ntun_ifname=tun92\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", pb); + config_ensure_keys_and_node_id(ca); config_ensure_keys_and_node_id(cb); + { struct utun_config* cfa = parse_config(ca); struct utun_config* cfb = parse_config(cb); + g_nid_a = cfa->global.my_node_id; g_nid_b = cfb->global.my_node_id; + free_config(cfa); free_config(cfb); } + char *p0 = gv(ca,"pub"), *r0 = gv(ca,"priv"), *p1 = gv(cb,"pub"); + wf(ca, "[global]\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.96.0.1/24\ntun_ifname=tun93\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[client: to_b]\nkeepalive=1\npeer_public_key=%s\nlink=s1:127.0.0.1:%d\n[allowed_keys]\nallow_all=1\n", r0, p0, pa, p1, pb); + u_free(p0); u_free(r0); u_free(p1); + + ua = uasync_create(); + g_a = utun_instance_create(ua, ca); g_b = utun_instance_create(ua, cb); + if (!g_a || !g_b) { fprintf(stderr,"create failed\n"); result=2; goto done; } + utun_instance_init(g_a); utun_instance_init(g_b); + uasync_call_soon(ua, NULL, test1); + ttimer = uasync_set_timeout(ua, TIMEOUT_TB, NULL, to_cb, "to"); + { int el = 0; while (!result && el < TIMEOUT_TB/10 + 500) { uasync_poll(ua, POLL_MS); el += POLL_MS; } } + +done: + if (ttimer) { uasync_cancel_timeout(ua, ttimer); } + g_a->running = 0; if (g_b) g_b->running = 0; + if (ua) uasync_destroy(ua, 1); // force cleanup + test_unlink(ca); test_unlink(cb); test_rmdir(tdir); + return (result == 1) ? 0 : 1; +} diff --git a/tests/test_media_delivery_full.c b/tests/test_media_delivery_full.c index f54470dd..fb5a4164 100644 --- a/tests/test_media_delivery_full.c +++ b/tests/test_media_delivery_full.c @@ -93,8 +93,10 @@ static int msend(struct UTUN_INSTANCE* inst, uint64_t dst, const uint8_t* data, struct ll_entry* e = queue_entry_new(0); if (!e) return -1; e->dgram = u_malloc(len + 1); if (!e->dgram) { queue_entry_free(e); return -1; } e->dgram[0] = ETCP_RT_ID_MEDIA_DELIVERY; memcpy(e->dgram + 1, data, len); e->len = (uint16_t)(len + 1); + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "[msend] 0x%016llx → 0x%016llx svc=0x%02x subcmd=0x%02x len=%zu", + (unsigned long long)inst->node_id, (unsigned long long)dst, e->dgram[0], data[0], len); int rc = etcp_route_send(inst, TOPO_GROUP_UTUN, dst, e, 1); - if (rc != 0) { u_free(e->dgram); queue_entry_free(e); } + if (rc != 0) { DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "[msend] FAILED rc=%d", rc); u_free(e->dgram); queue_entry_free(e); } return rc; } @@ -197,98 +199,106 @@ static void phase_a2_have_block_replication(void) { } } -/* ── Phase A3: admission control ── */ +/* ── helper: create test file + index before admission ── */ +static int create_test_file(void) { + char src_tmp[512]; snprintf(src_tmp, sizeof(src_tmp), "/tmp/mdl_test_%d.bin", getpid()); + char media_dir[512]; snprintf(media_dir, sizeof(media_dir), "%s/media", g_db_dir[0]); + utun_mkdir(media_dir, 0755); + char dst_path[512]; snprintf(dst_path, sizeof(dst_path), "%s/test_src.bin", media_dir); + char media_base[512]; snprintf(media_base, sizeof(media_base), "%s", g_db_dir[0]); + + uint8_t file_data[10240]; + for (int i = 0; i < 10240; i++) file_data[i] = (uint8_t)(i & 0xFF); + FILE* f = fopen(src_tmp, "wb"); + if (!f) return -1; + fwrite(file_data, 1, 10240, f); fclose(f); + + { FILE* fin = fopen(src_tmp, "rb"), *fout = fopen(dst_path, "wb"); + if (fin && fout) { uint8_t buf[4096]; size_t rd; while ((rd = fread(buf, 1, sizeof(buf), fin)) > 0) fwrite(buf, 1, rd, fout); } + if (fin) fclose(fin); if (fout) fclose(fout); } + + /* set media_base in ui_state */ + sqlite3_exec(g_inst[0]->topo_sqlite_db, "CREATE TABLE IF NOT EXISTS ui_state(key TEXT PRIMARY KEY, value TEXT)", NULL, NULL, NULL); + { sqlite3_stmt* us = NULL; + sqlite3_prepare_v2(g_inst[0]->topo_sqlite_db, "INSERT OR REPLACE INTO ui_state(key,value) VALUES('media_base',?)", -1, &us, NULL); + sqlite3_bind_text(us, 1, media_base, -1, SQLITE_STATIC); sqlite3_step(us); sqlite3_finalize(us); } + + /* index the file */ + uint8_t hash[32]; + { EVP_MD_CTX* ctx = EVP_MD_CTX_new(); EVP_DigestInit_ex(ctx, EVP_sha256(), NULL); + EVP_DigestUpdate(ctx, file_data, 10240); EVP_DigestFinal_ex(ctx, hash, NULL); EVP_MD_CTX_free(ctx); } + + media_index_init(g_inst[0]->topo_sqlite_db); + struct media_index_result result; memset(&result, 0, sizeof(result)); + media_index_generate_uuid(result.media_id); + memcpy(result.content_hash, hash, 32); + result.file_size = 10240; result.block_size = 5120; result.num_blocks = 2; + result.block_ids = u_malloc(32); result.block_sigs = u_malloc(128); + for (int i = 0; i < 2; i++) { + media_index_generate_uuid(result.block_ids + i * 16); + uint8_t smsg[5128]; size_t soff = 0; + memcpy(smsg + soff, file_data + i * 5120, 5120); soff += 5120; + uint64_t nid = g_nid[0]; memcpy(smsg + soff, &nid, 8); soff += 8; + sc_ed25519_sign(g_inst[0]->my_ed25519_privkey, smsg, soff, result.block_sigs + i * 64); + } + int rc = media_index_commit(g_inst[0]->topo_sqlite_db, &result, g_nid[0], + g_inst[0]->my_ed25519_privkey, "test_ch", dst_path, media_base); + if (rc != 0) { media_index_result_free(&result); return -1; } + memcpy(g_test_mid, result.media_id, 16); + memcpy(g_test_bid0, result.block_ids, 16); + memcpy(g_test_bid1, result.block_ids + 16, 16); + media_index_result_free(&result); + return 0; +} + +/* ── Phase A3: admission control (uses real file after create_test_file) ── */ static void phase_a3_admission(void) { + /* wait for UTUN-group routing to be established (NODEINFO exchange complete) */ + { int a = 0; while (a < 1000) { + struct TOPO_GROUP* g = topo_groups_find(g_inst[1]->topo_groups, TOPO_GROUP_UTUN); + struct ETCP_CONN* c = g ? topo_group_find_conn_for_node(g, g_nid[0]) : NULL; + if (c) break; + uasync_poll(g_ua, POLL_MS); a++; + }} + TEST("1st BLOCK_REQ → stream"); { - uint8_t mid[16], bid[16]; memset(mid, 0xCD, 16); memset(bid, 0xCE, 16); struct media_pkt_block_req req; memset(&req, 0, sizeof(req)); req.subcmd = MEDIA_SUBCMD_BLOCK_REQ; req.group_id = g_group_id; - memcpy(req.media_id, mid, 16); memcpy(req.block_id, bid, 16); + memcpy(req.media_id, g_test_mid, 16); memcpy(req.block_id, g_test_bid0, 16); msend(g_inst[1], g_nid[0], (const uint8_t*)&req, sizeof(req)); - int a = 0; while (a < 200) { uasync_poll(g_ua, POLL_MS); a++; } - if (g_inst[0]->md.active_streams >= 1) OK(); else FAIL("active=%d", g_inst[0]->md.active_streams); + int a = 0; while (a < 500) { uasync_poll(g_ua, POLL_MS); a++; } + if (g_inst[0]->md.streams_started >= 1) OK(); else FAIL("started=%d", g_inst[0]->md.streams_started); } TEST("fill → 3 → OVERLOADED"); { + int before = g_inst[0]->md.streams_started; for (int i = 0; i < 3; i++) { - uint8_t mid[16], bid[16]; memset(mid, i+10, 16); memset(bid, i+20, 16); struct media_pkt_block_req req; memset(&req, 0, sizeof(req)); req.subcmd = MEDIA_SUBCMD_BLOCK_REQ; req.group_id = g_group_id; req.chunk = (uint32_t)i; - memcpy(req.media_id, mid, 16); memcpy(req.block_id, bid, 16); + memcpy(req.media_id, g_test_mid, 16); memcpy(req.block_id, g_test_bid0, 16); msend(g_inst[1], g_nid[0], (const uint8_t*)&req, sizeof(req)); } - int a = 0; while (a < 200) { uasync_poll(g_ua, POLL_MS); a++; } - if (g_inst[0]->md.active_streams == 3) OK(); else FAIL("active=%d", g_inst[0]->md.active_streams); + int a = 0; while (a < 500) { uasync_poll(g_ua, POLL_MS); a++; } + int started = g_inst[0]->md.streams_started - before; + if (started == 3) OK(); else FAIL("started=%d (expected 3, total=%d)", started, g_inst[0]->md.streams_started); } g_inst[0]->md.active_streams = 0; } /* ══════════════════════════════════════════════════════════ - Phase B1: create file → register → stream → assemble → verify + Phase B1: BLOCK_REQ → stream → verify ══════════════════════════════════════════════════════════ */ -static void phase_b1_file_transfer(void) { - char src_tmp[512]; snprintf(src_tmp, sizeof(src_tmp), "/tmp/mdl_test_%d.bin", getpid()); - char media_dir[512]; snprintf(media_dir, sizeof(media_dir), "%s/media", g_db_dir[0]); - utun_mkdir(media_dir, 0755); - char dst_path[512]; snprintf(dst_path, sizeof(dst_path), "%s/test_src.bin", media_dir); - char media_base[512]; snprintf(media_base, sizeof(media_base), "%s", g_db_dir[0]); - - uint8_t file_data[2048]; - for (int i = 0; i < 2048; i++) file_data[i] = (uint8_t)(i & 0xFF); - FILE* f = fopen(src_tmp, "wb"); - if (f) { fwrite(file_data, 1, 2048, f); fclose(f); } - - { FILE* fin = fopen(src_tmp, "rb"), *fout = fopen(dst_path, "wb"); - if (fin && fout) { uint8_t buf[4096]; size_t rd; while ((rd = fread(buf, 1, sizeof(buf), fin)) > 0) fwrite(buf, 1, rd, fout); } - if (fin) fclose(fin); if (fout) fclose(fout); } - - TEST("media_index_commit + ui_state"); { - uint8_t hash[32]; - EVP_MD_CTX* ctx = EVP_MD_CTX_new(); - EVP_DigestInit_ex(ctx, EVP_sha256(), NULL); - EVP_DigestUpdate(ctx, file_data, 2048); - EVP_DigestFinal_ex(ctx, hash, NULL); - EVP_MD_CTX_free(ctx); - - media_index_init(g_inst[0]->topo_sqlite_db); - sqlite3_exec(g_inst[0]->topo_sqlite_db, "CREATE TABLE IF NOT EXISTS ui_state(key TEXT PRIMARY KEY, value TEXT)", NULL, NULL, NULL); - { sqlite3_stmt* us = NULL; - sqlite3_prepare_v2(g_inst[0]->topo_sqlite_db, "INSERT OR REPLACE INTO ui_state(key,value) VALUES('media_base',?)", -1, &us, NULL); - sqlite3_bind_text(us, 1, media_base, -1, SQLITE_STATIC); sqlite3_step(us); sqlite3_finalize(us); } - - struct media_index_result result; memset(&result, 0, sizeof(result)); - media_index_generate_uuid(result.media_id); - memcpy(result.content_hash, hash, 32); - result.file_size = 2048; result.block_size = 1024; result.num_blocks = 2; - result.block_ids = u_malloc(32); result.block_sigs = u_malloc(128); - for (int i = 0; i < 2; i++) { - media_index_generate_uuid(result.block_ids + i * 16); - uint8_t smsg[2048]; size_t soff = 0; - memcpy(smsg + soff, file_data + i * 1024, 1024); soff += 1024; - uint64_t nid = g_nid[0]; memcpy(smsg + soff, &nid, 8); soff += 8; - sc_ed25519_sign(g_inst[0]->my_ed25519_privkey, smsg, soff, result.block_sigs + i * 64); - } - int rc = media_index_commit(g_inst[0]->topo_sqlite_db, &result, g_nid[0], - g_inst[0]->my_ed25519_privkey, "test_ch", dst_path, media_base); - if (rc != 0) { FAIL("commit rc=%d", rc); media_index_result_free(&result); return; } - memcpy(g_test_mid, result.media_id, 16); - memcpy(g_test_bid0, result.block_ids, 16); - memcpy(g_test_bid1, result.block_ids + 16, 16); - int ndb = db_count(g_inst[0]->topo_sqlite_db, "media_files", NULL, 0); - media_index_result_free(&result); - if (ndb == 2) OK(); else FAIL("media_files has %d rows (expected 2)", ndb); - } - +static void phase_b1_stream_test(void) { TEST("n2→n1 BLOCK_REQ streams data"); { struct media_pkt_block_req req; memset(&req, 0, sizeof(req)); req.subcmd = MEDIA_SUBCMD_BLOCK_REQ; req.group_id = g_group_id; - memcpy(req.media_id, g_test_mid, 16); memcpy(req.block_id, g_test_bid0, 16); + memcpy(req.media_id, g_test_mid, 16); memcpy(req.block_id, g_test_bid1, 16); req.chunk = 0; g_inst[0]->md.stream_completed = 0; msend(g_inst[1], g_nid[0], (const uint8_t*)&req, sizeof(req)); int a = 0; - while (a < 2000 && !g_inst[0]->md.stream_completed) { uasync_poll(g_ua, POLL_MS); a++; } - if (g_inst[0]->md.stream_completed) OK(); else FAIL("stream not completed after %dms", 2000 * POLL_MS); + while (a < 3000 && !g_inst[0]->md.stream_completed) { uasync_poll(g_ua, POLL_MS); a++; } + if (g_inst[0]->md.stream_completed) OK(); else FAIL("stream not completed after %dms", 3000 * POLL_MS); } } @@ -353,8 +363,15 @@ int main(void) { phase_a1_super_hello(); phase_a2_have_block_replication(); + + /* create test file + index before admission tests (needed for real BLOCK_REQ to find files) */ + { TEST("media_index_commit + ui_state"); { + if (create_test_file() == 0 && db_count(g_inst[0]->topo_sqlite_db, "media_files", NULL, 0) == 2) OK(); + else FAIL("create_test_file failed"); + }} + phase_a3_admission(); - phase_b1_file_transfer(); + phase_b1_stream_test(); fflush(stdout); fflush(stderr); printf("\n%d/%d passed, %d failed\n", G_PASSED, G_TOTAL, G_FAILED); fflush(stdout);