@ -828,8 +828,9 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, uint64_t src, const
}
// ---- SEND_DATA batch helper ----
// Wire format: [id:8][ts:8][author:8][dlen:4][data][sig_len:1=64][sig:64]
static void si_send_data_batch ( struct DB_SYNC_INSTANCE * si , uint64_t dst , uint32_t from , uint32_t count , int allocated_buf )
// Wire format: [id:8][ts:8][node_id:8][dlen:4][data][sig_len:1=64][sig:64]
// Returns: number of records sent, or -1 on SQL error
static int si_send_data_batch ( struct DB_SYNC_INSTANCE * si , uint64_t dst , uint32_t from , uint32_t count , int allocated_buf )
{
uint8_t sbuf_stack [ 8192 ] ;
uint8_t * buf = allocated_buf ? u_malloc ( 8192 ) : NULL ;
@ -842,10 +843,16 @@ static void si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32
uint16_t * rcp = ( uint16_t * ) ( buf + off ) ; off + = 2 ;
sqlite3_stmt * stmt ;
si_prep ( si , & stmt ,
" SELECT id,timestamp,author ,data,author_signature "
int prep_rc = si_prep ( si , & stmt ,
" SELECT id,timestamp,node_id ,data,author_signature "
" FROM \" %s \" ORDER BY timestamp, author_signature "
" LIMIT ? OFFSET ? " ) ;
if ( prep_rc ! = SQLITE_OK | | ! stmt ) {
DEBUG_ERROR ( DEBUG_CATEGORY_DB_SYNC , " sync [%s:%04llX] SEND_DATA batch: si_prep failed rc=%d " ,
SI_SHRT ( si ) , ( unsigned long long ) ( dst > > 16 ) , prep_rc ) ;
if ( allocated_buf ) u_free ( buf ) ;
return - 1 ;
}
sqlite3_bind_int64 ( stmt , 1 , ( sqlite3_int64 ) count ) ;
sqlite3_bind_int64 ( stmt , 2 , ( sqlite3_int64 ) from ) ;
while ( sqlite3_step ( stmt ) = = SQLITE_ROW ) {
@ -873,7 +880,12 @@ static void si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32
db_sync_send ( si , dst , buf , off ) ;
DEBUG_INFO ( DEBUG_CATEGORY_DB_SYNC , " sync [%s:%04llX] → SEND DATA: %u records from pos=%u " ,
SI_SHRT ( si ) , ( unsigned long long ) ( dst > > 16 ) , rc , from ) ;
if ( rc = = 0 & & count > 0 ) {
DEBUG_ERROR ( DEBUG_CATEGORY_DB_SYNC , " sync [%s:%04llX] SEND_DATA batch: loaded 0 records from pos=%u count=%u — SQL error or empty range " ,
SI_SHRT ( si ) , ( unsigned long long ) ( dst > > 16 ) , from , count ) ;
}
if ( allocated_buf ) u_free ( buf ) ;
return ( int ) rc ;
}
static void db_handle_refine ( struct DB_SYNC_INSTANCE * si , uint64_t src , const uint8_t * p , size_t len )
@ -892,9 +904,16 @@ static void db_handle_refine(struct DB_SYNC_INSTANCE* si, uint64_t src, const ui
SI_SHRT ( si ) , ( unsigned long long ) ( src > > 16 ) , from , mc ) ;
return ;
}
si_send_data_batch ( si , src , from , scnt , 1 ) ;
DEBUG_INFO ( DEBUG_CATEGORY_DB_SYNC , " sync [%s:%04llX] ← REFINE: range narrowed to single pos → sending %u records from pos=%u " ,
SI_SHRT ( si ) , ( unsigned long long ) ( src > > 16 ) , scnt , from ) ;
int sent = si_send_data_batch ( si , src , from , scnt , 1 ) ;
if ( sent < 0 ) {
DEBUG_ERROR ( DEBUG_CATEGORY_DB_SYNC , " sync [%s:%04llX] ← REFINE: si_send_data_batch FAILED from=%u count=%u — resetting sync " ,
SI_SHRT ( si ) , ( unsigned long long ) ( src > > 16 ) , from , scnt ) ;
struct SI_PEER * sp2 = si_peer_find ( si , src ) ;
if ( sp2 ) sp2 - > sync_state = 0 ;
return ;
}
DEBUG_INFO ( DEBUG_CATEGORY_DB_SYNC , " sync [%s:%04llX] ← REFINE: range narrowed to single pos → sent %d records from pos=%u " ,
SI_SHRT ( si ) , ( unsigned long long ) ( src > > 16 ) , sent , from ) ;
return ;
}
@ -916,10 +935,21 @@ static void db_handle_refine(struct DB_SYNC_INSTANCE* si, uint64_t src, const ui
}
uint32_t mc = db_count ( si ) ;
uint32_t scnt = 4 ;
if ( fm > to | | fm > = mc ) scnt = 0 ;
si_send_data_batch ( si , src , fm , scnt , 0 ) ;
DEBUG_INFO ( DEBUG_CATEGORY_DB_SYNC , " sync [%s:%04llX] ← REFINE: range [%u..%u] %u checkpoints → last MATCH at pos=%u, first DIFF at pos=%u → sending %u records from pos=%u " ,
SI_SHRT ( si ) , ( unsigned long long ) ( src > > 16 ) , from , to , hc , last_match , first_diff , scnt , fm ) ;
if ( fm > to | | fm > = mc ) {
scnt = 0 ;
DEBUG_WARN ( DEBUG_CATEGORY_DB_SYNC , " sync [%s:%04llX] ← REFINE: range [%u..%u] → fm=%u out of bounds (to=%u mc=%u) — sending 0 records " ,
SI_SHRT ( si ) , ( unsigned long long ) ( src > > 16 ) , from , to , fm , to , mc ) ;
}
int sent2 = si_send_data_batch ( si , src , fm , scnt , 0 ) ;
if ( sent2 < 0 ) {
DEBUG_ERROR ( DEBUG_CATEGORY_DB_SYNC , " sync [%s:%04llX] ← REFINE: si_send_data_batch FAILED from=%u count=%u — resetting sync " ,
SI_SHRT ( si ) , ( unsigned long long ) ( src > > 16 ) , fm , scnt ) ;
struct SI_PEER * sp2 = si_peer_find ( si , src ) ;
if ( sp2 ) sp2 - > sync_state = 0 ;
return ;
}
DEBUG_INFO ( DEBUG_CATEGORY_DB_SYNC , " sync [%s:%04llX] ← REFINE: range [%u..%u] %u checkpoints → last MATCH at pos=%u, first DIFF at pos=%u → sent %d records from pos=%u " ,
SI_SHRT ( si ) , ( unsigned long long ) ( src > > 16 ) , from , to , hc , last_match , first_diff , sent2 , fm ) ;
}
// ---- Parse one record from SEND_DATA/PUSH wire format ----
@ -1010,9 +1040,15 @@ static void db_handle_send_data(struct DB_SYNC_INSTANCE* si, uint64_t src, const
if ( mc > pk & & sp & & sp - > sync_state = = 1 ) {
uint32_t scnt = mc - pk ; if ( scnt > DB_SEND_DATA_MAX ) scnt = DB_SEND_DATA_MAX ;
si_send_data_batch ( si , src , pk , scnt , 1 ) ;
DEBUG_INFO ( DEBUG_CATEGORY_DB_SYNC , " sync [%s:%04llX] → DATA: requesting next %u records from pos=%u (%u remaining behind) " ,
SI_SHRT ( si ) , ( unsigned long long ) ( src > > 16 ) , scnt , pk , mc - pk ) ;
int pushed = si_send_data_batch ( si , src , pk , scnt , 1 ) ;
if ( pushed < 0 ) {
DEBUG_ERROR ( DEBUG_CATEGORY_DB_SYNC , " sync [%s:%04llX] ← DATA: si_send_data_batch FAILED from=%u count=%u — resetting sync " ,
SI_SHRT ( si ) , ( unsigned long long ) ( src > > 16 ) , pk , scnt ) ;
sp - > sync_state = 0 ;
return ;
}
DEBUG_INFO ( DEBUG_CATEGORY_DB_SYNC , " sync [%s:%04llX] → DATA: sent %d records from pos=%u (%u remaining behind) " ,
SI_SHRT ( si ) , ( unsigned long long ) ( src > > 16 ) , pushed , pk , mc - pk ) ;
} else if ( mc > pk ) {
DEBUG_INFO ( DEBUG_CATEGORY_DB_SYNC , " sync [%s:%04llX] ← DATA: got %u records [%u..%u] but %u remaining — peer stopped sending (sync_state=%d) " ,
SI_SHRT ( si ) , ( unsigned long long ) ( src > > 16 ) , received , from , from + received - 1 , mc - pk ,
@ -1113,6 +1149,8 @@ static void db_handle_push(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint
int ret = db_record_insert ( si , rid , rts , rauthor , ( const char * ) rdata , rdlen , rsig , 0 ) ;
if ( ret = = 0 ) {
uint32_t ins_pos = si_find_pos ( si , rts , rsig ) ;
uint32_t total = db_count ( si ) ;
db_cascade_from ( si , ins_pos ) ;
int adj_count = 0 ;
for ( int j = 0 ; j < si - > peer_count ; j + + ) {
if ( si - > peers [ j ] . synced_pos > = ins_pos ) { si - > peers [ j ] . synced_pos = ins_pos > 0 ? ins_pos - 1 : 0 ; adj_count + + ; }
@ -1123,10 +1161,10 @@ static void db_handle_push(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint
memcpy ( ack + 1 , & rts , 8 ) ;
memcpy ( ack + 9 , & rauthor , 8 ) ;
db_sync_send ( si , src , ack , 17 ) ;
DEBUG_INFO ( DEBUG_CATEGORY_DB_SYNC , " sync [%s:%04llX] ← PUSH: msg id=%llu author=%016llX ts=%llu → inserted at pos=%u (of %u total), ACK sent " ,
DEBUG_INFO ( DEBUG_CATEGORY_DB_SYNC , " sync [%s:%04llX] ← PUSH: msg id=%llu author=%016llX ts=%llu → inserted at pos=%u (of %u total), cascade from pos=%u, ACK sent " ,
SI_SHRT ( si ) , ( unsigned long long ) ( src > > 16 ) ,
( unsigned long long ) rid , ( unsigned long long ) rauthor , ( unsigned long long ) rts , ins_pos , db_count ( si ) ) ;
if ( adj_count > 0 & & ins_pos < db_count ( si ) - 1 ) {
( unsigned long long ) rid , ( unsigned long long ) rauthor , ( unsigned long long ) rts , ins_pos , total , ins_pos ) ;
if ( adj_count > 0 & & ins_pos < total - 1 ) {
DEBUG_INFO ( DEBUG_CATEGORY_DB_SYNC , " sync [%s:----] PUSH: inserted at pos=%u NOT at tail → reset synced_pos of %d peers from ≥%u back to %u " ,
SI_SHRT ( si ) , ins_pos , adj_count , ins_pos , ins_pos > 0 ? ins_pos - 1 : 0 ) ;
}
@ -1217,17 +1255,6 @@ static void db_sync_on_conn_status(struct ETCP_CONN* conn, int status, void* arg
}
}
#if 0 /* replaced by db_sync_on_conn_status */
static void db_sync_on_new_conn ( struct ETCP_CONN * conn , void * arg )
{
( void ) arg ;
if ( ! conn | | ! conn - > instance | | ! conn - > instance - > db_sync ) return ;
DEBUG_TRACE ( DEBUG_CATEGORY_DB_SYNC , " conn=%p node=%016llx " , ( void * ) conn , ( unsigned long long ) conn - > peer_node_id ) ;
etcp_conn_add_up_cbk ( conn , db_sync_on_conn_up , NULL ) ;
etcp_conn_add_down_cbk ( conn , db_sync_on_conn_down , NULL ) ;
}
# endif
static void db_sync_on_conn_up ( struct ETCP_CONN * conn , void * arg )
{
( void ) arg ;
@ -1344,6 +1371,39 @@ static void db_verify_chain(struct DB_SYNC_INSTANCE* si)
}
}
int db_sync_chain_verify ( struct DB_SYNC_INSTANCE * si )
{
if ( ! si | | ! si - > enabled ) return 0 ;
uint32_t mc = db_count ( si ) ;
if ( mc = = 0 ) return 0 ;
uint8_t prev_ch [ 32 ] ; memset ( prev_ch , 0 , 32 ) ;
uint8_t exp_ch [ 32 ] , stored_ch [ 32 ] ;
sqlite3_stmt * stmt ;
if ( si_prep ( si , & stmt ,
" SELECT id,timestamp,node_id,author_signature,chain_hash FROM \" %s \" "
" ORDER BY timestamp, author_signature " ) ! = SQLITE_OK )
return - 1 ;
for ( uint32_t pos = 0 ; sqlite3_step ( stmt ) = = SQLITE_ROW ; pos + + ) {
uint64_t rid = ( uint64_t ) sqlite3_column_int64 ( stmt , 0 ) ;
uint64_t rts = ( uint64_t ) sqlite3_column_int64 ( stmt , 1 ) ;
uint64_t rauth = ( uint64_t ) sqlite3_column_int64 ( stmt , 2 ) ;
const void * sig_blob = sqlite3_column_blob ( stmt , 3 ) ;
uint8_t sig [ DB_SIG_SIZE ] ;
if ( sig_blob ) memcpy ( sig , sig_blob , DB_SIG_SIZE ) ; else memset ( sig , 0 , DB_SIG_SIZE ) ;
const void * b = sqlite3_column_blob ( stmt , 4 ) ;
if ( b & & sqlite3_column_bytes ( stmt , 4 ) > = 32 ) memcpy ( stored_ch , b , 32 ) ;
else memset ( stored_ch , 0 , 32 ) ;
db_chain_hash_compute ( prev_ch , rid , rts , rauth , sig , exp_ch ) ;
if ( memcmp ( exp_ch , stored_ch , 32 ) ! = 0 ) { sqlite3_finalize ( stmt ) ; return 1 ; }
memcpy ( prev_ch , exp_ch , 32 ) ;
}
sqlite3_finalize ( stmt ) ;
return 0 ;
}
// ============================================================
// Timers
// ============================================================