@ -974,6 +974,336 @@ static int run_stage5(void) {
return stage_failures ;
}
/* ── Stage 7: stress test — 10 spammers + 2 A↔B observers, star topology ── */
# define STRESS_N 13 /* 1 hub + 10 spammers + 2 observers */
# define STRESS_HUB 0
# define STRESS_SPAM_BEGIN 1
# define STRESS_SPAM_END 10 /* indices 1..10 */
# define STRESS_OBS_A 11
# define STRESS_OBS_B 12
# define STRESS_SPAM_MS 5000
# define STRESS_SYNC_TB 200000 /* 20s timeout for final sync */
# define STRESS_LINK_TB 300000 /* 30s for all links up */
struct str_state {
struct UASYNC * ua ;
struct UTUN_INSTANCE * inst [ STRESS_N ] ;
struct ms_data data [ STRESS_N ] ;
struct ms_test_ctx ctx [ STRESS_N ] ;
void * global_timeout ;
int phase ; /* 0=running, 1=stopping, 2=timeout */
int spam_active ;
int sync_pending ; /* atomic: count of outstanding sync ops */
} ;
static struct str_state * gs = NULL ;
static void _str_cleanup ( void ) ;
static int _str_cond_sync_done ( void ) ;
static int _str_cond_links_up ( void ) ;
static void _str_wait_sync_quiesce ( void ) ;
static void _str_sync_done ( uint64_t peer , const char * ns , int result , void * arg ) {
( void ) peer ; ( void ) ns ; ( void ) result ;
struct str_state * s = ( struct str_state * ) arg ;
if ( s ) __sync_fetch_and_sub ( & s - > sync_pending , 1 ) ;
}
static void _str_spam_cb ( void * arg ) ;
static void _str_spam_schedule ( struct str_state * s , int idx ) {
if ( ! s - > spam_active ) return ;
int delay_tb = ( ( rand ( ) % 30 ) + 1 ) * 10 ; /* 1-30ms → 10-300 timebase units */
uasync_set_timeout ( s - > ua , ( uint32_t ) delay_tb , ( void * ) ( intptr_t ) idx , _str_spam_cb , " spam " ) ;
}
static void _str_spam_cb ( void * arg ) {
if ( ! gs | | ! gs - > spam_active ) return ;
int idx = ( int ) ( intptr_t ) arg ;
if ( idx < STRESS_SPAM_BEGIN | | idx > STRESS_SPAM_END ) { _str_spam_schedule ( gs , idx ) ; return ; }
/* spammer: bump own member, sync to hub */
uint64_t key = ( ( uint64_t ) idx ) < < 60 ;
_data_insert ( & gs - > data [ idx ] , key , ( uint32_t ) ( gs - > data [ idx ] . count + 1 ) ) ;
merkle_sync_recompute_path ( gs - > inst [ idx ] , " stress " , key ) ;
__sync_fetch_and_add ( & gs - > sync_pending , 1 ) ;
merkle_sync_start ( gs - > inst [ idx ] , gs - > inst [ STRESS_HUB ] - > node_id , " stress " , _str_sync_done , gs ) ;
_str_spam_schedule ( gs , idx ) ;
}
static void _str_obs_spam_cb ( void * arg ) {
if ( ! gs | | ! gs - > spam_active ) return ;
int idx = ( int ) ( intptr_t ) arg ;
__sync_fetch_and_add ( & gs - > sync_pending , 1 ) ;
merkle_sync_start ( gs - > inst [ idx ] , gs - > inst [ STRESS_HUB ] - > node_id , " stress " , _str_sync_done , gs ) ;
int delay_tb = ( ( rand ( ) % 10 ) + 1 ) * 10 ; /* 1-10ms */
uasync_set_timeout ( gs - > ua , ( uint32_t ) delay_tb , arg , _str_obs_spam_cb , " obs_spam " ) ;
}
static void _str_global_timeout ( void * arg ) {
( void ) arg ;
if ( gs ) gs - > phase = 2 ;
}
static int _str_cond_links_up ( void ) {
if ( ! gs ) return 0 ;
int total_links = 0 ;
for ( int i = 0 ; i < STRESS_N ; i + + ) {
if ( ! gs - > inst [ i ] | | ! gs - > inst [ i ] - > connections ) return 0 ;
struct ll_entry * e = gs - > inst [ i ] - > connections - > head ;
while ( e ) {
struct conn_queue_entry * ce = ( struct conn_queue_entry * ) e - > data ;
struct ETCP_LINK * l = ce - > conn - > links ;
while ( l ) { if ( l - > initialized ) total_links + + ; l = l - > next ; }
e = e - > next ;
}
}
/* hub should have 12 links (one per spoke), spokes have 1 each */
return total_links > = STRESS_N * 2 - 2 ; /* 12 hub links + 12 spoke links = 24 total */
}
static int _str_wait_for ( const char * desc , int ( * cond ) ( void ) , int timeout_tb ) {
struct str_state * s = gs ;
uint64_t start = get_time_tb ( ) ;
while ( ! cond ( ) & & ( get_time_tb ( ) - start ) < ( uint64_t ) timeout_tb & & s - > phase = = 0 )
uasync_poll ( s - > ua , 1 ) ;
if ( cond ( ) ) return 1 ;
printf ( " STRESS TIMEOUT: %s \n " , desc ) ; s - > phase = 2 ;
return 0 ;
}
static int _str_init ( void ) {
struct str_state * s = u_calloc ( 1 , sizeof ( * s ) ) ;
if ( ! s ) return - 1 ;
gs = s ;
s - > spam_active = 0 ;
s - > phase = 0 ;
s - > sync_pending = 0 ;
/* --- temp dir --- */
char tmpld [ 128 ] = " /tmp/utun_str_XXXXXX " ;
if ( test_mkdtemp ( tmpld ) ! = 0 ) { printf ( " mkdtemp failed \n " ) ; u_free ( s ) ; gs = NULL ; return - 1 ; }
/* --- create configs --- */
int base_port = 44000 + ( getpid ( ) % 10000 ) ;
char cfgs [ STRESS_N ] [ 256 ] ; int ports [ STRESS_N ] ;
for ( int i = 0 ; i < STRESS_N ; i + + ) ports [ i ] = base_port + i ;
utun_instance_set_tun_init_enabled ( 0 ) ;
s - > ua = uasync_create ( ) ;
if ( ! s - > ua ) goto fail ;
/* hub config */
snprintf ( cfgs [ 0 ] , sizeof ( cfgs [ 0 ] ) , " %s/cfg_0.conf " , tmpld ) ;
{
char dbp [ 256 ] ; snprintf ( dbp , sizeof ( dbp ) , " %s/db_0 " , tmpld ) ; utun_mkdir ( dbp , 0755 ) ;
_write_cfg ( cfgs [ 0 ] ,
" [global] \n my_node_id=0xAAAAAAAAAAAAAAAA \n tun_ip=10.240.0.1/24 \n tun_ifname=tun240 \n keepalive_adaptive=0 \n db_path=%s \n \n "
" [server:s0] \n addr=127.0.0.1:%d \n type=public \n \n "
" [allowed_keys] \n allow_all=1 \n " , dbp , ports [ 0 ] ) ;
config_ensure_keys_and_node_id ( cfgs [ 0 ] ) ;
}
char * hub_pub = _get_pubkey ( cfgs [ 0 ] ) ; if ( ! hub_pub ) goto fail ;
/* spoke configs */
for ( int i = 1 ; i < STRESS_N ; i + + ) {
snprintf ( cfgs [ i ] , sizeof ( cfgs [ i ] ) , " %s/cfg_%d.conf " , tmpld , i ) ;
char dbp [ 256 ] ; snprintf ( dbp , sizeof ( dbp ) , " %s/db_%d " , tmpld , i ) ; utun_mkdir ( dbp , 0755 ) ;
char nid [ 32 ] ; snprintf ( nid , sizeof ( nid ) , " BBBBBBBBBBBB%02d00 " , i ) ;
_write_cfg ( cfgs [ i ] ,
" [global] \n my_node_id=0x%s \n tun_ip=10.240.0.%d/24 \n tun_ifname=tun240 \n keepalive_adaptive=0 \n db_path=%s \n \n "
" [server:s%d] \n addr=127.0.0.1:%d \n type=public \n \n "
" [client:to_hub] \n keepalive=1 \n peer_public_key=%s \n link=s%d:127.0.0.1:%d \n \n "
" [allowed_keys] \n allow_all=1 \n " ,
nid , i + 1 , dbp , i , ports [ i ] , hub_pub , i , ports [ 0 ] ) ;
config_ensure_keys_and_node_id ( cfgs [ i ] ) ;
}
u_free ( hub_pub ) ;
/* --- create instances --- */
for ( int i = 0 ; i < STRESS_N ; i + + ) {
s - > inst [ i ] = utun_instance_create ( s - > ua , cfgs [ i ] ) ;
if ( ! s - > inst [ i ] | | utun_instance_init ( s - > inst [ i ] ) ! = 0 ) {
printf ( " FAIL: instance %d init \n " , i ) ;
goto fail ;
}
}
printf ( " %d instances created, waiting for links... \n " , STRESS_N ) ;
s - > global_timeout = uasync_set_timeout ( s - > ua , STRESS_LINK_TB , NULL , _str_global_timeout , " str_timeout " ) ;
if ( ! _str_wait_for ( " links up " , _str_cond_links_up , STRESS_LINK_TB ) ) { printf ( " links never came up \n " ) ; goto fail ; }
printf ( " all links UP \n " ) ;
uasync_cancel_timeout ( s - > ua , s - > global_timeout ) ; s - > global_timeout = NULL ;
/* --- init merkle_sync on all --- */
for ( int i = 0 ; i < STRESS_N ; i + + ) {
memset ( & s - > data [ i ] , 0 , sizeof ( s - > data [ i ] ) ) ;
s - > ctx [ i ] . data = & s - > data [ i ] ;
s - > ctx [ i ] . inst = s - > inst [ i ] ;
s - > ctx [ i ] . tag = ( char ) ( ' A ' + i ) ;
if ( merkle_sync_init ( s - > inst [ i ] , MS_SVC_ID , & g_test_ops , & s - > ctx [ i ] ) ! = 0 ) {
printf ( " FAIL: merkle_sync_init %d \n " , i ) ; goto fail ;
}
}
printf ( " merkle_sync inited on all \n " ) ;
/* --- seed each instance with its own member; sync to hub --- */
for ( int i = 1 ; i < STRESS_N ; i + + ) {
uint64_t key = ( ( uint64_t ) i ) < < 60 ;
_data_insert ( & s - > data [ i ] , key , ( uint32_t ) i ) ;
merkle_sync_recompute_path ( s - > inst [ i ] , " stress " , key ) ;
}
/* initial sync: each spoke → hub */
s - > sync_pending = STRESS_N - 1 ;
for ( int i = 1 ; i < STRESS_N ; i + + )
merkle_sync_start ( s - > inst [ i ] , s - > inst [ STRESS_HUB ] - > node_id , " stress " , _str_sync_done , s ) ;
_str_wait_for ( " initial sync " , _str_cond_sync_done , STRESS_SYNC_TB ) ;
/* observers: sync with hub too */
return 0 ;
fail :
_str_cleanup ( ) ;
return - 1 ;
}
static int _str_cond_sync_done ( void ) {
return gs ? __sync_fetch_and_add ( & gs - > sync_pending , 0 ) < = 0 : 0 ;
}
static int _str_cond_sync_done_strict ( void ) {
return gs ? __sync_fetch_and_add ( & gs - > sync_pending , 0 ) < = 0 : 0 ;
}
static void _str_wait_sync_quiesce ( void ) {
struct str_state * s = gs ;
uint64_t start = get_time_tb ( ) ;
while ( s - > sync_pending > 0 & & ( get_time_tb ( ) - start ) < ( uint64_t ) STRESS_SYNC_TB )
uasync_poll ( s - > ua , 1 ) ;
}
static void _str_cleanup ( void ) {
struct str_state * s = gs ;
if ( ! s ) return ;
s - > spam_active = 0 ;
if ( s - > global_timeout & & s - > ua ) { uasync_cancel_timeout ( s - > ua , s - > global_timeout ) ; s - > global_timeout = NULL ; }
for ( int i = STRESS_N - 1 ; i > = 0 ; i - - ) {
if ( ! s - > inst [ i ] ) continue ;
merkle_sync_destroy ( s - > inst [ i ] ) ;
s - > inst [ i ] - > running = 0 ;
utun_instance_destroy ( s - > inst [ i ] ) ;
s - > inst [ i ] = NULL ;
}
if ( s - > ua ) { uasync_destroy ( s - > ua , 0 ) ; s - > ua = NULL ; }
u_free ( s ) ;
gs = NULL ;
}
static void test_stress_spam ( void ) {
TEST ( " stress: 10 spammers + 2 obs, 5s spam, verify merkle roots " ) ;
if ( _str_init ( ) ! = 0 ) { FAIL ( " init failed " ) ; _str_cleanup ( ) ; return ; }
struct str_state * s = gs ;
/* --- start spam --- */
s - > spam_active = 1 ;
srand ( 12345 ) ;
for ( int i = STRESS_SPAM_BEGIN ; i < = STRESS_SPAM_END ; i + + ) {
/* stagger initial delays */
int d = ( rand ( ) % 30 ) + 1 ;
uasync_set_timeout ( s - > ua , ( uint32_t ) ( d * 10 ) , ( void * ) ( intptr_t ) i , _str_spam_cb , " spam " ) ;
}
/* observers sync with hub aggressively */
uasync_set_timeout ( s - > ua , 10 , ( void * ) ( intptr_t ) STRESS_OBS_A , _str_obs_spam_cb , " obsA " ) ;
uasync_set_timeout ( s - > ua , 15 , ( void * ) ( intptr_t ) STRESS_OBS_B , _str_obs_spam_cb , " obsB " ) ;
/* --- run for 5 seconds --- */
printf ( " spamming for %dms... \n " , STRESS_SPAM_MS ) ;
uint64_t start = get_time_tb ( ) ;
while ( ( get_time_tb ( ) - start ) < ( uint64_t ) ( STRESS_SPAM_MS * 10 ) & & s - > phase = = 0 )
uasync_poll ( s - > ua , 1 ) ;
/* --- stop all spam --- */
s - > spam_active = 0 ;
/* let in-flight syncs settle: poll for up to 10s */
printf ( " sync_pending=%d, waiting for convergence... \n " , s - > sync_pending ) ;
{
uint64_t settle_start = get_time_tb ( ) ;
int last_pending = s - > sync_pending ;
while ( ( get_time_tb ( ) - settle_start ) < ( uint64_t ) ( STRESS_SYNC_TB ) ) {
uasync_poll ( s - > ua , 1 ) ;
int cur = __sync_fetch_and_add ( & s - > sync_pending , 0 ) ;
if ( cur = = 0 ) break ;
if ( cur ! = last_pending ) { last_pending = cur ; settle_start = get_time_tb ( ) ; }
}
}
printf ( " converged: sync_pending=%d phase=%d \n " , s - > sync_pending , s - > phase ) ;
if ( s - > phase = = 2 ) { FAIL ( " global timeout " ) ; _str_cleanup ( ) ; return ; }
/* --- final explicit sync from each spoke → hub, then wait --- */
for ( int i = 1 ; i < STRESS_N ; i + + )
merkle_sync_start ( s - > inst [ i ] , s - > inst [ STRESS_HUB ] - > node_id , " stress " , NULL , NULL ) ;
/* poll for data propagation */
for ( int i = 0 ; i < 500 ; i + + ) uasync_poll ( s - > ua , 5 ) ;
/* --- verify: tree sizes match pairwise --- */
int ok = 1 ;
int na = _db_count_rows ( s - > inst [ 0 ] - > topo_sqlite_db , " stress " ) ;
for ( int i = 1 ; i < STRESS_N & & ok ; i + + ) {
int nb = _db_count_rows ( s - > inst [ i ] - > topo_sqlite_db , " stress " ) ;
if ( na ! = nb ) { printf ( " tree size mismatch: hub=%d inst[%d]=%d \n " , na , i , nb ) ; ok = 0 ; }
}
/* --- verify: all root (level=1, prefix=0) hashes identical --- */
const uint8_t * root0 = merkle_sync_get_hash ( s - > inst [ 0 ] , " stress " , 1 , 0 ) ;
uint8_t root_copy [ MT_HASH_SIZE ] ; memcpy ( root_copy , root0 , MT_HASH_SIZE ) ;
for ( int i = 1 ; i < STRESS_N & & ok ; i + + ) {
const uint8_t * ri = merkle_sync_get_hash ( s - > inst [ i ] , " stress " , 1 , merkle_sync_level_prefix ( 0 , 1 ) ) ;
if ( memcmp ( root_copy , ri , MT_HASH_SIZE ) ! = 0 ) { printf ( " ROOT HASH MISMATCH inst[%d] \n " , i ) ; ok = 0 ; }
}
if ( ! ok ) { FAIL ( " hash mismatch " ) ; _str_cleanup ( ) ; return ; }
/* --- verify: recompute consistency for each instance --- */
for ( int i = 0 ; i < STRESS_N & & ok ; i + + ) {
sqlite3_stmt * st = NULL ;
sqlite3_prepare_v2 ( s - > inst [ i ] - > topo_sqlite_db ,
" SELECT level,prefix64 FROM merkle_tree_hash WHERE namespace=? ORDER BY level,prefix64 " ,
- 1 , & st , NULL ) ;
sqlite3_bind_text ( st , 1 , " stress " , - 1 , SQLITE_STATIC ) ;
while ( sqlite3_step ( st ) = = SQLITE_ROW & & ok ) {
uint8_t lv = ( uint8_t ) sqlite3_column_int ( st , 0 ) ;
uint64_t pf = ( uint64_t ) sqlite3_column_int64 ( st , 1 ) ;
EVP_MD_CTX * c = EVP_MD_CTX_new ( ) ;
EVP_DigestInit_ex ( c , EVP_sha256 ( ) , NULL ) ;
_bucket_hash ( & s - > ctx [ i ] , " stress " , lv , pf , c ) ;
uint8_t expected [ MT_HASH_SIZE ] ;
EVP_DigestFinal_ex ( c , expected , NULL ) ; EVP_MD_CTX_free ( c ) ;
const uint8_t * stored = merkle_sync_get_hash ( s - > inst [ i ] , " stress " , lv , pf ) ;
uint8_t scopy [ MT_HASH_SIZE ] ; memcpy ( scopy , stored , MT_HASH_SIZE ) ;
if ( memcmp ( expected , scopy , MT_HASH_SIZE ) ! = 0 ) {
printf ( " consistency fail inst[%d] L%d P%016llx \n " , i , lv , ( unsigned long long ) pf ) ;
ok = 0 ;
}
}
sqlite3_finalize ( st ) ;
}
if ( ! ok ) { FAIL ( " tree consistency " ) ; _str_cleanup ( ) ; return ; }
/* print summary */
printf ( " tree rows: %d data items per node: " , na ) ;
for ( int i = 0 ; i < STRESS_N & & i < 6 ; i + + ) printf ( " %d " , s - > data [ i ] . count ) ;
printf ( " .. \n " ) ;
PASS ( ) ;
_str_cleanup ( ) ;
}
static int run_stage6 ( void ) {
printf ( " --- Stage 6: protocol edge cases --- \n " ) ;
stage_failures = 0 ;
@ -982,6 +1312,13 @@ static int run_stage6(void) {
return stage_failures ;
}
static int run_stage7 ( void ) {
printf ( " --- Stage 7: stress test (10 spammers + 2 A↔B obs, 5s) --- \n " ) ;
stage_failures = 0 ;
test_stress_spam ( ) ;
return stage_failures ;
}
int main ( void ) {
printf ( " === test_merkle_sync === \n " ) ;
debug_config_init ( ) ;
@ -1005,6 +1342,9 @@ int main(void) {
if ( run_stage6 ( ) ! = 0 ) { printf ( " === Stage 6 FAILED === \n " ) ; return 1 ; }
printf ( " === Stage 6 PASSED (%d tests) === \n \n " , tests_run ) ;
if ( run_stage7 ( ) ! = 0 ) { printf ( " === Stage 7 FAILED === \n " ) ; return 1 ; }
printf ( " === Stage 7 PASSED (%d tests) === \n \n " , tests_run ) ;
printf ( " === ALL PASS (%d tests) === \n " , tests_run ) ;
return 0 ;
}