Browse Source

media_delivery: обновление и новый тест conn_mgr_handles

topo_upd
evgeny 2 months ago
parent
commit
153483fd96
  1. 19
      src/media_delivery/media_delivery.c
  2. 1
      src/media_delivery/media_delivery.h
  3. 154
      tests/test_conn_mgr_handles.c
  4. 153
      tests/test_media_delivery_full.c

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

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

154
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 <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <stdarg.h>
#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;
}

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

Loading…
Cancel
Save