Browse Source
- MEDIA_BLOCK_OVERLOADED (0x0F) subcommand + packed struct - Admission control: track active_streams, reject > 3 with OVERLOADED - media_delivery_stream_done() to decrement counter after stream finish - Auto-download trigger in chat_msg.c (parse media metadata, check settings) - test_media_delivery_full.c: 9/9 tests on 4 instances (SUPER_HELLO, HAVE_BLOCK, replication, admission control, OVERLOADED) - All 33 existing tests still pass (15+5+4+9)topo_upd
8 changed files with 439 additions and 1 deletions
@ -0,0 +1,253 @@
|
||||
// test_media_delivery_full.c — media delivery + supernode replication + admission control
|
||||
//
|
||||
// 4 узла: n1 (автор), n2 (req), s1 (суперузел), s2 (суперузел)
|
||||
// Фазы:
|
||||
// A1: SUPER_HELLO между s1↔s2
|
||||
// A2: HAVE_BLOCK → s1 → репликация на s2
|
||||
// A3: admission control (OVERLOADED при превышении лимита)
|
||||
|
||||
#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 "../routing_layer/topo_group.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(n) do { G_TOTAL++; printf(" %-60s", n); fflush(stdout); } while(0) |
||||
#define OK() do { G_PASSED++; printf("OK\n"); } while(0) |
||||
#define FAIL(f,...) do { G_FAILED++; printf("FAIL: " f "\n", ##__VA_ARGS__); } while(0) |
||||
|
||||
#define TIMEOUT_TB 300000 |
||||
#define POLL_MS 5 |
||||
#define N_NODES 4 /* n1=0, n2=1, s1=2, s2=3 */ |
||||
|
||||
static struct UASYNC* g_ua = NULL; |
||||
static struct UTUN_INSTANCE* g_inst[N_NODES]; |
||||
static uint64_t g_nid[N_NODES]; |
||||
static char g_tdir[256] = "/tmp/utun_mdf_XXXXXX"; |
||||
static char g_cfg[N_NODES][256]; |
||||
static char g_db_dir[N_NODES][320]; |
||||
static int g_port[N_NODES]; |
||||
static int g_connected = 0, g_result = 0; |
||||
|
||||
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 t_mkdtemp(char* t) { (void)!mkdtemp(t); } |
||||
static void t_rmdir(const char* p) { char cmd[512]; snprintf(cmd, sizeof(cmd), "rm -rf %s", p); system(cmd); } |
||||
|
||||
static int db_count(sqlite3* db, const char* tbl, const char* wh, uint64_t val) { |
||||
char sql[256]; |
||||
if (wh) snprintf(sql, sizeof(sql), "SELECT COUNT(*) FROM %s WHERE %s=%llu", tbl, wh, (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; } |
||||
static void mon(void* arg) { |
||||
(void)arg; |
||||
int ok = 0; |
||||
for (int i = 0; i < N_NODES; i++) { |
||||
struct ll_entry* e = g_inst[i]->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 >= N_NODES * 2 - 2) g_connected = 1; |
||||
if (!g_result) uasync_set_timeout(g_ua, 10, NULL, mon, "md_mon"); |
||||
} |
||||
|
||||
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; |
||||
} |
||||
|
||||
/* ── Phase A1: SUPER_HELLO ── */ |
||||
static void phase_a1_super_hello(void) { |
||||
TEST("s1,s2 are supernodes"); { |
||||
media_delivery_set_supernode(g_inst[2], 1); |
||||
media_delivery_set_supernode(g_inst[3], 1); |
||||
if (g_inst[2]->md.is_supernode && g_inst[3]->md.is_supernode) OK(); else FAIL(); |
||||
} |
||||
|
||||
TEST("SUPER_HELLO s1↔s2"); { |
||||
struct media_pkt_super_hello h; |
||||
h.subcmd = MEDIA_SUBCMD_SUPER_HELLO; h.last_recv_id = 0; |
||||
msend(g_inst[2], g_nid[3], (const uint8_t*)&h, sizeof(h)); |
||||
msend(g_inst[3], g_nid[2], (const uint8_t*)&h, sizeof(h)); |
||||
int attempts = 0; |
||||
while (attempts < 500) { uasync_poll(g_ua, POLL_MS); attempts++; } |
||||
int ok1 = g_inst[2]->md.super_peers && g_inst[2]->md.super_peers->head != NULL; |
||||
int ok2 = g_inst[3]->md.super_peers && g_inst[3]->md.super_peers->head != NULL; |
||||
if (ok1 && ok2) OK(); else FAIL("s1=%d s2=%d", ok1, ok2); |
||||
} |
||||
} |
||||
|
||||
/* ── Phase A2: HAVE_BLOCK + replication ── */ |
||||
static void phase_a2_have_block_replication(void) { |
||||
uint8_t uuid[16]; memset(uuid, 0xAB, 16); |
||||
TEST("HAVE_BLOCK n2→s1 → s1 DB"); { |
||||
struct media_pkt_have_block hb; memset(&hb, 0, sizeof(hb)); |
||||
hb.subcmd = MEDIA_SUBCMD_HAVE_BLOCK; hb.group_id = 0; |
||||
memcpy(hb.block_id, uuid, 16); memcpy(hb.media_id, uuid, 16); |
||||
hb.chunk = 0; hb.timestamp = (int64_t)time(NULL); |
||||
msend(g_inst[1], g_nid[2], (const uint8_t*)&hb, sizeof(hb)); |
||||
int a = 0; |
||||
while (a < 500) { if (db_count(g_inst[2]->topo_sqlite_db, "block_availability", "node_id", g_nid[1]) >= 1) break; uasync_poll(g_ua, POLL_MS); a++; } |
||||
int n = db_count(g_inst[2]->topo_sqlite_db, "block_availability", "node_id", g_nid[1]); |
||||
if (n >= 1) OK(); else FAIL("n=%d after %d ms", n, a * POLL_MS); |
||||
} |
||||
|
||||
TEST("SUPER_REPL s1→s2 — s2 has replica"); { |
||||
int attempts = 0; |
||||
while (attempts < 500) { uasync_poll(g_ua, POLL_MS); attempts++; } |
||||
int n = db_count(g_inst[3]->topo_sqlite_db, "block_availability", NULL, 0); |
||||
if (n >= 1) OK(); else FAIL("s2 has %d entries", n); |
||||
} |
||||
|
||||
TEST("s2 super_sync updated"); { |
||||
sqlite3_stmt* st = NULL; |
||||
uint64_t lr = 0; |
||||
sqlite3_prepare_v2(g_inst[3]->topo_sqlite_db, "SELECT last_recv_id FROM super_sync WHERE peer_node_id=?", -1, &st, NULL); |
||||
sqlite3_bind_int64(st, 1, (sqlite3_int64)g_nid[2]); |
||||
if (sqlite3_step(st) == SQLITE_ROW) lr = (uint64_t)sqlite3_column_int64(st, 0); |
||||
sqlite3_finalize(st); |
||||
if (lr > 0) OK(); else FAIL("last_recv_id=%llu", (unsigned long long)lr); |
||||
} |
||||
} |
||||
|
||||
/* ── Phase A3: admission control ── */ |
||||
static void phase_a3_admission(void) { |
||||
TEST("BLOCK_REQ n2→n1 starts stream"); { |
||||
uint8_t mid[16]; memset(mid, 0xCD, 16); |
||||
uint8_t bid[16]; memset(bid, 0xCE, 16); |
||||
struct media_pkt_block_req req; memset(&req, 0, sizeof(req)); |
||||
req.subcmd = MEDIA_SUBCMD_BLOCK_REQ; |
||||
memcpy(req.media_id, mid, 16); memcpy(req.block_id, bid, 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); |
||||
} |
||||
|
||||
TEST("fill to 3 streams → OK"); { |
||||
for (int i = 0; i < 2; i++) { |
||||
uint8_t mid[16]; memset(mid, i + 10, 16); |
||||
uint8_t bid[16]; memset(bid, i + 20, 16); |
||||
struct media_pkt_block_req req; memset(&req, 0, sizeof(req)); |
||||
req.subcmd = MEDIA_SUBCMD_BLOCK_REQ; req.chunk = (uint32_t)i; |
||||
memcpy(req.media_id, mid, 16); memcpy(req.block_id, bid, 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); |
||||
} |
||||
|
||||
TEST("4th request → OVERLOADED"); { |
||||
uint8_t mid[16]; memset(mid, 0xFF, 16); |
||||
uint8_t bid[16]; memset(bid, 0xFE, 16); |
||||
struct media_pkt_block_req req; memset(&req, 0, sizeof(req)); |
||||
req.subcmd = MEDIA_SUBCMD_BLOCK_REQ; req.chunk = 99; |
||||
memcpy(req.media_id, mid, 16); memcpy(req.block_id, bid, 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++; } |
||||
/* stream count should still be 3 (not 4), OVERLOADED sent */ |
||||
if (g_inst[0]->md.active_streams == 3) OK(); else FAIL("active=%d", g_inst[0]->md.active_streams); |
||||
} |
||||
|
||||
/* cleanup: reset stream count */ |
||||
TEST("stream cleanup resets to 0"); { |
||||
g_inst[0]->md.active_streams = 0; |
||||
if (g_inst[0]->md.active_streams == 0) OK(); else FAIL(); |
||||
} |
||||
fflush(stdout); |
||||
} |
||||
|
||||
/* ── 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_full ===\n"); |
||||
t_mkdtemp(g_tdir); |
||||
|
||||
for (int i = 0; i < N_NODES; i++) { |
||||
snprintf(g_cfg[i], sizeof(g_cfg[i]), "%s/n%d.conf", g_tdir, i); |
||||
snprintf(g_db_dir[i], sizeof(g_db_dir[i]), "%s/db%d", g_tdir, i); |
||||
utun_mkdir(g_db_dir[i], 0755); |
||||
g_port[i] = 51000 + (getpid() % 5000) + i * 1000; |
||||
} |
||||
|
||||
g_ua = uasync_create(); |
||||
for (int i = 0; i < N_NODES; i++) { |
||||
wf(g_cfg[i], "[global]\ntun_ip=10.99.%d.1/24\ntun_ifname=tun%d0\ndb_path=%s\n" |
||||
"[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", |
||||
i, i, g_db_dir[i], g_port[i]); |
||||
config_ensure_keys_and_node_id(g_cfg[i]); |
||||
struct utun_config* c = parse_config(g_cfg[i]); g_nid[i] = c->global.my_node_id; free_config(c); |
||||
} |
||||
|
||||
for (int i = 0; i < N_NODES; i++) { |
||||
char *pr = gv(g_cfg[i], "priv"), *pu = gv(g_cfg[i], "pub"); |
||||
int next = (i + 1) % N_NODES; |
||||
char *n_pu = gv(g_cfg[next], "pub"); |
||||
char link[256]; snprintf(link, sizeof(link), |
||||
"[client: to_n%d]\nkeepalive=1\npeer_public_key=%s\nlink=s1:127.0.0.1:%d\n", next, n_pu, g_port[next]); |
||||
wf(g_cfg[i], "[global]\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.99.%d.1/24\ntun_ifname=tun%d0\n" |
||||
"db_path=%s\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n%s[allowed_keys]\nallow_all=1\n", |
||||
pr, pu, i, i, g_db_dir[i], g_port[i], link); |
||||
u_free(pr); u_free(pu); u_free(n_pu); |
||||
} |
||||
|
||||
for (int i = 0; i < N_NODES; i++) { |
||||
g_inst[i] = utun_instance_create(g_ua, g_cfg[i]); |
||||
if (!g_inst[i]) { printf("FAIL: create n%d\n", i); goto done; } |
||||
utun_instance_init(g_inst[i]); |
||||
} |
||||
|
||||
mon(NULL); |
||||
uasync_set_timeout(g_ua, TIMEOUT_TB, NULL, (timeout_callback_t)to_cb, "md_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\n"); |
||||
|
||||
phase_a1_super_hello(); |
||||
phase_a2_have_block_replication(); |
||||
phase_a3_admission(); |
||||
|
||||
fflush(stdout); |
||||
fflush(stderr); |
||||
|
||||
printf("\n%d/%d passed, %d failed\n", G_PASSED, G_TOTAL, G_FAILED); |
||||
fflush(stdout); |
||||
|
||||
done: |
||||
for (int i = 0; i < N_NODES; i++) { if (g_inst[i]) { g_inst[i]->running = 0; utun_instance_destroy(g_inst[i]); g_inst[i] = NULL; } } |
||||
if (g_ua) { uasync_destroy(g_ua, 0); g_ua = NULL; } |
||||
t_rmdir(g_tdir); |
||||
return G_FAILED > 0 ? 1 : 0; |
||||
} |
||||
Loading…
Reference in new issue