Browse Source

db_sync: add INFO diagnostics for message sync flow tracing

- db_sync_on_conn_up: per-instance synced/skipped counters
- db_sync_peer_check_cb: summary INFO of check results
- db_sync_instance_add: peers_found + peers_synced counts
- db_handle_init_sync/initiate_sync: enhanced with tbl name
- db_handle_init_resp: tp + my_mc in sync complete/dh match
- db_handle_send_data: range + SYNC_DONE details; stop reason when mc>pk
- db_handle_sync_done: both counts + datahashes in mismatch; dh in confirm
- db_sync_recv_cb: INFO when instance NOT FOUND (timing race)
- cs_handle_channel_info_resp: sync readiness after channel ready
- cs_handle_welcome: note about db_sync peer_check timer
- on_db_sync_insert: counter (#1-3 then every 10) with ch/n/ts/dh/ct
topo_upd
Evgeny 3 months ago
parent
commit
2792586de7
  1. 92
      src/db_sync.c
  2. 6
      tools/chatgui/transport/chat_core.c
  3. 4
      tools/chatgui/transport/chat_sync.c

92
src/db_sync.c

@ -679,6 +679,9 @@ static void db_handle_init_sync(struct DB_SYNC_INSTANCE* si,
}
uint32_t pc = *(uint32_t*)p;
uint32_t mc = db_count(si);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"INIT_SYNC from %016llx peer_count=%u my_count=%u tbl=%s",
(unsigned long long)src, pc, mc, SI_TBL(si));
uint32_t tp = (pc < mc ? pc : mc);
if (tp > 0)
tp--;
@ -737,21 +740,23 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si,
if (sc == 0 && my_dh == peer_dh) {
struct SI_PEER* sp = si_peer_find(si, src);
uint32_t mc = db_count(si);
if (sp) {
sp->synced_pos = tp;
sp->sync_state = 2;
}
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"sync complete with %016llx",
(unsigned long long)src);
"sync complete with %016llx tp=%u my_mc=%u",
(unsigned long long)src, tp, mc);
return;
}
if (peer_dh == 0 && sc == 0) {
uint32_t mc = db_count(si);
uint32_t batches = (mc + DB_SEND_DATA_MAX - 1) / DB_SEND_DATA_MAX;
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"peer %016llx empty, sending all %u",
(unsigned long long)src, mc);
"peer %016llx empty, sending all %u records in %u batches",
(unsigned long long)src, mc, batches);
uint32_t sent = 0;
while (sent < mc) {
uint32_t b = mc - sent;
@ -814,14 +819,15 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si,
}
if (my_dh == peer_dh) {
uint32_t mc = db_count(si);
struct SI_PEER* sp = si_peer_find(si, src);
if (sp) {
sp->synced_pos = tp;
sp->sync_state = 2;
}
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"dh match with %016llx at %u",
(unsigned long long)src, tp);
"dh match with %016llx at tp=%u my_mc=%u",
(unsigned long long)src, tp, mc);
return;
}
@ -1065,8 +1071,16 @@ static void db_handle_send_data(struct DB_SYNC_INSTANCE* si,
uint32_t scnt = mc - pk;
if (scnt > DB_SEND_DATA_MAX)
scnt = DB_SEND_DATA_MAX;
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"send_data continue to %016llx from=%u cnt=%u mc=%u",
(unsigned long long)src, pk, scnt, mc);
si_send_data_batch(si, src, pk, scnt, 1);
}
else if (mc > pk) {
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"send_data stopping to %016llx — mc=%u pk=%u state=%d",
(unsigned long long)src, mc, pk, sp ? sp->sync_state : -1);
}
uint32_t nc = db_count(si);
uint64_t ldh = 0;
@ -1080,8 +1094,9 @@ static void db_handle_send_data(struct DB_SYNC_INSTANCE* si,
if (sp)
sp->sync_state = 2;
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"received %u records from %016llx",
received, (unsigned long long)src);
"received %u/%u records from %016llx range=[%u,%u] → SYNC_DONE nc=%u dh=%016llx",
received, count, (unsigned long long)src, from, from + received - 1,
nc, (unsigned long long)ldh);
}
static void db_handle_sync_done(struct DB_SYNC_INSTANCE* si,
@ -1098,9 +1113,10 @@ static void db_handle_sync_done(struct DB_SYNC_INSTANCE* si,
db_datahash_at(si, mc - 1, &mdh);
if (mc != pc || mdh != pdh) {
DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC,
"SYNC_DONE mismatch %016llx my=%u peer=%u",
(unsigned long long)src, mc, pc);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"SYNC_DONE mismatch with %016llx my_count=%u peer_count=%u my_dh=%016llx peer_dh=%016llx — re-initiating",
(unsigned long long)src, mc, pc,
(unsigned long long)mdh, (unsigned long long)pdh);
db_sync_initiate_sync(si, src);
return;
}
@ -1110,8 +1126,8 @@ static void db_handle_sync_done(struct DB_SYNC_INSTANCE* si,
sp->sync_state = 2;
}
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"sync confirmed with %016llx count=%u",
(unsigned long long)src, mc);
"sync confirmed with %016llx count=%u dh=%016llx",
(unsigned long long)src, mc, (unsigned long long)mdh);
}
static void db_handle_ack_push(struct DB_SYNC_INSTANCE* si,
@ -1214,6 +1230,9 @@ static void db_sync_recv_cb(struct ETCP_CONN* conn,
struct DB_SYNC_INSTANCE* si = db_instance_find(db, hash);
if (!si) {
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"recv msg type=0x%02x from %016llx hash=%016llx — instance NOT FOUND (not registered yet?); sending DB_ERR_NOT_FOUND",
type, (unsigned long long)src, (unsigned long long)hash);
uint8_t err[2];
err[0] = DB_MSG_ERROR;
err[1] = DB_ERR_NOT_FOUND;
@ -1293,19 +1312,22 @@ static void db_sync_on_conn_up(struct ETCP_CONN* conn, void* arg)
return;
db->last_connected_tb = get_time_tb();
int synced = 0, skipped_state = 0, skipped_not_ready = 0;
for (int i = 0; i < db->instance_count; i++) {
struct DB_SYNC_INSTANCE* si = &db->instances[i];
if (!si->enabled)
continue;
if (!si->enabled) continue;
struct SI_PEER* p = si_peer_add(si, pid);
if (p && p->sync_state == 0 && conn->initialized && conn->links_up) {
p->sync_state = 1;
db_sync_initiate_sync(si, pid);
}
if (!p) continue;
if (p->sync_state != 0) { skipped_state++; continue; }
if (!conn->initialized || !conn->links_up) { skipped_not_ready++; continue; }
p->sync_state = 1;
db_sync_initiate_sync(si, pid);
synced++;
}
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"peer up %016llx init=%d links=%d", (unsigned long long)pid,
conn->initialized, conn->links_up);
"conn_up peer=%016llx init=%d links=%d instances=%d synced=%d skipped_state=%d skipped_not_ready=%d",
(unsigned long long)pid, conn->initialized, conn->links_up,
db->instance_count, synced, skipped_state, skipped_not_ready);
}
static void db_sync_on_conn_down(struct ETCP_CONN* conn, void* arg)
@ -1351,8 +1373,8 @@ static void db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si,
{
uint32_t mc = db_count(si);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"initiate_sync to %016llx my=%u",
(unsigned long long)pid, mc);
"INIT_SYNC → %016llx my=%u tbl=%s",
(unsigned long long)pid, mc, SI_TBL(si));
uint8_t msg[5];
msg[0] = DB_MSG_INIT_SYNC;
@ -1426,11 +1448,13 @@ static void db_sync_peer_check_cb(void* arg)
return;
}
int total_synced = 0, total_skipped = 0;
for (int i = 0; i < db->instance_count; i++) {
struct DB_SYNC_INSTANCE* si = &db->instances[i];
if (!si->enabled)
continue;
int peers_found = 0;
struct ll_entry* e = g->senders_list->head;
while (e) {
struct TOPO_GROUP_CONN_ITEM* item =
@ -1439,8 +1463,10 @@ static void db_sync_peer_check_cb(void* arg)
&& item->conn->links_up && item->conn->initialized)
{
uint64_t pid = item->conn->peer_node_id;
if (pid != db->inst->node_id)
if (pid != db->inst->node_id) {
si_peer_add(si, pid);
peers_found++;
}
}
e = e->next;
}
@ -1456,9 +1482,16 @@ static void db_sync_peer_check_cb(void* arg)
if (best) {
best->sync_state = 1;
db_sync_initiate_sync(si, best->node_id);
total_synced++;
} else {
total_skipped++;
}
}
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"peer_check: instances=%d synced=%d skipped=%d",
db->instance_count, total_synced, total_skipped);
db->peer_check_timer =
uasync_set_timeout(db->inst->ua,
DB_SYNC_PEER_CHECK_INTERVAL * 10000u,
@ -1724,7 +1757,7 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst,
si, db_sync_instance_ttl_cb,
"db_sync_ttl");
// Initiate sync with already connected peers
int peers_found = 0, peers_synced = 0;
{
struct ll_entry* entry = inst->connections->head;
while (entry) {
@ -1733,10 +1766,12 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst,
if (pid != 0 && pid != inst->node_id
&& ce->conn->links_up > 0 && ce->conn->initialized)
{
peers_found++;
struct SI_PEER* p = si_peer_add(si, pid);
if (p && p->sync_state == 0) {
p->sync_state = 1;
db_sync_initiate_sync(si, pid);
peers_synced++;
}
}
entry = entry->next;
@ -1744,11 +1779,10 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst,
}
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"instance added name=%s id=%llx tbl=%s"
" hash=%016llx next_id=%llu",
"instance added name=%s id=%llx tbl=%s hash=%016llx next_id=%llu peers_found=%d peers_synced=%d mc=%u",
name, (unsigned long long)id, si->table_name,
(unsigned long long)hash,
(unsigned long long)si->next_id);
(unsigned long long)hash, (unsigned long long)si->next_id,
peers_found, peers_synced, db_count(si));
return si;
}

