@ -12,7 +12,6 @@
# include <sqlite3.h>
# define MS_ID "merkle_sync"
# define MS_SYNC_TIMEOUT_MS 10000
# define MS_BG_INTERVAL_MS 100
static struct merkle_sync * g_merkle = NULL ;
@ -25,8 +24,6 @@ struct ms_session {
uint64_t peer ;
uint8_t active ;
uint8_t synced ;
void * timer ;
uint8_t retries ;
merkle_sync_done_cb done_cb ;
void * cb_arg ;
} ;
@ -168,11 +165,19 @@ static int _get_level_hashes(struct merkle_sync* ms, const char* ns,
if ( level > = MT_MAX_LEVEL ) return 0 ;
sqlite3 * db = _db ( ms - > inst ) ; if ( ! db ) return - 1 ;
int parent_shift = 63 - ( int ) level * 5 ; if ( parent_shift < 0 ) parent_shift = 0 ;
int next_shift = 63 - ( ( int ) level + 1 ) * 5 ; if ( next_shift < 0 ) next_shift = 0 ;
uint64_t range_end = prefix | ( 0x1FULL < < parent_shift ) | ( 0x1FULL < < ( next_shift > 0 ? next_shift : 0 ) ) ;
int next_shift ;
if ( level = = 0 ) {
next_shift = 63 - 5 ;
} else {
next_shift = 63 - ( ( int ) level + 1 ) * 5 ; if ( next_shift < 0 ) next_shift = 0 ;
}
uint64_t range_end = ( level = = 0 ) ? UINT64_MAX
: prefix | ( 0x1FULL < < ( next_shift > 0 ? next_shift : 0 ) ) ;
sqlite3_stmt * stmt = NULL ;
int query_level = ( int ) ( level + 1 ) ;
if ( sqlite3_prepare_v2 ( db ,
" SELECT prefix64, hash FROM merkle_tree_hash "
" WHERE namespace=? AND level=? AND prefix64>=? AND prefix64<=? "
@ -181,7 +186,7 @@ static int _get_level_hashes(struct merkle_sync* ms, const char* ns,
return - 1 ;
}
sqlite3_bind_text ( stmt , 1 , ns , - 1 , SQLITE_STATIC ) ;
sqlite3_bind_int ( stmt , 2 , ( int ) ( level + 1 ) ) ;
sqlite3_bind_int ( stmt , 2 , query_level ) ;
sqlite3_bind_int64 ( stmt , 3 , ( sqlite3_int64 ) prefix ) ;
sqlite3_bind_int64 ( stmt , 4 , ( sqlite3_int64 ) range_end ) ;
while ( sqlite3_step ( stmt ) = = SQLITE_ROW ) {
@ -315,34 +320,6 @@ static int _send_batch(struct merkle_sync* ms, uint64_t peer, const char* ns,
/* ── Sessions ── */
static void _session_start_timer ( struct ms_session * s ) ;
static void _session_timeout_cb ( void * arg ) {
struct ms_session * s = ( struct ms_session * ) arg ;
if ( ! s | | ! s - > active | | ! g_merkle | | ! g_merkle - > initialized ) return ;
DEBUG_TRACE ( DEBUG_CATEGORY_DB_SYNC , " %s: timeout peer=%016llx ns=%s retry=%d/3 " , MS_ID , ( unsigned long long ) s - > peer , s - > ns , s - > retries ) ;
s - > timer = NULL ;
s - > retries + + ;
if ( s - > retries > 3 ) {
DEBUG_TRACE ( DEBUG_CATEGORY_DB_SYNC , " %s: → max retries, done(ERR_TIMEOUT) " , MS_ID ) ;
DEBUG_WARN ( DEBUG_CATEGORY_DB_SYNC , " %s: sync timeout peer=%016llx ns=%s " , MS_ID ,
( unsigned long long ) s - > peer , s - > ns ) ;
s - > active = 0 ;
if ( s - > done_cb ) s - > done_cb ( s - > peer , s - > ns , MT_ERR_TIMEOUT , s - > cb_arg ) ;
return ;
}
DEBUG_INFO ( DEBUG_CATEGORY_DB_SYNC , " %s: retry %d peer=%016llx ns=%s " , MS_ID ,
s - > retries , ( unsigned long long ) s - > peer , s - > ns ) ;
_send_hashes ( g_merkle , s - > peer , s - > ns , 1 , 0 , 1 , 0 ) ;
_session_start_timer ( s ) ;
}
static void _session_start_timer ( struct ms_session * s ) {
if ( ! g_merkle | | ! g_merkle - > inst ) return ;
s - > timer = uasync_set_timeout ( g_merkle - > inst - > ua , ( uint32_t ) ( MS_SYNC_TIMEOUT_MS * 10 ) ,
s , _session_timeout_cb , " ms_sync " ) ;
}
static struct ms_session * _session_find ( struct merkle_sync * ms , uint64_t peer , const char * ns ) {
for ( struct ms_session * s = ms - > sessions ; s ; s = s - > next )
if ( s - > peer = = peer & & strcmp ( s - > ns , ns ) = = 0 ) return s ;
@ -351,7 +328,6 @@ static struct ms_session* _session_find(struct merkle_sync* ms, uint64_t peer, c
static void _session_done ( struct ms_session * s , int result ) {
s - > active = 0 ;
if ( s - > timer ) { uasync_cancel_timeout ( g_merkle - > inst - > ua , s - > timer ) ; s - > timer = NULL ; }
if ( result = = MT_OK ) { s - > synced = 1 ; DEBUG_INFO ( DEBUG_CATEGORY_DB_SYNC , " %s: session SYNCED peer=%016llx ns=%s " , MS_ID , ( unsigned long long ) s - > peer , s - > ns ) ; }
else { s - > synced = 0 ; DEBUG_WARN ( DEBUG_CATEGORY_DB_SYNC , " %s: session FAILED peer=%016llx ns=%s result=%d " , MS_ID , ( unsigned long long ) s - > peer , s - > ns , result ) ; }
if ( s - > done_cb ) s - > done_cb ( s - > peer , s - > ns , result , s - > cb_arg ) ;
@ -370,7 +346,6 @@ static void _handle_hashes(struct merkle_sync* ms, uint64_t peer, const char* ns
DEBUG_TRACE ( DEBUG_CATEGORY_DB_SYNC , " %s: handle_hashes peer=%016llx ns=%s L%d/P%016llx is_data=%d " , MS_ID , ( unsigned long long ) peer , ns , level , ( unsigned long long ) prefix , is_data ) ;
struct ms_session * s = _session_find ( ms , peer , ns ) ;
if ( s & & s - > timer ) { uasync_cancel_timeout ( g_merkle - > inst - > ua , s - > timer ) ; s - > timer = NULL ; }
if ( is_data ) {
if ( paylen < 2 ) return ;
@ -447,7 +422,6 @@ static void _handle_hashes(struct merkle_sync* ms, uint64_t peer, const char* ns
_send_msg ( ms , peer , rbuf , ( size_t ) ( wr - rbuf ) ) ;
u_free ( rbuf ) ;
}
if ( s ) _session_start_timer ( s ) ;
} else {
if ( s ) _session_done ( s , MT_OK ) ;
}
@ -472,11 +446,6 @@ static void _handle_request(struct merkle_sync* ms, uint64_t peer, const char* n
}
if ( bc > 0 ) {
_send_batch ( ms , peer , ns , buckets , bc ) ;
struct ms_session * s = _session_find ( ms , peer , ns ) ;
if ( s & & s - > active ) {
if ( s - > timer ) { uasync_cancel_timeout ( ms - > inst - > ua , s - > timer ) ; s - > timer = NULL ; }
_session_start_timer ( s ) ;
}
}
}
@ -506,7 +475,7 @@ static void _handle_batch(struct merkle_sync* ms, uint64_t peer, const char* ns,
uint8_t sub_pl [ 4096 ] ; size_t sub_len = 0 ;
uint8_t next_lvl = ( uint8_t ) ( lvl < MT_MAX_LEVEL ? lvl + 1 : lvl ) ;
sub_pl [ sub_len + + ] = next_lvl ;
sub_pl [ sub_len + + ] = pb_i ;
sub_pl [ sub_len + + ] = merkle_sync_prefix_bytes ( next_lvl ) ;
_prefix_write ( sub_pl + sub_len , pr , pb_i ) ; sub_len + = pb_i ;
sub_pl [ sub_len + + ] = 0 ; /* is_data=0 */
uint32_t bm ; memcpy ( & bm , bp , 4 ) ; bp + = 4 ; rem - = 4 ;
@ -690,7 +659,6 @@ void merkle_sync_destroy(struct UTUN_INSTANCE* inst) {
struct ms_session * s = ms - > sessions ;
while ( s ) { struct ms_session * next = s - > next ;
if ( s - > timer ) uasync_cancel_timeout ( inst - > ua , s - > timer ) ;
u_free ( s ) ; s = next ; }
u_free ( ms ) ;
@ -715,14 +683,12 @@ int merkle_sync_start(struct UTUN_INSTANCE* inst, uint64_t peer,
ms - > sessions = s ;
} else {
/* replace callback if session already exists */
if ( s - > timer ) { uasync_cancel_timeout ( inst - > ua , s - > timer ) ; s - > timer = NULL ; }
s - > synced = 0 ;
}
s - > active = 1 ; s - > retries = 0 ;
s - > active = 1 ;
s - > done_cb = done_cb ; s - > cb_arg = arg ;
_send_hashes ( ms , peer , ns , 1 , 0 , 1 , 0 ) ;
_session_start_timer ( s ) ;
_send_hashes ( ms , peer , ns , 0 , 0 , 0 , 0 ) ;
return 0 ;
}
@ -734,7 +700,6 @@ void merkle_sync_cancel(struct UTUN_INSTANCE* inst, uint64_t peer, const char* n
while ( * p ) {
struct ms_session * s = * p ;
if ( s - > peer = = peer & & strcmp ( s - > ns , ns ) = = 0 ) {
if ( s - > timer ) { uasync_cancel_timeout ( inst - > ua , s - > timer ) ; s - > timer = NULL ; }
* p = s - > next ; u_free ( s ) ;
return ;
}