@ -856,7 +856,7 @@ static void test_randomized_two(void) {
if ( ! ok ) { FAIL ( " tree consistency iter %d " , iter ) ; _intg_cleanup ( ) ; break ; }
_intg_cleanup ( ) ;
}
if ( iter = = 50 ) PASS ( ) ;
if ( iter = = 2 5) PASS ( ) ;
}
/* ── Stage 5: randomized 3-instance star-topology sync ── */
@ -985,16 +985,26 @@ static int run_stage5(void) {
# 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 */
# define STRESS_MAX_ROUNDS 5
# define STRESS_DRAIN_TB 2000 /* 200ms drain after spam stop */
# define STRESS_ROUND_GAP_TB 500 /* 50ms gap between final rounds */
# define STRESS_SAFETY_TB 400000 /* 40s safety for spam+converge */
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 */
void * global_timeout ; /* setup: links-up timeout */
void * safety_timer ; /* spam+converge safety */
void * stop_timer ; /* spam stop */
void * drain_timer ; /* drain after spam stop */
void * round_timer ; /* gap between final rounds */
int phase ; /* 0=running, 1=pass, 2=fail */
int spam_active ;
int sync_pending ; /* atomic: count of outstanding sync ops */
int sync_pending ; /* initial sync only (balanced) */
int final_pending ; /* outstanding final-round syncs */
int final_round ; /* final round number */
} ;
static struct str_state * gs = NULL ;
@ -1002,7 +1012,14 @@ 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_stop_spam ( void * arg ) ;
static void _str_drain_cb ( void * arg ) ;
static void _str_round_cb ( void * arg ) ;
static void _str_start_final_round ( void ) ;
static void _str_final_done ( uint64_t peer , const char * ns , int result , void * arg ) ;
static void _str_check_roots ( void ) ;
static void _str_verify ( void ) ;
static void _str_fail_timeout ( void * arg ) ;
static void _str_sync_done ( uint64_t peer , const char * ns , int result , void * arg ) {
( void ) peer ; ( void ) ns ; ( void ) result ;
@ -1027,8 +1044,7 @@ static void _str_spam_cb(void* arg) {
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 ) ;
merkle_sync_start ( gs - > inst [ idx ] , gs - > inst [ STRESS_HUB ] - > node_id , " stress " , NULL , NULL ) ;
_str_spam_schedule ( gs , idx ) ;
}
@ -1037,8 +1053,7 @@ 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 ) ;
merkle_sync_start ( gs - > inst [ idx ] , gs - > inst [ STRESS_HUB ] - > node_id , " stress " , NULL , NULL ) ;
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 " ) ;
@ -1176,22 +1191,15 @@ 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 ; }
if ( s - > safety_timer & & s - > ua ) { uasync_cancel_timeout ( s - > ua , s - > safety_timer ) ; s - > safety_timer = NULL ; }
if ( s - > stop_timer & & s - > ua ) { uasync_cancel_timeout ( s - > ua , s - > stop_timer ) ; s - > stop_timer = NULL ; }
if ( s - > drain_timer & & s - > ua ) { uasync_cancel_timeout ( s - > ua , s - > drain_timer ) ; s - > drain_timer = NULL ; }
if ( s - > round_timer & & s - > ua ) { uasync_cancel_timeout ( s - > ua , s - > round_timer ) ; s - > round_timer = NULL ; }
for ( int i = STRESS_N - 1 ; i > = 0 ; i - - ) {
if ( ! s - > inst [ i ] ) continue ;
merkle_sync_destroy ( s - > inst [ i ] ) ;
@ -1204,74 +1212,60 @@ static void _str_cleanup(void) {
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 ) ;
static void _str_stop_spam ( void * arg ) {
( void ) arg ;
if ( ! gs ) return ;
gs - > stop_timer = NULL ;
gs - > spam_active = 0 ;
printf ( " spam stopped \n " ) ;
/* короткий дренаж — даём in-flight спам-синкам утихнуть перед финальными раундами */
gs - > drain_timer = uasync_set_timeout ( gs - > ua , STRESS_DRAIN_TB , NULL , _str_drain_cb , " str_drain " ) ;
}
if ( s - > phase = = 2 ) { FAIL ( " global timeout " ) ; _str_cleanup ( ) ; return ; }
static void _str_drain_cb ( void * arg ) {
( void ) arg ;
if ( ! gs ) return ;
gs - > drain_timer = NULL ;
_str_start_final_round ( ) ;
}
/* --- final explicit sync from each spoke → hub, then wait --- */
static void _str_start_final_round ( void ) {
struct str_state * s = gs ;
if ( ! s ) return ;
s - > final_round + + ;
s - > final_pending = STRESS_N - 1 ;
printf ( " final round %d: syncing %d spokes -> hub \n " , s - > final_round , s - > final_pending ) ;
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 ) ;
merkle_sync_start ( s - > inst [ i ] , s - > inst [ STRESS_HUB ] - > node_id , " stress " , _str_final_done , NULL ) ;
}
/* --- verify: tree sizes match pairwise --- */
static void _str_final_done ( uint64_t peer , const char * ns , int result , void * arg ) {
( void ) peer ; ( void ) ns ; ( void ) result ; ( void ) arg ;
if ( ! gs ) return ;
gs - > final_pending - - ;
if ( gs - > final_pending < = 0 ) _str_check_roots ( ) ;
}
static void _str_verify ( void ) {
struct str_state * s = gs ;
int ok = 1 ;
int na = _db_count_rows ( s - > inst [ 0 ] - > topo_sqlite_db , " stress " ) ;
/* tree sizes match pairwise */
int na = _db_count_rows ( s - > inst [ STRESS_HUB ] - > 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 ) ;
/* root (level=1, prefix=0) hashes identical */
const uint8_t * root0 = merkle_sync_get_hash ( s - > inst [ STRESS_HUB ] , " 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 ) ) ;
const uint8_t * ri = merkle_sync_get_hash ( s - > inst [ i ] , " stress " , 1 , 0 ) ;
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 --- */
/* 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 ,
@ -1295,12 +1289,82 @@ static void test_stress_spam(void) {
}
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 " ) ;
s - > phase = ok ? 1 : 2 ;
}
static void _str_check_roots ( void ) {
struct str_state * s = gs ;
if ( ! s ) return ;
const uint8_t * root0 = merkle_sync_get_hash ( s - > inst [ STRESS_HUB ] , " stress " , 1 , 0 ) ;
uint8_t root_copy [ MT_HASH_SIZE ] ;
memcpy ( root_copy , root0 , MT_HASH_SIZE ) ;
int converged = 1 ;
for ( int i = 1 ; i < STRESS_N ; i + + ) {
const uint8_t * ri = merkle_sync_get_hash ( s - > inst [ i ] , " stress " , 1 , 0 ) ;
if ( memcmp ( root_copy , ri , MT_HASH_SIZE ) ! = 0 ) { converged = 0 ; break ; }
}
if ( converged ) {
printf ( " roots converged after %d round(s) \n " , s - > final_round ) ;
_str_verify ( ) ;
} else if ( s - > final_round < STRESS_MAX_ROUNDS ) {
s - > round_timer = uasync_set_timeout ( s - > ua , STRESS_ROUND_GAP_TB , NULL , _str_round_cb , " str_round " ) ;
} else {
printf ( " no convergence after %d rounds \n " , s - > final_round ) ;
s - > phase = 2 ;
}
}
static void _str_round_cb ( void * arg ) {
( void ) arg ;
if ( ! gs ) return ;
gs - > round_timer = NULL ;
_str_start_final_round ( ) ;
}
static void _str_fail_timeout ( void * arg ) {
( void ) arg ;
if ( ! gs ) return ;
printf ( " SAFETY TIMEOUT \n " ) ;
gs - > phase = 2 ;
}
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 " ) ;
/* --- spam stop + safety timeout: событиями, без холостых ожиданий --- */
s - > stop_timer = uasync_set_timeout ( s - > ua , ( uint32_t ) ( STRESS_SPAM_MS * 10 ) , NULL , _str_stop_spam , " str_stop " ) ;
s - > safety_timer = uasync_set_timeout ( s - > ua , STRESS_SAFETY_TB , NULL , _str_fail_timeout , " str_safety " ) ;
printf ( " spamming for %dms... \n " , STRESS_SPAM_MS ) ;
/* Событийный цикл: фазы переключаются таймерами и done-коллбэками.
* poll с о г р а н и ч е н н ы м т а й м а у т о м ( 10 ms ) , т . к . uasync_poll ( ua , - 1 ) п р и у ж е
* п р о с р о ч е н н о м б л и ж а й ш е м т а й м е р е у х о д и т в б е с к о н е ч н ы й epoll_wait и н е
* о б р а б а т ы в а е т expired - т а й м е р ы . */
while ( s - > phase = = 0 )
uasync_poll ( s - > ua , 100 ) ;
if ( s - > phase ! = 1 ) { FAIL ( " convergence failed (phase=%d) " , s - > phase ) ; _str_cleanup ( ) ; return ; }
PASS ( ) ;
_str_cleanup ( ) ;
}