Browse Source

db_sync: replace all sync logs with unified narrative format sync [TBL:NODE]

- Unified format: 'sync [SHORT_TBL:PEER_HI] <-/->  action: key_data'
- Each log line explains what the algorithm is doing and why
- Added batch progress tracking in dump-all and tail SEND_DATA loops
- Track dropped records (duplicates, bad_sigs, parse_fails) in SEND_DATA recv
- Show synced_pos transitions (old->new) in SEND_DATA and SYNC_DONE
- Show all launched syncs in PEER CHECK log
- Convert all silent returns to ERROR logs with [TBL:NODE] context
- Convert old WARN/ERROR formats to new sync [TBL:NODE] format
- Add sync_state reset details on ERROR and timeout paths
topo_upd
Evgeny 2 months ago
parent
commit
a7ec45c4af
  1. 218
      src/db_sync.c

218
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;
}

Loading…
Cancel
Save