6
tools/chatgui/transport/chat_core.c

@ -1136,6 +1136,7 @@ static void si_register(struct DB_SYNC_INSTANCE* si, const char* ch_id) {
static void on_db_sync_insert(struct DB_SYNC_INSTANCE* si, const char* data, size_t len, uint64_t author, void* arg) {
const char* ch_id = (const char*)arg;
static int insert_count = 0;
/* parse JSON: {"n":<uint64>,"ch":"<str>","ct":"<str>","d":"<str>"} */
const char* p = data; const char* end = data + len; uint64_t jn = 0;
{ /* "n" */
@ -1173,6 +1174,11 @@ static void on_db_sync_insert(struct DB_SYNC_INSTANCE* si, const char* data, siz
sqlite3_bind_int(st,9,1); sqlite3_bind_int(st,10,0);
int rc=sqlite3_step(st); sqlite3_finalize(st);
if (rc==SQLITE_DONE) {
insert_count++;
if (insert_count <= 3 || insert_count % 10 == 0)
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: on_db_sync_insert #%d ch=%s n=%llu ts=%llu dh=%016llx ct=%.*s",
CC_ID, insert_count, ch_id, (unsigned long long)jn, (unsigned long long)ts,
(unsigned long long)dh, (int)jct_len, jct);
uint8_t evt[65]; uint8_t cl=(uint8_t)strlen(ch_id); evt[0]=cl; memcpy(evt+1,ch_id,cl);
gui_bridge_post(GUI_EVT_MSG_RECEIVED, evt, 1+cl);
} else if (rc == SQLITE_CONSTRAINT) {

4
tools/chatgui/transport/chat_sync.c

@ -889,6 +889,8 @@ static void cs_handle_channel_info_resp(struct chat_sync* cs, uint64_t peer,
topo_node_sqlite_channel_put(cs->inst->topo_groups->topo_sqlite_db,
ch_id, name, (int)is_dm, owner, x25519, NULL, ed_pub, NULL, ch_sig);
chat_core_ensure_channel_ready(ch_id);
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: channel ready for sync ch=%s, db_sync will pick up via periodic check or active conn",
CS_ID, ch_id);
/* verify inviter's join_sig AND save node info using ETCP-authenticated keys */
{ struct ETCP_CONN* inv_conn = cs_find_conn_for_node(cs->inst, peer);
@ -1212,7 +1214,7 @@ static void cs_handle_welcome(struct chat_sync* cs, uint64_t peer,
gui_bridge_post(GUI_EVT_MEMBERS_CHANGED, evt, 1 + ch_id_len);
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: WELCOME processed ch=%s peers=%d",
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: WELCOME processed ch=%s peers=%d — message sync triggered via db_sync (active conn or peer_check timer every 5s)",
CS_ID, ch_id, pc);
}

Loading…
Cancel
Save