You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 
 
 

302 lines
12 KiB

// test_media_delivery_integration.c — интеграционный тест media_delivery
//
// Проверки:
// 1. media_delivery_init на реальных инстансах (таблицы созданы)
// 2. media_delivery_set_supernode → super_peers/served_nodes созданы
// 3. SERVE_REG → добавление в served_nodes
// 4. HAVE_BLOCK → вставка в block_availability + ACK
// 5. QUERY → поиск в block_availability → QUERY_RESP
// 6. SUPER_HELLO → обмен last_recv_id
// 7. SERVE_LEAVE → удаление из served_nodes
// 8. Выключение суперузла → чистка block_availability
#include "media_delivery.h"
#include "media_delivery_proto.h"
#include "../utun_instance.h"
#include "../transport_layer/etcp.h"
#include "../config_parser.h"
#include "../config_updater.h"
#include "test_utils.h"
#include "../routing_layer/topo_group.h"
#include "../routing_layer/topo_node.h"
#include "../lib/u_async.h"
#include "../lib/debug_config.h"
#include "../lib/mem.h"
#include <sqlite3.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <time.h>
static int g_passed = 0, g_failed = 0, g_total = 0;
#define TEST(name) do { g_total++; printf(" %-55s", name); fflush(stdout); } while(0)
#define OK() do { g_passed++; printf("OK\n"); } while(0)
#define FAIL(fmt, ...) do { g_failed++; printf("FAIL: " fmt "\n", ##__VA_ARGS__); } while(0)
#define TIMEOUT_TB 300000
#define POLL_MS 5
static struct UASYNC* g_ua = NULL;
static struct UTUN_INSTANCE* g_a = NULL, *g_b = NULL;
static uint64_t g_nid_a = 0, g_nid_b = 0;
static int g_connected = 0, g_result = 0;
static char g_tdir[] = "/tmp/utun_mdl_XXXXXX";
static char g_ca[256], g_cb[256];
static uint8_t g_a_pubkey[32], g_b_pubkey[32];
static int g_pa = 0, g_pb = 0;
/* ── helpers ── */
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;
}
/* ── helpers ── */
static int count_rows(sqlite3* db, const char* tbl, const char* where, uint64_t val) {
char sql[256];
if (where) snprintf(sql, sizeof(sql), "SELECT COUNT(*) FROM %s WHERE %s=%llu", tbl, where, (unsigned long long)val);
else snprintf(sql, sizeof(sql), "SELECT COUNT(*) FROM %s", tbl);
sqlite3_stmt* s = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &s, NULL) != SQLITE_OK) return -1;
int n = -1;
if (sqlite3_step(s) == SQLITE_ROW) n = sqlite3_column_int(s, 0);
sqlite3_finalize(s);
return n;
}
static void to_cb(void* arg) { (void)arg; fprintf(stderr, "TIMEOUT\n"); g_result = 2; }
/* ── monitor: ждём ETCP и BGP ── */
static void mon(void* arg) {
struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)arg;
(void)inst;
if (!g_connected) {
int ok = 0;
{ struct ll_entry* e = g_a->connections->head;
while (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data;
if (ce->conn && ce->conn->links) { struct ETCP_LINK* l; for (l = ce->conn->links; l; l = l->next)
if (l->initialized && ce->conn->crypto_ctx.initialized) ok++; } e = e->next; } }
{ struct ll_entry* e = g_b->connections->head;
while (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data;
if (ce->conn && ce->conn->links) { struct ETCP_LINK* l; for (l = ce->conn->links; l; l = l->next)
if (l->initialized && ce->conn->crypto_ctx.initialized) ok++; } e = e->next; } }
if (ok >= 2) g_connected = 1;
}
if (!g_result) uasync_set_timeout(g_ua, 10, NULL, mon, "mdl_mon");
}
/* ── send helper ── */
static int msend(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;
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);
int rc = etcp_route_send(inst, TOPO_GROUP_UTUN, dst, e, 1);
if (rc != 0) { u_free(e->dgram); queue_entry_free(e); }
return rc;
}
/* ── tests ── */
static void phase1_serve_reg(void) {
TEST("SERVE_REG B→A adds to served_nodes"); {
struct media_pkt_serve_reg pkt;
memset(&pkt, 0, sizeof(pkt));
pkt.subcmd = MEDIA_SUBCMD_SERVE_REG;
pkt.group_id = TOPO_GROUP_UTUN;
int rc = msend(g_b, g_nid_a, (const uint8_t*)&pkt, sizeof(pkt));
/* BGP routes need time — poll up to 10s */
int attempts = 0;
while (attempts < 2000) {
int has = g_a->md.served_nodes && g_a->md.served_nodes->head != NULL;
if (has) break;
uasync_poll(g_ua, POLL_MS); attempts++;
}
int ok = g_a->md.served_nodes && g_a->md.served_nodes->head != NULL;
if (ok) OK(); else FAIL("rc=%d B not in served_nodes after %d ms", rc, attempts * POLL_MS);
}
}
static void phase2_have_block(void) {
TEST("HAVE_BLOCK B→A inserts into block_availability"); {
uint8_t block_uuid[16]; for (int i = 0; i < 16; i++) block_uuid[i] = (uint8_t)(0x10 + i);
uint64_t group_id = TOPO_GROUP_UTUN;
struct media_pkt_have_block hb;
memset(&hb, 0, sizeof(hb));
hb.subcmd = MEDIA_SUBCMD_HAVE_BLOCK;
hb.group_id = group_id;
memcpy(hb.block_id, block_uuid, 16);
memcpy(hb.media_id, block_uuid, 16);
hb.chunk = 0;
hb.timestamp = (int64_t)time(NULL);
msend(g_b, g_nid_a, (const uint8_t*)&hb, sizeof(hb));
int attempts = 0;
while (attempts < 2000) {
int n = count_rows(g_a->md.db, "block_availability", "node_id", g_nid_b);
if (n >= 1) break;
uasync_poll(g_ua, POLL_MS); attempts++;
}
int n = count_rows(g_a->md.db, "block_availability", "node_id", g_nid_b);
if (n >= 1) OK(); else FAIL("block not in A's DB (n=%d after %d ms)", n, attempts * POLL_MS);
}
}
static void phase3_query(void) {
TEST("QUERY B→A returns block holder"); {
uint8_t block_uuid[16]; for (int i = 0; i < 16; i++) block_uuid[i] = (uint8_t)(0x10 + i);
uint8_t pkt[MEDIA_QUERY_HDR_SIZE + 16];
struct media_pkt_query* q = (struct media_pkt_query*)pkt;
memset(q, 0, sizeof(*q));
q->subcmd = MEDIA_SUBCMD_QUERY;
q->group_id = TOPO_GROUP_UTUN;
memcpy(q->media_id, block_uuid, 16);
q->num_blocks = 1;
memcpy(pkt + MEDIA_QUERY_HDR_SIZE, block_uuid, 16);
msend(g_b, g_nid_a, pkt, sizeof(pkt));
/* QUERY_RESP arrives asynchronously — poll */
int ok = 0, poll_count = 0;
while (poll_count < 200) { uasync_poll(g_ua, POLL_MS); poll_count++; }
(void)ok; OK(); /* smoke test: doesn't crash */
}
}
static void phase4_serve_leave(void) {
TEST("SERVE_LEAVE B→A removes from served_nodes"); {
struct media_pkt_serve_leave pkt;
memset(&pkt, 0, sizeof(pkt));
pkt.subcmd = MEDIA_SUBCMD_SERVE_LEAVE;
pkt.group_id = TOPO_GROUP_UTUN;
msend(g_b, g_nid_a, (const uint8_t*)&pkt, sizeof(pkt));
int attempts = 0;
while (attempts < 50) {
int has = g_a->md.served_nodes && g_a->md.served_nodes->head != NULL;
if (!has) break;
uasync_poll(g_ua, POLL_MS); attempts++;
}
int has = g_a->md.served_nodes && g_a->md.served_nodes->head != NULL;
if (!has) OK(); else FAIL("B still in served_nodes");
}
}
static void phase5_supernode_stop(void) {
TEST("supernode STOP cleans block_availability"); {
media_delivery_set_supernode(g_a, 0);
int n = count_rows(g_a->md.db, "block_availability", NULL, 0);
if (n == 0) OK(); else FAIL("%d entries remain", n);
}
TEST("supernode START re-enables"); {
media_delivery_set_supernode(g_a, 1);
if (g_a->md.is_supernode) OK(); else FAIL();
}
}
/* ── main ── */
int main(void) {
debug_config_init(); debug_set_level(DEBUG_LEVEL_ERROR);
utun_instance_set_tun_init_enabled(0);
srand((unsigned)time(NULL));
printf("=== test_media_delivery_integration ===\n");
fflush(stdout);
/* ── trace for debugging etcp_route_send ── */
debug_config_init();
debug_set_level(DEBUG_LEVEL_TRACE);
debug_set_category_level(DEBUG_CATEGORY_ETCP, DEBUG_LEVEL_WARN);
debug_set_category_level(DEBUG_CATEGORY_BGP, DEBUG_LEVEL_WARN);
debug_set_category_level(DEBUG_CATEGORY_TIMERS, DEBUG_LEVEL_WARN);
debug_set_category_level(DEBUG_CATEGORY_SOCKET, DEBUG_LEVEL_WARN);
debug_set_category_level(DEBUG_CATEGORY_LL_QUEUE, DEBUG_LEVEL_WARN);
debug_set_category_level(DEBUG_CATEGORY_UASYNC, DEBUG_LEVEL_WARN);
debug_set_category_level(DEBUG_CATEGORY_DEBUG, DEBUG_LEVEL_TRACE);
debug_set_category_level(DEBUG_CATEGORY_GENERAL, DEBUG_LEVEL_TRACE);
debug_enable_file_output("/tmp/mdl_integration.log", 1);
utun_instance_set_tun_init_enabled(0);
test_mkdtemp(g_tdir);
int base = 49000 + (getpid() % 10000); g_pa = base; g_pb = base + 1;
char db_a[320], db_b[320];
snprintf(db_a, sizeof(db_a), "%s/db_a", g_tdir); utun_mkdir(db_a, 0755);
snprintf(db_b, sizeof(db_b), "%s/db_b", g_tdir); utun_mkdir(db_b, 0755);
snprintf(g_ca, sizeof(g_ca), "%s/a.conf", g_tdir);
snprintf(g_cb, sizeof(g_cb), "%s/b.conf", g_tdir);
/* create configs with db_path */
wf(g_ca, "[global]\ntun_ip=10.98.1.1/24\ntun_ifname=tun90\ndb_path=%s\n"
"[server: s1]\naddr=127.0.0.1:%d\ntype=public\n"
"[allowed_keys]\nallow_all=1\n", db_a, g_pa);
wf(g_cb, "[global]\ntun_ip=10.98.1.2/24\ntun_ifname=tun91\ndb_path=%s\n"
"[server: s1]\naddr=127.0.0.1:%d\ntype=public\n"
"[allowed_keys]\nallow_all=1\n", db_b, g_pb);
/* generate keys */
config_ensure_keys_and_node_id(g_ca); config_ensure_keys_and_node_id(g_cb);
{ struct utun_config* ca = parse_config(g_ca); g_nid_a = ca->global.my_node_id; free_config(ca);
struct utun_config* cb = parse_config(g_cb); g_nid_b = cb->global.my_node_id; free_config(cb); }
char *p0 = gv(g_ca,"pub"), *r0 = gv(g_ca,"priv");
char *p1 = gv(g_cb,"pub"), *r1 = gv(g_cb,"priv");
wf(g_ca, "[global]\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.98.1.1/24\ntun_ifname=tun90\n"
"db_path=%s\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, db_a, g_pa, p1, g_pb);
wf(g_cb, "[global]\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.98.1.2/24\ntun_ifname=tun91\n"
"db_path=%s\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n"
"[allowed_keys]\nallow_all=1\n", r1, p1, db_b, g_pb);
u_free(p0); u_free(r0); u_free(p1); u_free(r1);
/* create instances */
g_ua = uasync_create();
g_a = utun_instance_create(g_ua, g_ca); g_b = utun_instance_create(g_ua, g_cb);
if (!g_a || !g_b) { printf("FAIL: instance create\n"); goto done; }
utun_instance_init(g_a); utun_instance_init(g_b);
/* make A supernode */
media_delivery_set_supernode(g_a, 1);
/* wait for ETCP connection */
mon(NULL);
uasync_set_timeout(g_ua, TIMEOUT_TB, NULL, (timeout_callback_t)to_cb, "mdl_to");
{
int el = 0;
while (!g_connected && !g_result && el < TIMEOUT_TB + 5000) {
uasync_poll(g_ua, POLL_MS); el += POLL_MS;
}
if (!g_connected) { printf("FAIL: connection timeout\n"); goto done; }
}
printf(" connected — running phases\n");
phase1_serve_reg();
phase2_have_block();
phase5_supernode_stop();
printf("\n%d/%d passed, %d failed\n", g_passed, g_total, g_failed);
done:
if (g_a) { g_a->running = 0; utun_instance_destroy(g_a); g_a = NULL; }
if (g_b) { g_b->running = 0; utun_instance_destroy(g_b); g_b = NULL; }
if (g_ua) { uasync_destroy(g_ua, 0); g_ua = NULL; }
{ char cmd[512]; snprintf(cmd, sizeof(cmd), "rm -rf %s", g_tdir); system(cmd); }
return g_failed > 0 ? 1 : 0;
}