diff --git a/src/db_sync.c b/src/db_sync.c index 431ae87f..4b16bd12 100644 --- a/src/db_sync.c +++ b/src/db_sync.c @@ -73,6 +73,7 @@ struct DB_SYNC { #define SI_DB(si) ((si)->db_sync->db) #define SI_TBL(si) ((si)->table_name) +#define SI_SHRT(si) ({ const char* _t = SI_TBL(si); const char* _u = strrchr(_t, '_'); _u ? _u + 1 : _t; }) static void si_sql(char* buf, size_t sz, struct DB_SYNC_INSTANCE* si, const char* fmt) { @@ -521,7 +522,7 @@ static int db_sync_send_hash(struct DB_SYNC* db, uint64_t node_id, uint64_t hash struct ll_entry* entry = queue_entry_new(0); if (!entry) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "queue_entry_new"); return -1; } uint8_t* buf = u_malloc(plen + 9); - if (!buf) { queue_entry_free(entry); return -1; } + if (!buf) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync_send_hash: u_malloc(%zu) failed for dst=%016llx", plen + 9, (unsigned long long)node_id); queue_entry_free(entry); return -1; } buf[0] = ETCP_RT_ID_DB_SYNC; uint64_t hb = htobe64(hash); memcpy(buf + 1, &hb, 8); @@ -535,7 +536,8 @@ static int db_sync_send_hash(struct DB_SYNC* db, uint64_t node_id, uint64_t hash if (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; if (ce->conn->links_up) conn = ce->conn; } } if (!conn) { - DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "db_sync: no direct conn to %016llx, dropping", (unsigned long long)node_id); + DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "sync: send DROP — no direct conn to %04llX (hash=%016llx, type=%02x)", + (unsigned long long)(node_id >> 16), hash, plen > 0 ? payload[0] : 0); queue_entry_free(entry); return -1; } @@ -596,10 +598,14 @@ static void db_handle_error(struct DB_SYNC_INSTANCE* si, uint64_t src, const uin 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); + "sync [%s:%04llX] ← ERROR: code=%u (%s) — peer rejected our sync", + SI_SHRT(si), (unsigned long long)(src >> 16), p[0], what); struct SI_PEER* sp = si_peer_find(si, src); - if (sp) sp->sync_state = 0; + if (sp) { + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← ERROR: reset sync_state 2→0, synced_pos stays at %u", + SI_SHRT(si), (unsigned long long)(src >> 16), sp->synced_pos); + 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) @@ -607,9 +613,15 @@ static void db_handle_init_sync(struct DB_SYNC_INSTANCE* si, uint64_t src, const if (len < 4) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "INIT_SYNC too short %zu from %016llx", len, (unsigned long long)src); return; } uint32_t pc = *(uint32_t*)p; uint32_t mc = db_count(si); - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "INIT_SYNC from=%016llx peer_cnt=%u my_cnt=%u tbl=%s → tp_offset=%u", (unsigned long long)src, pc, mc, SI_TBL(si), (pc < mc ? pc : mc)); uint32_t tp = (pc < mc ? pc : mc); if (tp > 0) tp--; + uint64_t my_ch8_at_tp = 0; + if (mc > 0) db_chain_hash8_at(si, tp, &my_ch8_at_tp); + uint32_t min_pow2 = tp; + int sparse_cnt = 0; + { uint32_t s = 0; for (; s < 16; s++) { if (min_pow2 < (1u << s)) break; sparse_cnt++; } } + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← INIT_SYNC: peer=%u my=%u → tp=%u my_ch8_at_tp=%016llX, dispatching INIT_RESP (~%d sparse)", + SI_SHRT(si), (unsigned long long)(src >> 16), pc, mc, tp, my_ch8_at_tp, sparse_cnt); uint8_t resp[4096]; uint32_t off = 0; @@ -636,7 +648,6 @@ static void db_handle_init_sync(struct DB_SYNC_INSTANCE* si, uint64_t src, const } resp[sp] = (uint8_t)scnt; db_sync_send(si, src, resp, off); - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "INIT_RESP to %016llx tp=%u sparse=%d", (unsigned long long)src, tp, scnt); } static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len) @@ -652,11 +663,11 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, uint64_t src, const if (peer_ch8 == 0 && sc == 0) { uint32_t mc = db_count(si); - DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ peer empty, sending all %u records", mc); 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 records in %u batches", - (unsigned long long)src, mc, batches); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← INIT_RESP: peer empty (0 records) — I'm the source, sending all %u records in %u batches", + SI_SHRT(si), (unsigned long long)(src >> 16), mc, batches); uint32_t sent = 0; + int batch_nr = 0; while (sent < mc) { uint32_t b = mc - sent; if (b > DB_SEND_DATA_MAX) b = DB_SEND_DATA_MAX; uint8_t sdbuf[8192]; @@ -665,6 +676,7 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, uint64_t src, const memcpy(sdbuf + off, &sent, 4); off += 4; uint16_t rc = 0; uint16_t* rcp = (uint16_t*)(sdbuf + off); off += 2; + batch_nr++; sqlite3_stmt* stmt; si_prep(si, &stmt, @@ -697,6 +709,8 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, uint64_t src, const *rcp = rc; sent += rc; db_sync_send(si, src, sdbuf, off); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → SEND DATA: %u records from pos=%u (batch %d/%u, %u/%u sent)", + SI_SHRT(si), (unsigned long long)(src >> 16), rc, sent - rc, batch_nr, batches, sent, mc); if (rc == 0) break; } return; @@ -707,9 +721,12 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, uint64_t src, const struct SI_PEER* sp = si_peer_find(si, src); uint32_t tail = tp + 1; if (mc > tail) { - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "ch8 match at tp=%u my_mc=%u — sending tail [%u,%u) to %016llx", - tp, mc, tail, mc, (unsigned long long)src); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← INIT_RESP: tp=%u hash MATCH (%016llX=%016llX) — chains agree up to %u, sending tail [%u..%u)", + SI_SHRT(si), (unsigned long long)(src >> 16), tp, my_ch8, peer_ch8, tp, tail, mc); uint32_t sent = tail; + int btch = 0; + uint32_t tail_total = mc - tail; + uint32_t tail_batches = (tail_total + DB_SEND_DATA_MAX - 1) / DB_SEND_DATA_MAX; while (sent < mc) { uint32_t b = mc - sent; if (b > DB_SEND_DATA_MAX) b = DB_SEND_DATA_MAX; uint8_t sdbuf[8192]; @@ -718,6 +735,7 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, uint64_t src, const memcpy(sdbuf + off, &sent, 4); off += 4; uint16_t rc = 0; uint16_t* rcp = (uint16_t*)(sdbuf + off); off += 2; + btch++; sqlite3_stmt* stmt; si_prep(si, &stmt, "SELECT id,timestamp,node_id,data,author_signature" @@ -748,12 +766,14 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, uint64_t src, const *rcp = rc; sent += rc; db_sync_send(si, src, sdbuf, off); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → SEND DATA: %u records from pos=%u (batch %d/%u, %u/%u sent)", + SI_SHRT(si), (unsigned long long)(src >> 16), rc, sent - rc, btch, tail_batches, sent - tail, tail_total); if (rc == 0) break; } } if (sp) { sp->synced_pos = tp; sp->sync_state = 2; } - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync complete with %016llx tp=%u my_mc=%u", - (unsigned long long)src, tp, mc); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← INIT_RESP: tp=%u hash MATCH — both have %u records, nothing to send → synced_pos=%u", + SI_SHRT(si), (unsigned long long)(src >> 16), tp, mc, tp); return; } @@ -769,7 +789,8 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, uint64_t src, const } spr += 12; } - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "divergence with %016llx range [%u,%u]", (unsigned long long)src, ds, de); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← INIT_RESP: tp=%u hashes DIFFER (my=%016llX vs peer=%016llX) — sparse check: matched up to pos=%u, diverged from pos=%u → range [%u..%u]", + SI_SHRT(si), (unsigned long long)(src >> 16), tp, my_ch8, peer_ch8, ds > 0 ? ds - 1 : 0, de, ds, de); if (de - ds <= 1) { uint8_t ref[512]; @@ -779,6 +800,8 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, uint64_t src, const memcpy(ref + roff, &de, 4); roff += 4; ref[roff++] = 0; db_sync_send(si, src, ref, roff); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → REFINE: range [%u..%u] small enough — requesting data directly", + SI_SHRT(si), (unsigned long long)(src >> 16), ds, de); } else { uint8_t ref[512]; uint32_t roff = 0; @@ -788,15 +811,19 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint32_t rng = de - ds; uint8_t cnt = rng < DB_REFINE_HASHES ? (uint8_t)rng : DB_REFINE_HASHES; ref[roff++] = cnt; + uint64_t first_id = 0, last_id = 0; int chk_cnt = 0; for (uint8_t i = 0; i < cnt; i++) { uint32_t pos = ds + (rng * i / cnt); uint64_t ch8; if (db_chain_hash8_at(si, pos, &ch8) == 0) { memcpy(ref + roff, &pos, 4); roff += 4; memcpy(ref + roff, &ch8, 8); roff += 8; + if (chk_cnt == 0) first_id = pos; last_id = pos; chk_cnt++; } } db_sync_send(si, src, ref, roff); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → REFINE: probing range [%u..%u] (%u records) with %u checkpoints [%u..%u]", + SI_SHRT(si), (unsigned long long)(src >> 16), ds, de, de - ds, chk_cnt, (uint32_t)first_id, (uint32_t)last_id); } } @@ -844,6 +871,8 @@ static void si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32 *rcp = rc; 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 (allocated_buf) u_free(buf); } @@ -856,24 +885,32 @@ static void db_handle_refine(struct DB_SYNC_INSTANCE* si, uint64_t src, const ui DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "from=%016llx range=[%u..%u] hc=%d", (unsigned long long)src, from, to, hc); if (hc == 0) { - DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ hc=0, sending data batch from %u", from); uint32_t mc = db_count(si); uint32_t scnt = (to - from + 1) < DB_SEND_DATA_MAX ? (to - from + 1) : DB_SEND_DATA_MAX; - if (from >= mc) return; + if (from >= mc) { + DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← REFINE: from=%u >= my_count=%u — peer asked for data beyond our range, ignoring", + 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); return; } uint32_t fm = to + 1; const uint8_t* hp = p + 9; + uint32_t last_match = from, first_diff = to + 1; for (uint8_t i = 0; i < hc && hp + 12 <= p + len; i++) { uint32_t pos = *(uint32_t*)hp; uint64_t pch8 = *(uint64_t*)(hp + 4); uint64_t mch8; if (db_chain_hash8_at(si, pos, &mch8) == 0 && mch8 == pch8) { if (pos >= from && pos + 1 < fm) fm = pos + 1; + if (pos > last_match) last_match = pos; } else { if (pos < fm) fm = pos; + if (pos < first_diff) first_diff = pos; } hp += 12; } @@ -881,7 +918,8 @@ static void db_handle_refine(struct DB_SYNC_INSTANCE* si, uint64_t src, const ui 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, "SEND_DATA to %016llx from=%u count=%u", (unsigned long long)src, fm, scnt); + 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); } // ---- Parse one record from SEND_DATA/PUSH wire format ---- @@ -915,24 +953,28 @@ static int si_parse_record(const uint8_t** pp, const uint8_t* end, static void db_handle_send_data(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len) { - if (len < 6) return; + if (len < 6) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← SEND_DATA: truncated len=%zu (need >=6)", SI_SHRT(si), (unsigned long long)(src >> 16), len); return; } uint32_t from = *(uint32_t*)p; uint16_t count = *(uint16_t*)(p + 4); DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "from=%016llx from=%u count=%u", (unsigned long long)src, from, count); const uint8_t* ptr = p + 6; struct SI_PEER* sp = si_peer_find(si, src); - uint32_t fix_from = sp ? sp->synced_pos : 0; - uint16_t received = 0; + uint32_t old_synced_pos = sp ? sp->synced_pos : 0; + uint32_t fix_from = old_synced_pos; + uint16_t received = 0, parse_fails = 0, duplicates = 0, bad_sigs = 0; + uint64_t first_id = 0, last_id = 0; for (uint16_t i = 0; i < count; i++) { uint64_t rid, rts, rauthor; uint32_t rdlen; const uint8_t* rdata, *rsig; int rsiglen; - if (si_parse_record(&ptr, p + len, &rid, &rts, &rauthor, &rdlen, &rdata, &rsig, &rsiglen) != 0) break; + if (si_parse_record(&ptr, p + len, &rid, &rts, &rauthor, &rdlen, &rdata, &rsig, &rsiglen) != 0) { parse_fails++; break; } int ret = db_record_insert(si, rid, rts, rauthor, (const char*)rdata, rdlen, rsig, 0); - if (ret >= 0) received++; + if (ret >= 0) { if (received == 0) first_id = rid; last_id = rid; received++; } + if (ret == 1) duplicates++; + if (ret == -2) bad_sigs++; if (ret == 0 && si->on_insert) si->on_insert(si, rts, (const char*)rdata, rdlen, src, si->on_insert_arg); } @@ -944,17 +986,40 @@ static void db_handle_send_data(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint32_t mc = db_count(si); uint32_t pk = from + received; + if (received > 0) { + uint16_t dropped = count - received - parse_fails; + const char* extra = ""; + char extra_buf[128] = ""; + if (duplicates || bad_sigs || parse_fails || dropped) { + int pos = 0; + if (duplicates) pos += snprintf(extra_buf + pos, sizeof(extra_buf) - (size_t)pos, "%u dup", duplicates); + if (bad_sigs) pos += snprintf(extra_buf + pos, sizeof(extra_buf) - (size_t)pos, "%s%u bad_sig", pos ? ", " : "", bad_sigs); + if (parse_fails)pos += snprintf(extra_buf + pos, sizeof(extra_buf) - (size_t)pos, "%s%u parse_err", pos ? ", " : "", parse_fails); + if (dropped) pos += snprintf(extra_buf + pos, sizeof(extra_buf) - (size_t)pos, "%s%u unk", pos ? ", " : "", dropped); + extra = extra_buf; + } + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← DATA: from=%u count=%u → inserted %u records [id=%llu..%llu]%s%s, cascade from %u, synced_pos %u→%u", + SI_SHRT(si), (unsigned long long)(src >> 16), from, count, received, + (unsigned long long)first_id, (unsigned long long)last_id, + extra[0] ? " (" : "", extra, + fix_from, old_synced_pos, sp ? sp->synced_pos : 0); + } else if (count > 0) { + DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← DATA: from=%u count=%u → inserted 0 records (%u dup, %u bad_sig, %u parse_err) — everything rejected", + SI_SHRT(si), (unsigned long long)(src >> 16), from, count, duplicates, bad_sigs, parse_fails); + } + if (mc > pk && sp && sp->sync_state == 1) { uint32_t scnt = mc - pk; if (scnt > DB_SEND_DATA_MAX) scnt = DB_SEND_DATA_MAX; - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync SEND_DATA batch: %u records [%u..%u/%u] → still %u remaining", - received, from, from + received - 1, mc, mc - pk); 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); } else if (mc > pk) { - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync SEND_DATA batch: %u records [%u..%u/%u] — stopping (state=%d)", - received, from, from + received - 1, mc, sp ? sp->sync_state : -1); + 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, + sp ? sp->sync_state : -1); } else { - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync SEND_DATA batch: %u records [%u..%u/%u] — all caught up (state=%d)", - received, from, from + received - 1, mc, sp ? sp->sync_state : -1); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← DATA: all caught up at %u records — verifying final state", + SI_SHRT(si), (unsigned long long)(src >> 16), mc); } uint32_t nc = db_count(si); @@ -966,14 +1031,13 @@ static void db_handle_send_data(struct DB_SYNC_INSTANCE* si, uint64_t src, const memcpy(sd + 5, &lch8, 8); db_sync_send(si, src, sd, 13); if (sp) sp->sync_state = 2; - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, - "received %u/%u records from %016llx range=[%u,%u] → SYNC_DONE nc=%u ch8=%016llx", - received, count, (unsigned long long)src, from, from + received - 1, nc, (unsigned long long)lch8); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → SYNC_DONE: count=%u last_ch8=%016llX — confirming sync complete", + SI_SHRT(si), (unsigned long long)(src >> 16), nc, lch8); } static void db_handle_sync_done(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len) { - if (len < 12) return; + if (len < 12) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← SYNC_DONE: truncated len=%zu (need >=12)", SI_SHRT(si), (unsigned long long)(src >> 16), len); return; } uint32_t pc = *(uint32_t*)p; uint64_t pch8; memcpy(&pch8, p + 4, 8); uint32_t mc = db_count(si); @@ -985,30 +1049,33 @@ static void db_handle_sync_done(struct DB_SYNC_INSTANCE* si, uint64_t src, const if (mc != pc || mch8 != pch8) { struct SI_PEER* sp = si_peer_find(si, src); + uint32_t old_spos = sp ? sp->synced_pos : 0; + uint8_t old_retry = sp ? sp->sync_retry_count : 0; if (sp) sp->sync_retry_count++; if (sp && sp->sync_retry_count > 3) { DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, - "SYNC_DONE mismatch retry limit exceeded (%u/3) with %016llx — giving up", - sp->sync_retry_count, (unsigned long long)src); + "sync [%s:%04llX] ← SYNC_DONE: mismatch retry limit (%u/3) — giving up, accepting partial sync synced_pos %u→%u", + SI_SHRT(si), (unsigned long long)(src >> 16), sp->sync_retry_count, old_spos, (mc > 0) ? mc - 1 : 0); sp->sync_state = 2; sp->synced_pos = (mc > 0) ? mc - 1 : 0; return; } DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, - "SYNC_DONE mismatch with %016llx my_count=%u peer_count=%u my_ch8=%016llx peer_ch8=%016llx — re-initiating (retry %u/3)", - (unsigned long long)src, mc, pc, (unsigned long long)mch8, (unsigned long long)pch8, - sp ? sp->sync_retry_count : 0); + "sync [%s:%04llX] ← SYNC_DONE: MISMATCH my=%u/%016llX vs peer=%u/%016llX, synced_pos %u stays — re-initiating (retry %u→%u/3)", + SI_SHRT(si), (unsigned long long)(src >> 16), mc, mch8, pc, pch8, old_spos, + old_retry, sp ? sp->sync_retry_count : 0); db_sync_initiate_sync(si, src); return; } struct SI_PEER* sp = si_peer_find(si, src); + uint32_t old_spos = sp ? sp->synced_pos : 0; if (sp) { sp->synced_pos = (mc > 0) ? mc - 1 : 0; sp->sync_state = 2; sp->sync_retry_count = 0; } - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync confirmed with %016llx count=%u ch8=%016llx", - (unsigned long long)src, mc, (unsigned long long)mch8); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← SYNC_DONE: counts match (%u=%u), hashes match (%016llX=%016llX) ✓ synced_pos %u→%u", + SI_SHRT(si), (unsigned long long)(src >> 16), mc, pc, mch8, pch8, old_spos, mc > 0 ? mc - 1 : 0); } static void db_handle_ack_push(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len) { - if (len < 16) return; + if (len < 16) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← PUSH ACK: truncated len=%zu (need >=16)", SI_SHRT(si), (unsigned long long)(src >> 16), len); return; } uint64_t ts = *(uint64_t*)p; uint64_t author = *(uint64_t*)(p + 8); DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "from=%016llx ts=%llu author=%016llx", (unsigned long long)src, (unsigned long long)ts, (unsigned long long)author); @@ -1026,15 +1093,15 @@ static void db_handle_ack_push(struct DB_SYNC_INSTANCE* si, uint64_t src, const int changed = sqlite3_changes(SI_DB(si)); sqlite3_finalize(stmt); if (changed > 0) - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "ACK_PUSH mark sent author=%016llx ts=%llu from %016llx", - (unsigned long long)author, (unsigned long long)ts, (unsigned long long)src); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← PUSH ACK: ts=%llu — delivery confirmed by peer", + SI_SHRT(si), (unsigned long long)(src >> 16), (unsigned long long)ts); } si_delivery_update(si, ts, author, src); } static void db_handle_push(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len) { - if (len < 30) return; + if (len < 30) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← PUSH: truncated len=%zu (need >=30)", SI_SHRT(si), (unsigned long long)(src >> 16), len); return; } const uint8_t* ptr = p; uint64_t rid, rts, rauthor; uint32_t rdlen; @@ -1046,8 +1113,9 @@ 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); + 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; + if (si->peers[j].synced_pos >= ins_pos) { si->peers[j].synced_pos = ins_pos > 0 ? ins_pos - 1 : 0; adj_count++; } } if (si->on_insert) si->on_insert(si, rts, (const char*)rdata, rdlen, src, si->on_insert_arg); uint8_t ack[17]; @@ -1055,8 +1123,13 @@ 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, "PUSH inserted from %016llx id=%llu author=%016llx ts=%llu", - (unsigned long long)src, (unsigned long long)rid, (unsigned long long)rauthor, (unsigned long long)rts); + 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", + 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) { + 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); + } } } @@ -1124,7 +1197,8 @@ static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) 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(si, src, payload, plen); break; default: - DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "unknown msg type 0x%02x from %016llx", type, (unsigned long long)src); + DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "sync: unknown msg type 0x%02x from %04llX (hash=%016llx) — dropped", + type, (unsigned long long)(src >> 16), hash); break; } queue_dgram_free(entry); @@ -1165,24 +1239,24 @@ static void db_sync_on_conn_up(struct ETCP_CONN* conn, void* arg) DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "peer=%016llx init=%d links=%d", (unsigned long long)pid, conn->initialized, conn->links_up); db->last_connected_tb = get_time_tb(); - int synced = 0, skipped_state = 0, skipped_not_ready = 0, skipped_no_edkey = 0; + int synced = 0; + char tbl_list[256] = ""; for (int i = 0; i < db->instance_count; i++) { struct DB_SYNC_INSTANCE* si = &db->instances[i]; if (!si->enabled) continue; struct SI_PEER* p = si_peer_add(si, pid); if (!p) continue; - if (p->sync_state != 0) { skipped_state++; continue; } - if (!conn->initialized || !conn->links_up) { skipped_not_ready++; continue; } - { uint8_t ek[32]; if (db_get_ed25519_pubkey(db, pid, ek) != 0) { skipped_no_edkey++; continue; } } + if (p->sync_state != 0) continue; + if (!conn->initialized || !conn->links_up) continue; + { uint8_t ek[32]; if (db_get_ed25519_pubkey(db, pid, ek) != 0) continue; } p->sync_state = 1; p->sync_retry_count = 0; db_sync_initiate_sync(si, pid); synced++; + if (tbl_list[0]) { size_t tl = strlen(tbl_list); snprintf(tbl_list + tl, sizeof(tbl_list) - tl, ",%s", SI_SHRT(si)); } + else snprintf(tbl_list, sizeof(tbl_list), "%s", SI_SHRT(si)); } - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, - "conn_up peer=%016llx init=%d links=%d instances=%d synced=%d skipped_state=%d skipped_not_ready=%d skipped_no_edkey=%d tbls=[%s]", - (unsigned long long)pid, conn->initialized, conn->links_up, - db->instance_count, synced, skipped_state, skipped_not_ready, skipped_no_edkey, - (synced > 0 && db->instance_count > 0) ? db->instances[0].table_name : "none"); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync: CONN UP peer=%04llX → %d tables: [%s] — initiating sync for %d", + (unsigned long long)(pid >> 16), db->instance_count, synced > 0 ? tbl_list : "none", synced); } static void db_sync_on_conn_down(struct ETCP_CONN* conn, void* arg) @@ -1208,7 +1282,8 @@ static void db_sync_on_conn_down(struct ETCP_CONN* conn, void* arg) } cd_done: if (!any) db->last_connected_tb = 0; - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "peer down %016llx", (unsigned long long)pid); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync: CONN DOWN peer=%04llX → reset sync state for %d tables", + (unsigned long long)(pid >> 16), db->instance_count); } // ============================================================ @@ -1226,7 +1301,8 @@ static void db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si, uint64_t pid) msg[0] = DB_MSG_INIT_SYNC; memcpy(msg + 1, &mc, 4); db_sync_send(si, pid, msg, 5); - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "INIT_SYNC → %016llx my=%u tbl=%s", (unsigned long long)pid, mc, SI_TBL(si)); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → sent INIT_SYNC: my_count=%u — \"let's compare chains\"", + SI_SHRT(si), (unsigned long long)(pid >> 16), mc); } static void db_verify_chain(struct DB_SYNC_INSTANCE* si) @@ -1286,6 +1362,7 @@ static void db_sync_peer_check_cb(void* arg) } int total_synced = 0, total_skipped = 0, any_peers = 0; + char launched_list[256] = ""; for (int i = 0; i < db->instance_count; i++) { struct DB_SYNC_INSTANCE* si = &db->instances[i]; if (!si->enabled) continue; @@ -1313,6 +1390,10 @@ static void db_sync_peer_check_cb(void* arg) uint8_t ek[32]; if (db_get_ed25519_pubkey(db, best->node_id, ek) == 0) { best->sync_state = 1; db_sync_initiate_sync(si, best->node_id); total_synced++; + { size_t tl = strlen(launched_list); + snprintf(launched_list + tl, sizeof(launched_list) - tl, + "%s%s:%04llX[sp=%u]", tl ? "," : "", + SI_SHRT(si), (unsigned long long)(best->node_id >> 16), min_pos); } } else { DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "peer_check skip peer=%016llx — no Ed25519 pubkey yet", (unsigned long long)best->node_id); best = NULL; total_skipped++; @@ -1331,18 +1412,18 @@ static void db_sync_peer_check_cb(void* arg) 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)); + "sync [%s:%04llX] TIMEOUT: no response for %llu ms, sync_state 1→0 (synced_pos was %u)", + SI_SHRT(si), (unsigned long long)(p->node_id >> 16), + (unsigned long long)((now - p->sync_start_tb) / 10), p->synced_pos); p->sync_state = 0; } } } } - if (any_peers) { - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "peer_check: instances=%d synced=%d skipped=%d", - db->instance_count, total_synced, total_skipped); + if (total_synced > 0) { + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync: PEER CHECK → %d tables, launched sync for %d: %s", + db->instance_count, total_synced, launched_list); } db->peer_check_timer = uasync_set_timeout(db->inst->ua, DB_SYNC_PEER_CHECK_INTERVAL * 10000u, @@ -1592,10 +1673,9 @@ int db_sync_insert_signed(struct DB_SYNC_INSTANCE* si, const char* json_data, si { si_delivery_update(si, ts, author_node_id, si->peers[i].node_id); push_count++; } } } - DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ pushed to %d synced peers", push_count); - - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "insert id=%llu author=%016llx ts=%llu len=%zu", - (unsigned long long)id, (unsigned long long)author_node_id, (unsigned long long)ts, len); + if (push_count > 0) + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:ME--] → PUSH out: new msg id=%llu ts=%llu → forwarded to %d synced peers", + SI_SHRT(si), (unsigned long long)id, (unsigned long long)ts, push_count); return 0; }