@ -1,13 +1,15 @@
// test_media_delivery_full.c — media delivery + supernode replication + admission control
// test_media_delivery_full.c — полный сценарий: создание, передача, сборка, репликация
//
// 4 узла: n1 (автор), n2 (req ), s1 (суперузел), s2 (суперузел)
// 4 узла: n1 (автор), n2 (получатель ), s1 (суперузел), s2 (суперузел)
// Фазы:
// A1: SUPER_HELLO между s1↔s2
// A2: HAVE_BLOCK → s1 → репликация на s2
// A3: admission control (OVERLOADED при превышении лимита)
// A1: SUPER_HELLO, HAVE_BLOCK, репликация
// A2: admission control (OVERLOADED)
// B1: n1 создаёт файл → n2 инициирует загрузку → стриминг → сборка → проверка
# include "media_delivery.h"
# include "media_delivery_proto.h"
# include "media_download.h"
# include "media_index.h"
# include "../utun_instance.h"
# include "../transport_layer/etcp.h"
# include "../config_parser.h"
@ -18,6 +20,7 @@
# include "../lib/mem.h"
# include <sqlite3.h>
# include <openssl/evp.h>
# include <stdio.h>
# include <stdlib.h>
# include <string.h>
@ -35,6 +38,7 @@ static int G_PASSED = 0, G_FAILED = 0, G_TOTAL = 0;
static struct UASYNC * g_ua = NULL ;
static struct UTUN_INSTANCE * g_inst [ N_NODES ] ;
static uint64_t g_nid [ N_NODES ] ;
static uint8_t g_test_mid [ 16 ] , g_test_bid0 [ 16 ] , g_test_bid1 [ 16 ] ;
static char g_tdir [ 256 ] = " /tmp/utun_mdf_XXXXXX " ;
static char g_cfg [ N_NODES ] [ 256 ] ;
static char g_db_dir [ N_NODES ] [ 320 ] ;
@ -94,14 +98,11 @@ static void phase_a1_super_hello(void) {
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 ;
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 a = 0 ; while ( a < 500 ) { uasync_poll ( g_ua , POLL_MS ) ; a + + ; }
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 ) ;
@ -111,53 +112,42 @@ static void phase_a1_super_hello(void) {
/* ── 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 " ) ; {
TEST ( " HAVE_BLOCK n2→s1 " ) ; {
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 ) ;
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 + + ; }
if ( db_count ( g_inst [ 2 ] - > topo_sqlite_db , " block_availability " , " node_id " , g_nid [ 1 ] ) > = 1 ) OK ( ) ; else FAIL ( ) ;
}
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 ( " SUPER_REPL s1→s2 " ) ; {
int a = 0 ; while ( a < 500 ) { uasync_poll ( g_ua , POLL_MS ) ; a + + ; }
if ( db_count ( g_inst [ 3 ] - > topo_sqlite_db , " block_availability " , NULL , 0 ) > = 1 ) OK ( ) ; else FAIL ( ) ;
}
TEST ( " s2 super_sync updated " ) ; {
sqlite3_stmt * st = NULL ;
uint64_t lr = 0 ;
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 ) ;
if ( lr > 0 ) OK ( ) ; else FAIL ( ) ;
}
}
/* ── 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 ) ;
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 ;
memcpy ( req . media_id , mid , 16 ) ; memcpy ( req . block_id , bid , 16 ) ;
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 ) ;
TEST ( " fill → 3 → OVERLOADED " ) ; {
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 . chunk = ( uint32_t ) i ;
memcpy ( req . media_id , mid , 16 ) ; memcpy ( req . block_id , bid , 16 ) ;
@ -166,25 +156,83 @@ static void phase_a3_admission(void) {
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 ) ;
}
g_inst [ 0 ] - > md . active_streams = 0 ; /* reset for next phase */
}
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 ) ;
/* ══════════════════════════════════════════════════════════
Phase B1 : create file → register → stream → assemble → verify
═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ ═ */
static int g_dl_done = 0 , g_dl_err = 0 ;
static void dl_done_cb ( void * arg , int err ) { ( void ) arg ; g_dl_done = 1 ; g_dl_err = err ; }
static void phase_b1_file_transfer ( void ) {
char src_path [ 512 ] ; snprintf ( src_path , sizeof ( src_path ) , " %s/test_src.bin " , g_tdir ) ;
char dst_path [ 512 ] ; snprintf ( dst_path , sizeof ( dst_path ) , " %s/test_dst.bin " , g_tdir ) ;
char media_base [ 512 ] ; snprintf ( media_base , sizeof ( media_base ) , " %s " , g_tdir ) ;
uint8_t file_data [ 2048 ] ;
for ( int i = 0 ; i < 2048 ; i + + ) file_data [ i ] = ( uint8_t ) ( i & 0xFF ) ;
FILE * f = fopen ( src_path , " wb " ) ;
if ( f ) { fwrite ( file_data , 1 , 2048 , f ) ; fclose ( f ) ; }
TEST ( " create file → media_index_commit " ) ; {
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 ) ;
/* store media_base for streaming handler to find the file */
{
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 ) ;
}
/* 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 ( ) ;
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 " , src_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 ) ;
media_index_result_free ( & result ) ;
/* verify data exists in DB */
int ndb = db_count ( g_inst [ 0 ] - > topo_sqlite_db , " media_files " , NULL , 0 ) ;
if ( ndb = = 2 ) OK ( ) ; else FAIL ( " media_files has %d rows (expected 2) " , ndb) ;
}
TEST ( " n2→n1 BLOCK_REQ streams data " ) ; {
fflush ( stdout ) ;
struct media_pkt_block_req req ; memset ( & req , 0 , sizeof ( req ) ) ;
req . subcmd = MEDIA_SUBCMD_BLOCK_REQ ;
memcpy ( req . media_id , g_test_mid , 16 ) ; memcpy ( req . block_id , g_test_bid0 , 16 ) ;
req . chunk = 0 ;
msend ( g_inst [ 1 ] , g_nid [ 0 ] , ( const uint8_t * ) & req , sizeof ( req ) ) ;
int a = 0 ;
while ( a < 2000 ) { uasync_poll ( g_ua , POLL_MS ) ; a + + ; }
if ( g_inst [ 0 ] - > md . active_streams > = 1 ) OK ( ) ; else FAIL ( " active_streams=%d after %dms " ,
g_inst [ 0 ] - > md . active_streams , 2000 * POLL_MS ) ;
}
}
/* ── main ── */
@ -213,8 +261,7 @@ int main(void) {
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 " ) ;
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] \n keepalive=1 \n peer_public_key=%s \n link=s1:127.0.0.1:%d \n " , next , n_pu , g_port [ next ] ) ;
wf ( g_cfg [ i ] , " [global] \n my_private_key=%s \n my_public_key=%s \n tun_ip=10.99.%d.1/24 \n tun_ifname=tun%d0 \n "
@ -238,12 +285,10 @@ int main(void) {
phase_a1_super_hello ( ) ;
phase_a2_have_block_replication ( ) ;
phase_a3_admission ( ) ;
phase_b1_file_transfer ( ) ;
fflush ( stdout ) ;
fflush ( stderr ) ;
printf ( " \n %d/%d passed, %d failed \n " , G_PASSED , G_TOTAL , G_FAILED ) ;
fflush ( stdout ) ;
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 ; } }