@ -38,6 +38,7 @@ struct SI_PEER {
uint64_t node_id ;
uint32_t synced_pos ;
uint8_t sync_state ;
uint64_t sync_start_tb ;
} ;
struct DB_SYNC_INSTANCE {
@ -224,6 +225,7 @@ static struct SI_PEER* si_peer_add(struct DB_SYNC_INSTANCE* si, uint64_t node_id
p - > node_id = node_id ;
p - > synced_pos = 0 ;
p - > sync_state = 0 ;
p - > sync_start_tb = 0 ;
return p ;
}
@ -581,12 +583,16 @@ static void si_delivery_update(struct DB_SYNC_INSTANCE* si, uint64_t ts, uint64_
// Sync protocol handlers
// ============================================================
static void db_handle_error ( struct DB_SYNC * db , uint64_t src , const uint8_t * p , size_t len )
static void db_handle_error ( struct DB_SYNC_INSTANCE * si , uint64_t src , const uint8_t * p , size_t len )
{
if ( len < 1 ) return ;
DEBUG_TRACE ( DEBUG_CATEGORY_DB_SYNC , " from=%016llx code=%u " , ( unsigned long long ) src , p [ 0 ] ) ;
DEBUG_WARN ( DEBUG_CATEGORY_DB_SYNC , " db_sync: ERROR from %016llx code=%u " , ( unsigned long long ) src , p [ 0 ] ) ;
( void ) db ;
if ( len < 1 | | ! si ) return ;
const char * what = ( p [ 0 ] = = DB_ERR_NOT_FOUND ) ? " NOT_FOUND " :
( p [ 0 ] = = DB_ERR_DISABLED ) ? " DISABLED " : " UNKNOWN " ;
DEBUG_ERROR ( DEBUG_CATEGORY_DB_SYNC ,
" db_sync: ERROR from %016llx code=%u (%s) tbl=%s hash=%016llx " ,
( unsigned long long ) src , p [ 0 ] , what , SI_TBL ( si ) , ( unsigned long long ) si - > hash ) ;
struct SI_PEER * sp = si_peer_find ( si , src ) ;
if ( sp ) sp - > sync_state = 0 ;
}
static void db_handle_init_sync ( struct DB_SYNC_INSTANCE * si , uint64_t src , const uint8_t * p , size_t len )
@ -1075,7 +1081,7 @@ static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry)
case DB_MSG_PUSH : DEBUG_TRACE ( DEBUG_CATEGORY_DB_SYNC , " → handle PUSH " ) ; db_handle_push ( si , src , payload , plen ) ; break ;
case DB_MSG_ACK_PUSH : DEBUG_TRACE ( DEBUG_CATEGORY_DB_SYNC , " → handle ACK_PUSH " ) ; db_handle_ack_push ( si , src , payload , plen ) ; break ;
case DB_MSG_SYNC_DONE : DEBUG_TRACE ( DEBUG_CATEGORY_DB_SYNC , " → handle SYNC_DONE " ) ; db_handle_sync_done ( si , src , payload , plen ) ; break ;
case DB_MSG_ERROR : DEBUG_TRACE ( DEBUG_CATEGORY_DB_SYNC , " → handle ERROR " ) ; db_handle_error ( db , src , payload , plen ) ; break ;
case DB_MSG_ERROR : DEBUG_TRACE ( DEBUG_CATEGORY_DB_SYNC , " → handle ERROR " ) ; db_handle_error ( si , src , payload , plen ) ; break ;
default :
DEBUG_WARN ( DEBUG_CATEGORY_DB_SYNC , " unknown msg type 0x%02x from %016llx " , type , ( unsigned long long ) src ) ;
break ;
@ -1160,6 +1166,8 @@ static void db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si, uint64_t pid)
{
uint32_t mc = db_count ( si ) ;
DEBUG_TRACE ( DEBUG_CATEGORY_DB_SYNC , " pid=%016llx my=%u tbl=%s " , ( unsigned long long ) pid , mc , SI_TBL ( si ) ) ;
struct SI_PEER * p = si_peer_find ( si , pid ) ;
if ( p ) p - > sync_start_tb = get_time_tb ( ) ;
uint8_t msg [ 5 ] ;
msg [ 0 ] = DB_MSG_INIT_SYNC ;
@ -1252,6 +1260,25 @@ static void db_sync_peer_check_cb(void* arg)
else { total_skipped + + ; }
}
{
uint64_t now = get_time_tb ( ) ;
uint64_t to_tb = DB_SYNC_SYNC_TIMEOUT * 10000u ;
for ( int i = 0 ; i < db - > instance_count ; i + + ) {
struct DB_SYNC_INSTANCE * si = & db - > instances [ i ] ;
if ( ! si - > enabled ) continue ;
for ( int j = 0 ; j < si - > peer_count ; j + + ) {
struct SI_PEER * p = & si - > peers [ j ] ;
if ( p - > sync_state = = 1 & & p - > sync_start_tb > 0 & & now - p - > sync_start_tb > to_tb ) {
DEBUG_WARN ( DEBUG_CATEGORY_DB_SYNC ,
" db_sync: no response for peer=%016llx tbl=%s elapsed=%llu ms, resetting sync_state " ,
( unsigned long long ) p - > node_id , SI_TBL ( si ) ,
( unsigned long long ) ( ( now - p - > sync_start_tb ) / 10 ) ) ;
p - > sync_state = 0 ;
}
}
}
}
DEBUG_INFO ( DEBUG_CATEGORY_DB_SYNC , " peer_check: instances=%d synced=%d skipped=%d " ,
db - > instance_count , total_synced , total_skipped ) ;