@ -107,9 +107,11 @@ static struct channel_cache* cs_find(struct chat_sync* cs, const char* ch_id) {
static void cs_handle_init_sync ( struct chat_sync * cs , uint64_t peer ,
const char * ch_id , const uint8_t * pl , size_t len ) {
if ( len < 36 ) return ;
if ( len < 36 ) { DEBUG_WARN ( DEBUG_CATEGORY_CONNECTIVITY , " %s: INIT_SYNC too short len=%zu peer=%016llx " , CS_ID , len , ( unsigned long long ) peer ) ; return ; }
uint32_t peer_count ; memcpy ( & peer_count , pl , 4 ) ;
uint8_t peer_last_ch [ 32 ] ; memcpy ( peer_last_ch , pl + 4 , 32 ) ;
DEBUG_INFO ( DEBUG_CATEGORY_CONNECTIVITY , " %s: RECV INIT_SYNC peer=%016llx ch=%s peer_count=%u my_count=%u " ,
CS_ID , ( unsigned long long ) peer , ch_id , peer_count , chat_core_count ( ch_id ) ) ;
uint32_t my_count = chat_core_count ( ch_id ) ;
uint32_t tp = peer_count < my_count ? peer_count : my_count ;
@ -138,47 +140,74 @@ static void cs_handle_init_sync(struct chat_sync* cs, uint64_t peer,
static void cs_handle_init_resp ( struct chat_sync * cs , uint64_t peer ,
const char * ch_id , const uint8_t * pl , size_t len ) {
if ( len < 37 ) return ;
if ( len < 5 ) {
DEBUG_WARN ( DEBUG_CATEGORY_CONNECTIVITY , " %s: INIT_RESP too short len=%zu peer=%016llx ch=%s " ,
CS_ID , len , ( unsigned long long ) peer , ch_id ) ;
return ;
}
uint32_t peer_count = * ( const uint32_t * ) pl ;
uint32_t test_pos = * ( const uint32_t * ) ( pl + 4 ) ;
uint8_t peer_ch [ 32 ] ; memcpy ( peer_ch , pl + 8 , 32 ) ;
uint8_t sparse_count = pl [ 40 ] ; ( void ) sparse_count ;
uint8_t sparse_count = pl [ 4 ] ;
uint8_t my_ch [ 32 ] ;
chat_core_chain_hash_at ( ch_id , test_pos , my_ch ) ;
struct channel_cache * ch = cs_find ( cs , ch_id ) ;
uint32_t my_count = chat_core_count ( ch_id ) ;
if ( memcmp ( my_ch , peer_ch , 32 ) = = 0 ) {
struct channel_cache * ch = cs_find ( cs , ch_id ) ;
if ( ch & & peer_count > test_pos + 1 ) {
uint32_t from = test_pos + 1 ;
uint8_t snd [ 7 ] ; snd [ 0 ] = CS_MSG_SEND_DATA ;
memcpy ( snd + 1 , & from , 4 ) ; uint16_t z = 0 ; memcpy ( snd + 5 , & z , 2 ) ;
cs_send ( cs , ch_id , peer , snd , 7 ) ;
} else {
if ( ch ) ch - > synced = CS_SYNC_DONE ;
DEBUG_INFO ( DEBUG_CATEGORY_CONNECTIVITY , " %s: RECV INIT_RESP peer=%016llx ch=%s peer_count=%u my_count=%u len=%zu " ,
CS_ID , ( unsigned long long ) peer , ch_id , peer_count , my_count , len ) ;
if ( len > = 41 ) {
/* extended format: peer_count(4) + test_pos(4) + peer_ch(32) + sparse_count(1) */
uint32_t test_pos = * ( const uint32_t * ) ( pl + 4 ) ;
uint8_t peer_ch [ 32 ] ; memcpy ( peer_ch , pl + 8 , 32 ) ;
sparse_count = pl [ 40 ] ;
uint8_t my_ch [ 32 ] ; chat_core_chain_hash_at ( ch_id , test_pos , my_ch ) ;
if ( memcmp ( my_ch , peer_ch , 32 ) = = 0 ) {
if ( ch & & peer_count > test_pos + 1 ) {
uint32_t from = test_pos + 1 ;
uint8_t snd [ 7 ] ; snd [ 0 ] = CS_MSG_SEND_DATA ;
memcpy ( snd + 1 , & from , 4 ) ; uint16_t z = 0 ; memcpy ( snd + 5 , & z , 2 ) ;
cs_send ( cs , ch_id , peer , snd , 7 ) ;
DEBUG_INFO ( DEBUG_CATEGORY_CONNECTIVITY , " %s: INIT_RESP chain_match, request data from=%u " , CS_ID , from ) ;
} else {
if ( ch ) ch - > synced = CS_SYNC_DONE ;
DEBUG_INFO ( DEBUG_CATEGORY_CONNECTIVITY , " %s: INIT_RESP chain_match, synced " , CS_ID ) ;
}
return ;
}
DEBUG_INFO ( DEBUG_CATEGORY_CONNECTIVITY , " %s: INIT_RESP chain_mismatch, request from start " , CS_ID ) ;
uint32_t from = 0 ;
uint8_t snd [ 7 ] ; snd [ 0 ] = CS_MSG_SEND_DATA ;
memcpy ( snd + 1 , & from , 4 ) ; uint16_t z = 0 ; memcpy ( snd + 5 , & z , 2 ) ;
cs_send ( cs , ch_id , peer , snd , 7 ) ;
return ;
}
/* mismatch — request from start */
uint32_t from = 0 ;
uint8_t snd [ 7 ] ; snd [ 0 ] = CS_MSG_SEND_DATA ;
memcpy ( snd + 1 , & from , 4 ) ; uint16_t z = 0 ; memcpy ( snd + 5 , & z , 2 ) ;
cs_send ( cs , ch_id , peer , snd , 7 ) ;
/* short format (synced): peer_count(4) + sparse_count(1) */
if ( peer_count > my_count ) {
uint32_t from = my_count ;
uint8_t snd [ 7 ] ; snd [ 0 ] = CS_MSG_SEND_DATA ;
memcpy ( snd + 1 , & from , 4 ) ; uint16_t z = 0 ; memcpy ( snd + 5 , & z , 2 ) ;
cs_send ( cs , ch_id , peer , snd , 7 ) ;
DEBUG_INFO ( DEBUG_CATEGORY_CONNECTIVITY , " %s: INIT_RESP peer ahead (peer=%u my=%u), request from=%u " ,
CS_ID , peer_count , my_count , from ) ;
} else {
if ( ch ) ch - > synced = CS_SYNC_DONE ;
DEBUG_INFO ( DEBUG_CATEGORY_CONNECTIVITY , " %s: INIT_RESP synced (peer=%u <= my=%u) " ,
CS_ID , peer_count , my_count ) ;
}
}
/* ── SEND_DATA handler ── */
static void cs_handle_send_data ( struct chat_sync * cs , uint64_t peer ,
const char * ch_id , const uint8_t * pl , size_t len ) {
if ( len < 6 ) return ;
if ( len < 6 ) { DEBUG_WARN ( DEBUG_CATEGORY_CONNECTIVITY , " %s: SEND_DATA too short len=%zu peer=%016llx " , CS_ID , len , ( unsigned long long ) peer ) ; return ; }
uint32_t from = * ( const uint32_t * ) pl ;
uint16_t count = * ( const uint16_t * ) ( pl + 4 ) ;
if ( count = = 0 ) {
/* peer requests OUR data starting from 'from' */
uint32_t cur = chat_core_cursor_open ( ch_id ) ;
if ( ! cur ) return ;
if ( ! cur ) { DEBUG_WARN ( DEBUG_CATEGORY_CONNECTIVITY , " %s: cursor_open failed ch=%s " , CS_ID , ch_id ) ; return ; }
uint16_t sent = 0 ;
uint8_t rec_buf [ 4096 ] ; size_t rec_len ;
@ -196,10 +225,13 @@ static void cs_handle_send_data(struct chat_sync* cs, uint64_t peer,
sent + + ;
}
chat_core_cursor_close ( cur ) ;
DEBUG_INFO ( DEBUG_CATEGORY_CONNECTIVITY , " %s: SEND_DATA sent=%u records to=%016llx ch=%s from=%u " ,
CS_ID , sent , ( unsigned long long ) peer , ch_id , from ) ;
return ;
}
/* peer sent US data — insert each record */
uint16_t inserted = 0 ;
const uint8_t * ptr = pl + 6 ;
size_t remain = len - 6 ;
@ -218,9 +250,12 @@ static void cs_handle_send_data(struct chat_sync* cs, uint64_t peer,
ptr + = 32 ; remain - = 32 ; /* chain_hash */
size_t reclen = ( size_t ) ( ptr - rec_start ) ;
chat_core_insert_record ( ch_id , rec_start , reclen ) ;
if ( chat_core_insert_record ( ch_id , rec_start , reclen ) = = 0 ) inserted + + ;
}
DEBUG_INFO ( DEBUG_CATEGORY_CONNECTIVITY , " %s: SEND_DATA recv=%u/%u from=%u to=%016llx ch=%s " ,
CS_ID , inserted , count , from , ( unsigned long long ) peer , ch_id ) ;
uint32_t next = from + count ;
uint8_t snd [ 7 ] ; snd [ 0 ] = CS_MSG_SEND_DATA ;
memcpy ( snd + 1 , & next , 4 ) ; uint16_t z = 0 ; memcpy ( snd + 5 , & z , 2 ) ;
@ -234,7 +269,7 @@ static void cs_handle_send_data(struct chat_sync* cs, uint64_t peer,
static void cs_handle_push ( struct chat_sync * cs , uint64_t peer ,
const char * ch_id , const uint8_t * pl , size_t len ) {
if ( len < 24 + 1 + 4 ) return ;
if ( len < 24 + 1 + 4 ) { DEBUG_WARN ( DEBUG_CATEGORY_CONNECTIVITY , " %s: PUSH too short len=%zu peer=%016llx " , CS_ID , len , ( unsigned long long ) peer ) ; return ; }
uint64_t ts , dh ;
memcpy ( & ts , pl , 8 ) ; memcpy ( & dh , pl + 8 , 8 ) ;
@ -261,7 +296,7 @@ static void cs_handle_push(struct chat_sync* cs, uint64_t peer,
static void cs_handle_ack_push ( struct chat_sync * cs , uint64_t peer ,
const char * ch_id , const uint8_t * pl , size_t len ) {
( void ) peer ;
if ( len < 16 ) return ;
if ( len < 16 ) { DEBUG_WARN ( DEBUG_CATEGORY_CONNECTIVITY , " %s: ACK_PUSH too short len=%zu peer=%016llx " , CS_ID , len , ( unsigned long long ) peer ) ; return ; }
uint64_t ts , dh ; memcpy ( & dh , pl , 8 ) ; memcpy ( & ts , pl + 8 , 8 ) ;
chat_core_mark_sent ( ch_id , ts , dh , cs - > inst - > node_id ) ;
}
@ -271,7 +306,7 @@ static void cs_handle_ack_push(struct chat_sync* cs, uint64_t peer,
static void cs_handle_sync_done ( struct chat_sync * cs , uint64_t peer ,
const char * ch_id , const uint8_t * pl , size_t len ) {
( void ) peer ;
if ( len < 4 ) return ;
if ( len < 4 ) { DEBUG_WARN ( DEBUG_CATEGORY_CONNECTIVITY , " %s: SYNC_DONE too short len=%zu peer=%016llx " , CS_ID , len , ( unsigned long long ) peer ) ; return ; }
uint32_t pc ; memcpy ( & pc , pl , 4 ) ;
struct channel_cache * ch = cs_find ( cs , ch_id ) ;
if ( ch ) { ch - > msg_count = pc ; ch - > synced = CS_SYNC_DONE ; }
@ -389,6 +424,8 @@ static void cs_on_conn_up(struct ETCP_CONN* conn, void* arg) {
if ( peer = = 0 | | peer = = g_cs - > inst - > node_id ) return ;
if ( g_cs - > pending_invite_ch_id ! = 0 & & ! g_cs - > info_req_timer ) {
DEBUG_INFO ( DEBUG_CATEGORY_CONNECTIVITY , " %s: conn_up invite path peer=%016llx ch=%llu " ,
CS_ID , ( unsigned long long ) peer , ( unsigned long long ) g_cs - > pending_invite_ch_id ) ;
g_cs - > pending_invite_node_id = peer ;
char ch_id_str [ 64 ] ;
snprintf ( ch_id_str , sizeof ( ch_id_str ) , " %llu " ,
@ -749,8 +786,8 @@ static void cs_handle_channel_info_req(struct chat_sync* cs, uint64_t peer,
static void cs_handle_channel_info_resp ( struct chat_sync * cs , uint64_t peer ,
const char * ch_id , const uint8_t * pl , size_t len ) {
if ( cs - > info_req_timer ) { uasync_cancel_timeout ( cs - > inst - > ua , cs - > info_req_timer ) ; cs - > info_req_timer = NULL ; }
if ( len < 1 ) return ;
uint8_t nl = pl [ 0 ] ; if ( 1 + nl + 8 + 1 + 32 + 32 + 64 + 64 > len ) return ;
if ( len < 1 ) { DEBUG_WARN ( DEBUG_CATEGORY_CONNECTIVITY , " %s: CHANNEL_INFO_RESP too short len=%zu peer=%016llx " , CS_ID , len , ( unsigned long long ) peer ) ; return ; }
uint8_t nl = pl [ 0 ] ; if ( 1 + nl + 8 + 1 + 32 + 32 + 64 + 64 > len ) { DEBUG_WARN ( DEBUG_CATEGORY_CONNECTIVITY , " %s: CHANNEL_INFO_RESP truncated len=%zu need=%d " , CS_ID , len , 1 + nl + 8 + 1 + 32 + 32 + 64 + 64 ) ; return ; }
const uint8_t * p = pl + 1 ;
char name [ 128 ] ; memcpy ( name , p , nl ) ; name [ nl ] = ' \0 ' ; p + = nl ;
uint64_t owner ; memcpy ( & owner , p , 8 ) ; p + = 8 ;
@ -867,7 +904,7 @@ static void cs_handle_channel_info_resp(struct chat_sync* cs, uint64_t peer,
static void cs_handle_channel_join ( struct chat_sync * cs , uint64_t peer ,
const char * ch_id , const uint8_t * pl , size_t len ) {
if ( len < 8 + 32 + 32 + 64 + 1 ) return ;
if ( len < 8 + 32 + 32 + 64 + 1 ) { DEBUG_WARN ( DEBUG_CATEGORY_CONNECTIVITY , " %s: CHANNEL_JOIN too short len=%zu peer=%016llx " , CS_ID , len , ( unsigned long long ) peer ) ; return ; }
const uint8_t * p = pl ;
uint64_t node_id ; memcpy ( & node_id , p , 8 ) ; p + = 8 ;
const uint8_t * x25519 = p ; p + = 32 ;
@ -963,7 +1000,7 @@ static void cs_handle_welcome(struct chat_sync* cs, uint64_t peer,
const char * ch_id , const uint8_t * pl , size_t len ) {
if ( cs - > join_timer ) { uasync_cancel_timeout ( cs - > inst - > ua , cs - > join_timer ) ; cs - > join_timer = NULL ; }
( void ) peer ;
if ( len < 2 ) return ;
if ( len < 2 ) { DEBUG_WARN ( DEBUG_CATEGORY_CONNECTIVITY , " %s: WELCOME too short len=%zu peer=%016llx " , CS_ID , len , ( unsigned long long ) peer ) ; return ; }
sqlite3 * db = cs - > inst - > topo_groups - > topo_sqlite_db ;
const uint8_t * p = pl ;
@ -1038,7 +1075,7 @@ static void cs_handle_welcome(struct chat_sync* cs, uint64_t peer,
static void cs_handle_peer_upsert ( struct chat_sync * cs , uint64_t peer ,
const char * ch_id , const uint8_t * pl , size_t len ) {
if ( len < 8 + 32 + 32 + 64 + 1 ) return ;
if ( len < 8 + 32 + 32 + 64 + 1 ) { DEBUG_WARN ( DEBUG_CATEGORY_CONNECTIVITY , " %s: PEER_UPSERT too short len=%zu " , CS_ID , len ) ; return ; }
const uint8_t * p = pl ;
uint64_t node_id ; memcpy ( & node_id , p , 8 ) ; p + = 8 ;
const uint8_t * x25519 = p ; p + = 32 ;
@ -1099,7 +1136,7 @@ static void cs_handle_peer_upsert(struct chat_sync* cs, uint64_t peer,
static void cs_handle_peer_remove ( struct chat_sync * cs , uint64_t peer ,
const char * ch_id , const uint8_t * pl , size_t len ) {
if ( len < 8 ) return ;
if ( len < 8 ) { DEBUG_WARN ( DEBUG_CATEGORY_CONNECTIVITY , " %s: PEER_REMOVE too short len=%zu " , CS_ID , len ) ; return ; }
uint64_t node_id ; memcpy ( & node_id , pl , 8 ) ;
sqlite3 * db = cs - > inst - > topo_groups - > topo_sqlite_db ;