diff --git a/lib/tcp_io.c b/lib/tcp_io.c index 2b0a5210..bea103c9 100644 --- a/lib/tcp_io.c +++ b/lib/tcp_io.c @@ -170,10 +170,6 @@ static void tcp_conn_handle_error(struct tcp_conn* tc, int err) if (tc->sock != SOCKET_INVALID) { uasync_remove_socket_t(tc->ua, tc->sock); tc->socket_id = NULL; -#if HAS_EPOLL - if (tc->ua && tc->ua->use_epoll && tc->ua->epoll_fd >= 0) - epoll_ctl(tc->ua->epoll_fd, EPOLL_CTL_DEL, (int)tc->sock, NULL); -#endif socket_close_wrapper(tc->sock); tc->sock = SOCKET_INVALID; // ДО on_error: read/write в том же epoll event увидят INVALID } diff --git a/lib/u_async.c b/lib/u_async.c index 8ceb0c72..b6dcfbff 100644 --- a/lib/u_async.c +++ b/lib/u_async.c @@ -55,7 +55,7 @@ struct socket_node { int active; // 1 if socket is active, 0 if freed (for reuse) int enable_read; // 1 if read monitoring is enabled int enable_write; // 1 if write monitoring is enabled - uint16_t gen; // generation counter for epoll event validation + uint32_t gen; // generation counter for epoll event validation }; // Array-based socket management for O(1) operations @@ -67,7 +67,7 @@ struct socket_array { int capacity; // Total allocated capacity int count; // Number of active sockets int max_fd; // Maximum FD for bounds checking - uint16_t gen_counter; // incrementing generation for epoll stale-event detection + uint32_t gen_counter; // incrementing generation for epoll stale-event detection }; static struct socket_array* socket_array_create(int initial_capacity); @@ -792,7 +792,6 @@ err_t uasync_set_socket_read(struct UASYNC* ua, void* s_id, int enable) { if (node->type == SOCKET_NODE_TYPE_SOCK) { if (node->read_cbk_sock && node->enable_read) ev.events |= EPOLLIN; if (node->write_cbk_sock && node->enable_write) ev.events |= EPOLLOUT; - ev.data.fd = node->sock; } else { if (node->read_cbk && node->enable_read) ev.events |= EPOLLIN; if (node->write_cbk && node->enable_write) ev.events |= EPOLLOUT; @@ -929,7 +928,7 @@ static void process_epoll_events(struct UASYNC* ua, struct epoll_event* events, } int fd = (int)(events[i].data.u64 & 0xFFFFFFFF); - uint16_t ev_gen = (uint16_t)(events[i].data.u64 >> 32); + uint32_t ev_gen = (uint32_t)(events[i].data.u64 >> 32); struct socket_node* node = socket_array_get(ua->sockets, fd); if (!node || !node->active) continue; diff --git a/src/db_sync.c b/src/db_sync.c index 8d4abbd4..978228c4 100644 --- a/src/db_sync.c +++ b/src/db_sync.c @@ -29,6 +29,7 @@ static void db_sync_instance_ttl_cb(void* arg); static int db_sync_send(struct DB_SYNC_INSTANCE* si, uint64_t dst_node_id, const uint8_t* payload, size_t len); static int db_sync_send_hash(struct DB_SYNC* db, uint64_t node_id, uint64_t hash, const uint8_t* payload, size_t len); static void db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si, uint64_t peer_node_id); +static int si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32_t from, uint32_t count, uint32_t vp, int allocated_buf); // ============================================================ // Structures @@ -37,8 +38,8 @@ static void db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si, uint64_t peer_nod struct SI_PEER { uint64_t node_id; uint32_t synced_pos; + uint32_t verified_pos; uint8_t sync_state; - uint8_t sync_retry_count; uint64_t sync_start_tb; }; @@ -228,10 +229,8 @@ static struct SI_PEER* si_peer_add(struct DB_SYNC_INSTANCE* si, uint64_t node_id si->peer_capacity = nc; } p = &si->peers[si->peer_count++]; + memset(p, 0, sizeof(*p)); p->node_id = node_id; - p->synced_pos = 0; - p->sync_state = 0; - p->sync_start_tb = 0; return p; } @@ -617,28 +616,29 @@ static void db_handle_init_sync(struct DB_SYNC_INSTANCE* si, uint64_t src, const 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 has_tail = (mc > tp + 1) ? 1 : 0; + static const uint16_t iv[16] = {1, 1, 1, 2, 2, 4, 4, 4, 4, 4, 4, 4, 4, 4, 4, 4}; + uint32_t accum = 0; int sparse_cnt = 0; + for (int i = 0; i < 16; i++) { if (tp < accum + iv[i]) break; accum += iv[i]; sparse_cnt++; } + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← INIT_SYNC: peer=%u my=%u → tp=%u ch8=%016llX has_tail=%u sparse=%d", + SI_SHRT(si), (unsigned long long)(src >> 16), pc, mc, tp, my_ch8_at_tp, has_tail, sparse_cnt); uint8_t resp[4096]; uint32_t off = 0; resp[off++] = DB_MSG_INIT_RESP; memcpy(resp + off, &tp, 4); off += 4; - uint64_t my_ch8 = 0; if (mc > 0) db_chain_hash8_at(si, tp, &my_ch8); memcpy(resp + off, &my_ch8, 8); off += 8; - + resp[off++] = has_tail; uint32_t sp = off; off++; int scnt = 0; - for (uint32_t k = 0; k < 16; k++) { - uint32_t step = (uint32_t)(1u << k); - if (tp < step) break; - uint32_t pos = tp - step; + accum = 0; + for (int i = 0; i < 16; i++) { + if (tp < accum + iv[i]) break; + accum += iv[i]; + uint32_t pos = tp - accum; uint64_t sch8; if (db_chain_hash8_at(si, pos, &sch8) != 0) break; if (off + 12 > sizeof(resp)) break; @@ -652,186 +652,90 @@ static void db_handle_init_sync(struct DB_SYNC_INSTANCE* si, uint64_t src, const static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len) { - if (len < 13) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "INIT_RESP too short %zu from %016llx", len, (unsigned long long)src); return; } + if (len < 14) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "INIT_RESP too short %zu from %016llx", len, (unsigned long long)src); return; } uint32_t tp = *(uint32_t*)p; uint64_t peer_ch8; memcpy(&peer_ch8, p + 4, 8); - uint8_t sc = p[12]; - DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "from=%016llx tp=%u peer_ch8=%016llx sc=%d", (unsigned long long)src, tp, (unsigned long long)peer_ch8, sc); + uint8_t has_tail = p[12]; + uint8_t sc = p[13]; + const uint8_t* spr = p + 14; uint64_t my_ch8 = 0; - uint32_t mc_init = db_count(si); - if (mc_init > 0) db_chain_hash8_at(si, tp, &my_ch8); + uint32_t mc = db_count(si); + if (mc > 0 && tp < mc) db_chain_hash8_at(si, tp, &my_ch8); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← INIT_RESP: tp=%u my=%u my_ch8=%016llX peer_ch8=%016llX has_tail=%u sc=%d", + SI_SHRT(si), (unsigned long long)(src >> 16), tp, mc, my_ch8, peer_ch8, has_tail, sc); + + struct SI_PEER* sp = si_peer_find(si, src); if (peer_ch8 == 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, "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 vp = (uint32_t)-1; + if (sp) sp->verified_pos = vp; 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]; - uint32_t off = 0; - sdbuf[off++] = DB_MSG_SEND_DATA; - 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, - "SELECT id,timestamp,node_id,data,author_signature" - " FROM \"%s\" ORDER BY timestamp, author_signature" - " LIMIT ? OFFSET ?"); - sqlite3_bind_int64(stmt, 1, (sqlite3_int64)b); - sqlite3_bind_int64(stmt, 2, (sqlite3_int64)sent); - while (sqlite3_step(stmt) == SQLITE_ROW) { - 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 uint8_t* rd = (const uint8_t*)sqlite3_column_blob(stmt, 3); - uint32_t rdl = (uint32_t)sqlite3_column_bytes(stmt, 3); if (!rd) rdl = 0; - const uint8_t* rsig = (const uint8_t*)sqlite3_column_blob(stmt, 4); - int rsl = sqlite3_column_bytes(stmt, 4); if (!rsig) rsl = 0; - // Wire: [id:8][ts:8][author:8][dlen:4][data][sig_len:1][sig] - int rec_sz = 28 + rdl + 1 + rsl; - if (off + rec_sz > (int)sizeof(sdbuf)) break; - memcpy(sdbuf + off, &rid, 8); off += 8; - memcpy(sdbuf + off, &rts, 8); off += 8; - memcpy(sdbuf + off, &rauth, 8); off += 8; - memcpy(sdbuf + off, &rdl, 4); off += 4; - if (rdl > 0) { memcpy(sdbuf + off, rd, rdl); off += rdl; } - sdbuf[off++] = (uint8_t)rsl; - if (rsl > 0) { memcpy(sdbuf + off, rsig, rsl); off += rsl; } - rc++; - } - sqlite3_finalize(stmt); - *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; + int pushed = si_send_data_batch(si, src, sent, b, vp, 1); + if (pushed <= 0) break; + sent += (uint32_t)pushed; } + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← INIT_RESP: peer empty → sent %u/%u records vp=%u", + SI_SHRT(si), (unsigned long long)(src >> 16), sent, mc, vp); return; } if (my_ch8 == peer_ch8) { - uint32_t mc = db_count(si); - struct SI_PEER* sp = si_peer_find(si, src); - uint32_t tail = tp + 1; - if (mc > tail) { - 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]; - uint32_t off = 0; - sdbuf[off++] = DB_MSG_SEND_DATA; - 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" - " FROM \"%s\" ORDER BY timestamp, author_signature" - " LIMIT ? OFFSET ?"); - sqlite3_bind_int64(stmt, 1, (sqlite3_int64)b); - sqlite3_bind_int64(stmt, 2, (sqlite3_int64)sent); - while (sqlite3_step(stmt) == SQLITE_ROW) { - 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 uint8_t* rd = (const uint8_t*)sqlite3_column_blob(stmt, 3); - uint32_t rdl = (uint32_t)sqlite3_column_bytes(stmt, 3); if (!rd) rdl = 0; - const uint8_t* rsig = (const uint8_t*)sqlite3_column_blob(stmt, 4); - int rsl = sqlite3_column_bytes(stmt, 4); if (!rsig) rsl = 0; - int rec_sz = 28 + rdl + 1 + rsl; - if (off + rec_sz > (int)sizeof(sdbuf)) break; - memcpy(sdbuf + off, &rid, 8); off += 8; - memcpy(sdbuf + off, &rts, 8); off += 8; - memcpy(sdbuf + off, &rauth, 8); off += 8; - memcpy(sdbuf + off, &rdl, 4); off += 4; - if (rdl > 0) { memcpy(sdbuf + off, rd, rdl); off += rdl; } - sdbuf[off++] = (uint8_t)rsl; - if (rsl > 0) { memcpy(sdbuf + off, rsig, rsl); off += rsl; } - rc++; - } - sqlite3_finalize(stmt); - *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; - } + uint32_t old_sp = sp ? sp->synced_pos : 0; + if (!has_tail) { + if (sp) { sp->verified_pos = tp; sp->synced_pos = tp; sp->sync_state = 2; sp->sync_start_tb = 0; } + uint8_t sd[13]; + sd[0] = DB_MSG_SYNC_DONE; memcpy(sd + 1, &mc, 4); memcpy(sd + 5, &my_ch8, 8); + db_sync_send(si, src, sd, 13); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← INIT_RESP: MATCH+no_tail ✓ synced_pos %u→%u → SYNC_DONE", + SI_SHRT(si), (unsigned long long)(src >> 16), old_sp, tp); + return; } - if (sp) { sp->synced_pos = tp; sp->sync_state = 2; } - 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); + uint32_t vp = tp; + if (sp) sp->verified_pos = vp; + uint8_t req[11]; + req[0] = DB_MSG_SEND_DATA; uint32_t rfrom = tp + 1; memcpy(req + 1, &rfrom, 4); + uint16_t z = 0; memcpy(req + 5, &z, 2); memcpy(req + 7, &vp, 4); + db_sync_send(si, src, req, 11); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← INIT_RESP: MATCH but has_tail=1 → requesting SEND_DATA from=%u vp=%u", + SI_SHRT(si), (unsigned long long)(src >> 16), rfrom, vp); return; } - uint32_t ds = 0, de = tp; - const uint8_t* spr = p + 13; + uint32_t ds = 0, de = tp, fm = tp; for (int i = 0; i < sc && spr + 12 <= p + len; i++) { - uint32_t pos = *(uint32_t*)spr; - uint64_t pch8; memcpy(&pch8, spr + 4, 8); + uint32_t pos = *(uint32_t*)spr; uint64_t pch8; memcpy(&pch8, spr + 4, 8); spr += 12; uint64_t mch8; if (db_chain_hash8_at(si, pos, &mch8) == 0) { if (mch8 == pch8) { if (pos + 1 > ds) ds = pos + 1; } - else { if (pos < de) de = pos; } - } - spr += 12; + else { if (pos < de) de = pos; if (pos < fm) fm = pos; } + } else { if (pos < fm) fm = pos; } } - 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); + uint32_t vp = fm > 0 ? fm - 1 : (uint32_t)-1; + if (sp) sp->verified_pos = vp; + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← INIT_RESP: MISMATCH tp=%u ds=%u de=%u fm=%u vp=%u", + SI_SHRT(si), (unsigned long long)(src >> 16), tp, ds, de, fm, vp); - if (de - ds <= 1) { - uint8_t ref[512]; - uint32_t roff = 0; - ref[roff++] = DB_MSG_REFINE; - memcpy(ref + roff, &ds, 4); roff += 4; - 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; - ref[roff++] = DB_MSG_REFINE; - memcpy(ref + roff, &ds, 4); roff += 4; - memcpy(ref + roff, &de, 4); roff += 4; - 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++; - } + { + uint32_t scnt = mc - fm; if (scnt > DB_SEND_DATA_MAX) scnt = DB_SEND_DATA_MAX; + if (scnt > 0) si_send_data_batch(si, src, fm, scnt, vp, 0); + else { + uint8_t req[11]; + req[0] = DB_MSG_SEND_DATA; memcpy(req + 1, &fm, 4); uint16_t z = 0; memcpy(req + 5, &z, 2); memcpy(req + 7, &vp, 4); + db_sync_send(si, src, req, 11); } - 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); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → SEND_DATA from=%u vp=%u count=%u (ds=%u de=%u fm=%u)", + SI_SHRT(si), (unsigned long long)(src >> 16), fm, vp, scnt, ds, de, fm); } } // ---- SEND_DATA batch helper ---- -// Wire format: [id:8][ts:8][node_id:8][dlen:4][data][sig_len:1=64][sig:64] +// Wire format: [from:4][count:2][vp:4][records...] +// Record: [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) +static int si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32_t from, uint32_t count, uint32_t vp, int allocated_buf) { uint8_t sbuf_stack[8192]; uint8_t* buf = allocated_buf ? u_malloc(8192) : NULL; @@ -842,6 +746,7 @@ static int si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32_ memcpy(buf + off, &from, 4); off += 4; uint16_t rc = 0; uint16_t* rcp = (uint16_t*)(buf + off); off += 2; + memcpy(buf + off, &vp, 4); off += 4; sqlite3_stmt* stmt; int prep_rc = si_prep(si, &stmt, @@ -889,69 +794,6 @@ static int si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32_ return (int)rc; } -static void db_handle_refine(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len) -{ - if (len < 9) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "REFINE too short %zu", len); return; } - uint32_t from = *(uint32_t*)p; - uint32_t to = *(uint32_t*)(p + 4); - uint8_t hc = p[8]; - DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "from=%016llx range=[%u..%u] hc=%d", (unsigned long long)src, from, to, hc); - - if (hc == 0) { - 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) { - 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; - } - 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; - } - - 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; - } - uint32_t mc = db_count(si); - uint32_t scnt = 4; - 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 ---- // Wire: [id:8][ts:8][author:8][dlen:4][data][sig_len:1][sig] @@ -982,106 +824,210 @@ static int si_parse_record(const uint8_t** pp, const uint8_t* end, return 0; } +// ---- Cascade chain_hash for a specific range, then full cascade from end of range ---- +static void db_cascade_range(struct DB_SYNC_INSTANCE* si, uint32_t range_start, uint32_t range_end) +{ + uint32_t mc = db_count(si); + DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "tbl=%s start=%u end=%u mc=%u", SI_TBL(si), range_start, range_end, mc); + if (range_start >= range_end || range_start >= mc) return; + sqlite3* db = SI_DB(si); + int rc = sqlite3_exec(db, "BEGIN IMMEDIATE", NULL, NULL, NULL); + if (rc != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "cascade_range BEGIN: %s", sqlite3_errmsg(db)); return; } + + uint8_t prev_ch[32]; memset(prev_ch, 0, 32); + if (range_start > 0) { sqlite3_stmt* s; if (si_prep(si, &s, "SELECT chain_hash FROM \"%s\" ORDER BY timestamp, author_signature LIMIT 1 OFFSET ?") == SQLITE_OK) { sqlite3_bind_int64(s, 1, (sqlite3_int64)(range_start - 1)); if (sqlite3_step(s) == SQLITE_ROW) { const void* b = sqlite3_column_blob(s, 0); if (b && sqlite3_column_bytes(s, 0) >= 32) memcpy(prev_ch, b, 32); } sqlite3_finalize(s); } } + + sqlite3_stmt* sel, *upd; + if (si_prep(si, &sel, "SELECT id,timestamp,node_id,author_signature FROM \"%s\" ORDER BY timestamp, author_signature LIMIT -1 OFFSET ?") != SQLITE_OK) { sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return; } + sqlite3_bind_int64(sel, 1, (sqlite3_int64)range_start); + if (si_prep(si, &upd, "UPDATE \"%s\" SET chain_hash=? WHERE timestamp=? AND author_signature=?") != SQLITE_OK) { sqlite3_finalize(sel); sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return; } + + uint8_t ch[32]; uint32_t count = 0; + while (sqlite3_step(sel) == SQLITE_ROW) { + sqlite3_int64 n_id = sqlite3_column_int64(sel, 0), n_ts = sqlite3_column_int64(sel, 1), n_auth = sqlite3_column_int64(sel, 2); + const void* sig_blob = sqlite3_column_blob(sel, 3); + uint8_t sig[DB_SIG_SIZE]; if (sig_blob) memcpy(sig, sig_blob, DB_SIG_SIZE); else memset(sig, 0, DB_SIG_SIZE); + db_chain_hash_compute(prev_ch, (uint64_t)n_id, (uint64_t)n_ts, (uint64_t)n_auth, sig, ch); + sqlite3_reset(upd); sqlite3_bind_blob(upd, 1, ch, 32, SQLITE_STATIC); sqlite3_bind_int64(upd, 2, n_ts); sqlite3_bind_blob(upd, 3, sig, DB_SIG_SIZE, SQLITE_STATIC); + sqlite3_step(upd); + memcpy(prev_ch, ch, 32); count++; + } + sqlite3_finalize(upd); sqlite3_finalize(sel); + rc = sqlite3_exec(db, "COMMIT", NULL, NULL, NULL); + if (rc != SQLITE_OK) DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "cascade_range COMMIT: %s", sqlite3_errmsg(db)); + DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "cascade_range [%s] from=%u → %u records recalculated", SI_TBL(si), range_start, count); +} + static void db_handle_send_data(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len) { - 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; } + if (len < 10) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← SEND_DATA: truncated len=%zu (need >=10)", 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; + uint32_t vp = *(uint32_t*)(p + 6); struct SI_PEER* sp = si_peer_find(si, src); + + // count=0: peer is requesting OUR data from position "from" + if (count == 0) { + uint32_t mc = db_count(si); + uint32_t eff_from = from; + if (vp != (uint32_t)-1 && vp + 1 > from) eff_from = vp + 1; + if (eff_from >= mc) return; + uint32_t scnt = mc - eff_from; if (scnt > DB_SEND_DATA_MAX) scnt = DB_SEND_DATA_MAX; + int sent = si_send_data_batch(si, src, eff_from, scnt, vp, 0); + uint64_t nc = db_count(si); uint64_t lch8 = 0; if (nc > 0) db_chain_hash8_at(si, nc - 1, &lch8); + uint8_t sd[13]; + sd[0] = DB_MSG_SYNC_DONE; memcpy(sd + 1, &nc, 4); memcpy(sd + 5, &lch8, 8); + db_sync_send(si, src, sd, 13); + if (sp) { sp->synced_pos = nc > 0 ? nc - 1 : 0; sp->sync_state = 2; sp->sync_start_tb = 0; } + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← SEND_DATA request from=%u (eff=%u) → sent %d records mc=%u + SYNC_DONE", + SI_SHRT(si), (unsigned long long)(src >> 16), from, eff_from, sent, mc); + return; + } + + DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "from=%016llx from=%u count=%u vp=%u", (unsigned long long)src, from, count, vp); + const uint8_t* ptr = p + 10; + 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; + uint64_t recv_ts[32], recv_sig8[32]; uint16_t recv_cnt = 0; - for (uint16_t i = 0; i < count; i++) { - uint64_t rid, rts, rauthor; - uint32_t rdlen; - const uint8_t* rdata, *rsig; - int rsiglen; + for (uint16_t i = 0; i < count && i < 32; i++) { + const uint8_t* save = ptr; + 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) { parse_fails++; break; } int ret = db_record_insert(si, rid, rts, rauthor, (const char*)rdata, rdlen, rsig, 0); + if (ret >= 0) { recv_ts[recv_cnt] = rts; if (rsig && rsiglen >= 8) memcpy(&recv_sig8[recv_cnt], rsig, 8); else recv_sig8[recv_cnt] = 0; recv_cnt++; } 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); + if (ret == 0 && si->on_insert) si->on_insert(si, rts, (const char*)rdata, rdlen, src, si->on_insert_arg); } if (received > 0) { - db_cascade_from(si, fix_from); - + uint32_t rst = (from < old_synced_pos) ? from : old_synced_pos; + db_cascade_range(si, rst, from + received); + if (sp) sp->synced_pos = from + received - 1; for (int j = 0; j < si->peer_count; j++) { - if (si->peers[j].sync_state != 1 && si->peers[j].synced_pos >= fix_from && &si->peers[j] != sp) { - uint32_t old = si->peers[j].synced_pos; - si->peers[j].synced_pos = fix_from > 0 ? fix_from - 1 : 0; - DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] cascade changed chain from pos=%u → peer=%016llx synced_pos %u→%u (state=%u)", - SI_SHRT(si), (unsigned long long)(src >> 16), fix_from, - (unsigned long long)si->peers[j].node_id, old, si->peers[j].synced_pos, si->peers[j].sync_state); + if (si->peers[j].synced_pos >= rst && &si->peers[j] != sp) + si->peers[j].synced_pos = rst > 0 ? rst - 1 : 0; + } + uint32_t nc = db_count(si); + char extra_buf[64] = ""; if (duplicates) snprintf(extra_buf, sizeof(extra_buf), " (%u dup)", duplicates); + if (bad_sigs) { size_t el = strlen(extra_buf); snprintf(extra_buf+el, sizeof(extra_buf)-el, "%s%u bad_sig", extra_buf[0]?", ":" (", bad_sigs); if (!extra_buf[0]) extra_buf[0]=' '; } + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← DATA: from=%u count=%u vp=%u → inserted %u [id=%llu..%llu]%s cascade=[%u..%u] synced=%u→%u", + SI_SHRT(si), (unsigned long long)(src >> 16), from, count, vp, received, (unsigned long long)first_id, (unsigned long long)last_id, + extra_buf, rst, from+received-1, old_synced_pos, sp?sp->synced_pos:0); + + int relay_count = 0; + uint8_t push_buf[2560]; uint32_t push_off; + sqlite3_stmt* pstmt; + if (si_prep(si, &pstmt, + "SELECT id,timestamp,node_id,data,author_signature" + " FROM \"%s\" ORDER BY timestamp, author_signature" + " LIMIT ? OFFSET ?") == SQLITE_OK) + { + sqlite3_bind_int64(pstmt, 1, (sqlite3_int64)received); + sqlite3_bind_int64(pstmt, 2, (sqlite3_int64)from); + while (sqlite3_step(pstmt) == SQLITE_ROW) { + push_off = 1; + uint64_t prid = (uint64_t)sqlite3_column_int64(pstmt, 0); + uint64_t prts = (uint64_t)sqlite3_column_int64(pstmt, 1); + uint64_t prauth = (uint64_t)sqlite3_column_int64(pstmt, 2); + const uint8_t* prd = (const uint8_t*)sqlite3_column_blob(pstmt, 3); + uint32_t prdl = (uint32_t)sqlite3_column_bytes(pstmt, 3); if (!prd) prdl = 0; + const uint8_t* prsig = (const uint8_t*)sqlite3_column_blob(pstmt, 4); + int prsl = sqlite3_column_bytes(pstmt, 4); if (!prsig) prsl = 0; + memcpy(push_buf + push_off, &prid, 8); push_off += 8; + memcpy(push_buf + push_off, &prts, 8); push_off += 8; + memcpy(push_buf + push_off, &prauth, 8); push_off += 8; + memcpy(push_buf + push_off, &prdl, 4); push_off += 4; + if (prdl > 0) { memcpy(push_buf + push_off, prd, prdl); push_off += prdl; } + push_buf[push_off++] = (uint8_t)prsl; + if (prsl > 0) { memcpy(push_buf + push_off, prsig, prsl); push_off += prsl; } + push_buf[0] = DB_MSG_PUSH; + for (int j = 0; j < si->peer_count; j++) { + if (si->peers[j].sync_state >= 1 && si->peers[j].node_id != src + && si->peers[j].node_id != si->db_sync->inst->node_id) { + if (db_sync_send(si, si->peers[j].node_id, push_buf, push_off) >= 0) { + si_delivery_update(si, prts, prauth, si->peers[j].node_id); relay_count++; + } + } + } } + sqlite3_finalize(pstmt); } + if (relay_count > 0) + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← DATA: relayed %u records to %d peers", + SI_SHRT(si), (unsigned long long)(src >> 16), received, relay_count / (int)received); + } else { + uint32_t nc = db_count(si); + if (old_synced_pos < nc) db_cascade_range(si, old_synced_pos, nc); + DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← DATA: from=%u count=%u vp=%u → 0 inserted (%u dup %u bad_sig %u parse_err) cascade_from=%u", + SI_SHRT(si), (unsigned long long)(src >> 16), from, count, vp, duplicates, bad_sigs, parse_fails, old_synced_pos); } - if (sp && count > 0) sp->synced_pos = from + count - 1; - + // Send back own records after vp EXCEPT those just received 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; + uint32_t resp_from = vp; + uint32_t resp_start = (vp == (uint32_t)-1) ? 0 : (vp + 1 > from + received ? vp + 1 : from + received); + int resp_count = 0; + if (resp_start < mc) { + uint8_t rbuf[8192]; + uint32_t roff = 0; + rbuf[roff++] = DB_MSG_SEND_DATA; + memcpy(rbuf + roff, &resp_start, 4); roff += 4; + uint16_t* rrcp = (uint16_t*)(rbuf + roff); roff += 2; + memcpy(rbuf + roff, &vp, 4); roff += 4; + sqlite3_stmt* stmt; + if (si_prep(si, &stmt, + "SELECT id,timestamp,node_id,data,author_signature" + " FROM \"%s\" ORDER BY timestamp, author_signature" + " LIMIT ? OFFSET ?") == SQLITE_OK) + { + sqlite3_bind_int64(stmt, 1, (sqlite3_int64)(mc - resp_start)); + sqlite3_bind_int64(stmt, 2, (sqlite3_int64)resp_start); + while (sqlite3_step(stmt) == SQLITE_ROW && resp_count < DB_SEND_DATA_MAX) { + uint64_t rts = (uint64_t)sqlite3_column_int64(stmt, 1); + const uint8_t* rsig = (const uint8_t*)sqlite3_column_blob(stmt, 4); + uint64_t chk_sig8 = 0; if (rsig) memcpy(&chk_sig8, rsig, 8); + int is_dup = 0; + for (int k = 0; k < recv_cnt; k++) { if (recv_ts[k] == rts && recv_sig8[k] == chk_sig8) { is_dup = 1; break; } } + if (is_dup) continue; + uint64_t rid = (uint64_t)sqlite3_column_int64(stmt, 0); + uint64_t rauth = (uint64_t)sqlite3_column_int64(stmt, 2); + const uint8_t* rd = (const uint8_t*)sqlite3_column_blob(stmt, 3); + uint32_t rdl = (uint32_t)sqlite3_column_bytes(stmt, 3); if (!rd) rdl = 0; + int rsl = sqlite3_column_bytes(stmt, 4); if (!rsig) rsl = 0; + int rec_sz = 28 + rdl + 1 + rsl; + if (roff + rec_sz > 8000) break; + memcpy(rbuf + roff, &rid, 8); roff += 8; + memcpy(rbuf + roff, &rts, 8); roff += 8; + memcpy(rbuf + roff, &rauth, 8); roff += 8; + memcpy(rbuf + roff, &rdl, 4); roff += 4; + if (rdl > 0) { memcpy(rbuf + roff, rd, rdl); roff += rdl; } + rbuf[roff++] = (uint8_t)rsl; + if (rsl > 0) { memcpy(rbuf + roff, rsig, rsl); roff += rsl; } + resp_count++; + } + sqlite3_finalize(stmt); } - 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; - 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; + *rrcp = resp_count; + if (resp_count > 0) { + db_sync_send(si, src, rbuf, roff); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → RESP DATA: %u own records after vp=%u (from=%u)", + SI_SHRT(si), (unsigned long long)(src >> 16), resp_count, vp, resp_start); } - 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, - sp ? sp->sync_state : -1); - } else { - 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); - uint64_t lch8 = 0; - if (nc > 0) db_chain_hash8_at(si, nc - 1, &lch8); + uint64_t lch8 = 0; if (mc > 0) db_chain_hash8_at(si, mc - 1, &lch8); uint8_t sd[13]; - sd[0] = DB_MSG_SYNC_DONE; - memcpy(sd + 1, &nc, 4); - memcpy(sd + 5, &lch8, 8); + sd[0] = DB_MSG_SYNC_DONE; memcpy(sd + 1, &mc, 4); memcpy(sd + 5, &lch8, 8); db_sync_send(si, src, sd, 13); - if (sp) sp->sync_state = 2; - 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); + if (sp) { sp->synced_pos = mc > 0 ? mc - 1 : 0; sp->sync_state = 2; sp->sync_start_tb = 0; } + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → SYNC_DONE: fc=%u ch8=%016llX synced=%u⇥2", + SI_SHRT(si), (unsigned long long)(src >> 16), mc, lch8, sp ? sp->synced_pos : 0); } static void db_handle_sync_done(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len) @@ -1092,34 +1038,46 @@ static void db_handle_sync_done(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint32_t mc = db_count(si); uint64_t mch8 = 0; if (mc > 0) db_chain_hash8_at(si, mc - 1, &mch8); - DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "from=%016llx my_cnt=%u peer_cnt=%u my_ch8=%016llx peer_ch8=%016llx → %s", - (unsigned long long)src, mc, pc, (unsigned long long)mch8, (unsigned long long)pch8, - (mc == pc && mch8 == pch8) ? "matched" : "MISMATCH"); - - 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 [%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; - } + + struct SI_PEER* sp = si_peer_find(si, src); + uint32_t old_spos = sp ? sp->synced_pos : 0; + + if (mc == pc && mch8 == pch8) { + if (sp) { sp->synced_pos = mc > 0 ? mc - 1 : 0; sp->sync_state = 2; } DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, - "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); + "sync [%s:%04llX] ← SYNC_DONE: matched (%u=%u, %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); 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 [%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); + + if (mc < pc) { + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, + "sync [%s:%04llX] ← SYNC_DONE: count diff my=%u peer=%u — requesting tail from=%u", + SI_SHRT(si), (unsigned long long)(src >> 16), mc, pc, mc); + uint32_t vp_send = mc > 0 ? mc - 1 : (uint32_t)-1; + uint8_t req[11]; + req[0] = DB_MSG_SEND_DATA; memcpy(req + 1, &mc, 4); uint16_t z = 0; memcpy(req + 5, &z, 2); memcpy(req + 7, &vp_send, 4); + db_sync_send(si, src, req, 11); + if (sp) sp->sync_state = 1; + return; + } + + if (mc > pc) { + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, + "sync [%s:%04llX] ← SYNC_DONE: count diff my=%u peer=%u — sending tail from=%u", + SI_SHRT(si), (unsigned long long)(src >> 16), mc, pc, pc); + uint32_t vp_send = pc > 0 ? pc - 1 : (uint32_t)-1; + uint32_t scnt = mc - pc; if (scnt > DB_SEND_DATA_MAX) scnt = DB_SEND_DATA_MAX; + si_send_data_batch(si, src, pc, scnt, vp_send, 0); + if (sp) sp->sync_state = 1; + return; + } + + DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, + "sync [%s:%04llX] ← SYNC_DONE: HASH MISMATCH my=%u/%016llX vs peer=%u/%016llX (same count) — re-initiating sync", + SI_SHRT(si), (unsigned long long)(src >> 16), mc, mch8, pc, pch8); + if (sp) { sp->sync_state = 1; sp->sync_start_tb = get_time_tb(); } + db_sync_initiate_sync(si, src); } static void db_handle_ack_push(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len) @@ -1174,9 +1132,19 @@ 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), cascade from pos=%u, ACK sent", + uint8_t rbuf[2560]; + rbuf[0] = DB_MSG_PUSH; + memcpy(rbuf + 1, p, len); + int relay_count = 0; + for (int i = 0; i < si->peer_count; i++) { + if (si->peers[i].sync_state >= 1 && si->peers[i].node_id != src && si->peers[i].node_id != si->db_sync->inst->node_id) { + if (db_sync_send(si, si->peers[i].node_id, rbuf, len + 1) >= 0) + { si_delivery_update(si, rts, rauthor, si->peers[i].node_id); relay_count++; } + } + } + 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 → relayed to %d peers", SI_SHRT(si), (unsigned long long)(src >> 16), - (unsigned long long)rid, (unsigned long long)rauthor, (unsigned long long)rts, ins_pos, total, ins_pos); + (unsigned long long)rid, (unsigned long long)rauthor, (unsigned long long)rts, ins_pos, total, ins_pos, relay_count); 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); @@ -1223,7 +1191,7 @@ static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) sqlite3_finalize(st); } if (!si) { - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, + DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "recv msg type=0x%02x from %016llx hash=%016llx — instance NOT FOUND; sending DB_ERR_NOT_FOUND", type, (unsigned long long)src, (unsigned long long)hash); DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ instance not found for hash=%016llx", (unsigned long long)hash); @@ -1241,7 +1209,6 @@ static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) switch (type) { case DB_MSG_INIT_SYNC: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ handle INIT_SYNC"); db_handle_init_sync(si, src, payload, plen); break; case DB_MSG_INIT_RESP: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ handle INIT_RESP"); db_handle_init_resp(si, src, payload, plen); break; - case DB_MSG_REFINE: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ handle REFINE"); db_handle_refine(si, src, payload, plen); break; case DB_MSG_SEND_DATA: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ handle SEND_DATA"); db_handle_send_data(si, src, payload, plen); break; case DB_MSG_PUSH: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ handle PUSH"); db_handle_push(si, src, payload, plen); break; case DB_MSG_ACK_PUSH: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ handle ACK_PUSH"); db_handle_ack_push(si, src, payload, plen); break; @@ -1289,7 +1256,7 @@ static void db_sync_on_conn_up(struct ETCP_CONN* conn, void* arg) 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; + p->sync_state = 1; 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)); } @@ -1333,9 +1300,9 @@ cd_done: static void db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si, uint64_t pid) { uint32_t mc = db_count(si); - DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "pid=%016llx my=%u tbl=%s", (unsigned long long)pid, mc, SI_TBL(si)); struct SI_PEER* p = si_peer_find(si, pid); - if (p) { p->sync_start_tb = get_time_tb(); p->sync_retry_count = 0; } + DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "pid=%016llx my=%u tbl=%s", (unsigned long long)pid, mc, SI_TBL(si)); + if (p) { p->sync_start_tb = get_time_tb(); } uint8_t msg[5]; msg[0] = DB_MSG_INIT_SYNC; @@ -1426,7 +1393,18 @@ void db_sync_peer_set_state(struct DB_SYNC_INSTANCE* si, uint64_t node_id, uint8 void db_sync_reinitiate(struct DB_SYNC_INSTANCE* si, uint64_t node_id) { struct SI_PEER* p = si_peer_find(si, node_id); - if (p) { p->sync_state = 0; db_sync_initiate_sync(si, node_id); } + if (p) { + p->sync_state = 1; p->sync_start_tb = get_time_tb(); + db_sync_initiate_sync(si, node_id); + } +} + +int db_sync_last_chain_hash8(struct DB_SYNC_INSTANCE* si, uint64_t* out_hash8) +{ + if (!si || !out_hash8) return -1; + uint32_t mc = db_count(si); + if (mc == 0) { *out_hash8 = 0; return 0; } + return db_chain_hash8_at(si, mc - 1, out_hash8); } // ============================================================ @@ -1474,7 +1452,8 @@ static void db_sync_peer_check_cb(void* arg) if (best) { 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++; + 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 ? "," : "", diff --git a/src/db_sync.h b/src/db_sync.h index 21bc592b..ae722d78 100644 --- a/src/db_sync.h +++ b/src/db_sync.h @@ -57,7 +57,6 @@ struct DB_SYNC_INSTANCE; // Message types #define DB_MSG_INIT_SYNC 0x01 #define DB_MSG_INIT_RESP 0x02 -#define DB_MSG_REFINE 0x03 #define DB_MSG_SEND_DATA 0x04 #define DB_MSG_PUSH 0x05 #define DB_MSG_ACK_PUSH 0x06 @@ -79,7 +78,6 @@ struct DB_SYNC_INSTANCE; #define DB_REC_FLAG_WAS_SENT 0x01 // Sync protocol constants -#define DB_REFINE_HASHES 16 #define DB_SEND_DATA_MAX 32 // Ed25519 signature size @@ -129,6 +127,9 @@ void db_sync_peer_set_state(struct DB_SYNC_INSTANCE* si, uint64_t node_id, uint8 // Force re-initiate sync to a specific peer (for testing) void db_sync_reinitiate(struct DB_SYNC_INSTANCE* si, uint64_t node_id); +// Get last chain hash8 (for cross-peer consistency check in tests) +int db_sync_last_chain_hash8(struct DB_SYNC_INSTANCE* si, uint64_t* out_hash8); + #ifdef __cplusplus } #endif diff --git a/src/transport_layer/etcp_connections.c b/src/transport_layer/etcp_connections.c index 0e6e48f7..4c70a87a 100644 --- a/src/transport_layer/etcp_connections.c +++ b/src/transport_layer/etcp_connections.c @@ -368,9 +368,13 @@ static int insert_link_queue(struct ETCP_SOCKET* e_sock, struct ETCP_LINK* link) uint8_t key[LINK_ADDR_KEY_SIZE]; sockaddr_to_key(&link->remote_addr, key); - if (queue_find_data_by_index(e_sock->links_queue, key)) { - DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "insert_link_queue: DUP addr in [%s]", e_sock->name); - return -1; + struct ll_entry* dup_qe = queue_find_data_by_index(e_sock->links_queue, key); + if (dup_qe) { + struct link_queue_entry* dup_lqe = (struct link_queue_entry*)dup_qe->data; + DEBUG_WARN(DEBUG_CATEGORY_CONNECTION, "insert_link_queue: replacing stale DUP addr in [%s] old_link=%p", e_sock->name, dup_lqe->link); + if (dup_lqe->link) dup_lqe->link->link_queue_entry = NULL; + queue_remove_data(e_sock->links_queue, dup_qe); + queue_entry_free(dup_qe); } struct ll_entry* qe = queue_entry_new(sizeof(struct link_queue_entry)); diff --git a/tests/test_db_sync.c b/tests/test_db_sync.c index e09e36bf..3786fa24 100644 --- a/tests/test_db_sync.c +++ b/tests/test_db_sync.c @@ -28,8 +28,8 @@ #include "../lib/debug_config.h" #include "../lib/mem.h" -#define TEST_TIMEOUT_TB 30000 // 3s -#define PHASE_TIMEOUT_TB 25000 // 2.5s per phase +#define TEST_TIMEOUT_TB 120000 // 12s +#define PHASE_TIMEOUT_TB 100000 // 10s per phase #define POLL_INTERVAL_MS 1 static struct UTUN_INSTANCE* inst_a = NULL; @@ -165,6 +165,7 @@ static void remove_si(struct DB_SYNC_INSTANCE** psi) { int main(void) { printf("=== test_db_sync ===\n"); debug_config_init(); debug_set_level(DEBUG_LEVEL_ERROR); + debug_set_category_level(DEBUG_CATEGORY_DB_SYNC, DEBUG_LEVEL_DEBUG); if (create_temp_configs() != 0) { cleanup_temp_configs(); return 1; } utun_instance_set_tun_init_enabled(0); @@ -186,7 +187,7 @@ int main(void) { si_b = db_sync_instance_add(inst_b, "test", 1, 0); if (!si_a || !si_b) { test_phase = 2; goto done; } if (insert_many(si_a, inst_a, 0, 5) != 0) { test_phase = 2; goto done; } - db_sync_reinitiate(si_a, inst_b->node_id); + db_sync_reinitiate(si_b, inst_a->node_id); cb_target = 5; if (!wait_for("B=5", _cond_cb, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } if (db_sync_count(si_b) != 5 || db_sync_chain_verify(si_a) || db_sync_chain_verify(si_b)) @@ -201,7 +202,7 @@ int main(void) { printf("Phase 2: hash_MATCH tail-send\n"); db_sync_peer_set_state(si_a, inst_b->node_id, 0); if (insert_many(si_a, inst_a, 5, 3) != 0) { test_phase = 2; goto done; } - db_sync_reinitiate(si_a, inst_b->node_id); + db_sync_reinitiate(si_b, inst_a->node_id); cb_target = 8; if (!wait_for("B=8", _cond_cb, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } if (db_sync_count(si_b) != 8 || db_sync_chain_verify(si_a) || db_sync_chain_verify(si_b)) @@ -219,7 +220,6 @@ int main(void) { if (!si_a || !si_b) { test_phase = 2; goto done; } if (insert_many(si_a, inst_a, 0, 3) != 0 || insert_many(si_b, inst_b, 10, 2) != 0) { test_phase = 2; goto done; } - db_sync_reinitiate(si_a, inst_b->node_id); db_sync_reinitiate(si_b, inst_a->node_id); ca_target = 5; cb_target = 5; if (!wait_for("both=5", _cond_both, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } @@ -237,7 +237,6 @@ int main(void) { if (!si_a || !si_b) { test_phase = 2; goto done; } if (insert_many(si_a, inst_a, 0, 2) != 0 || insert_many(si_b, inst_b, 10, 2) != 0) { test_phase = 2; goto done; } - db_sync_reinitiate(si_a, inst_b->node_id); db_sync_reinitiate(si_b, inst_a->node_id); ca_target = 4; cb_target = 4; if (!wait_for("both=4", _cond_both, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } @@ -285,9 +284,8 @@ int main(void) { if (!si_a || !si_b || !si_c) { test_phase = 2; goto done; } if (insert_many(si_a, inst_a, 0, 2) != 0 || insert_many(si_b, inst_b, 10, 1) != 0 || insert_many(si_c, inst_c, 20, 1) != 0) { test_phase = 2; goto done; } - db_sync_reinitiate(si_a, inst_b->node_id); db_sync_reinitiate(si_a, inst_c->node_id); - db_sync_reinitiate(si_b, inst_a->node_id); db_sync_reinitiate(si_b, inst_c->node_id); - db_sync_reinitiate(si_c, inst_a->node_id); db_sync_reinitiate(si_c, inst_b->node_id); + db_sync_reinitiate(si_b, inst_a->node_id); + db_sync_reinitiate(si_c, inst_a->node_id); ca_target = 4; cb_target = 4; cc_target = 4; if (!wait_for("all=4", _cond_all, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } if (db_sync_chain_verify(si_a) || db_sync_chain_verify(si_b) || db_sync_chain_verify(si_c)) @@ -319,6 +317,80 @@ int main(void) { { test_phase = 2; goto done; } printf(" mid PUSH cascade=6 PASS\n"); + // =================================================================== + // Phase 7: randomized triple — cold start divergence, live PUSH, + // disconnect/reconnect with additional random inserts. + // =================================================================== + printf("Phase 7: randomized triple\n"); + remove_si(&si_a); remove_si(&si_b); remove_si(&si_c); + srand((unsigned)time(NULL)); + printf("seed=%u\n", (unsigned)time(NULL)); + + si_a = db_sync_instance_add(inst_a, "rnd_triple", 80, 0); + si_b = db_sync_instance_add(inst_b, "rnd_triple", 80, 0); + si_c = db_sync_instance_add(inst_c, "rnd_triple", 80, 0); + if (!si_a || !si_b || !si_c) { test_phase = 2; goto done; } + + // -- Round 1: cold start divergence merge with 0..20 random records per peer -- + int r1_a = rand() % 2, r1_b = rand() % 2, r1_c = rand() % 2; + int total = r1_a + r1_b + r1_c; + printf(" R1: A=%d B=%d C=%d (total=%d)\n", r1_a, r1_b, r1_c, total); + if ((r1_a > 0 && insert_many(si_a, inst_a, 0, r1_a) != 0) + || (r1_b > 0 && insert_many(si_b, inst_b, 100, r1_b) != 0) + || (r1_c > 0 && insert_many(si_c, inst_c, 200, r1_c) != 0)) { test_phase = 2; goto done; } + db_sync_reinitiate(si_b, inst_a->node_id); + db_sync_reinitiate(si_c, inst_a->node_id); + ca_target = cb_target = cc_target = (uint32_t)total; + if (!wait_for("R1 converge", _cond_all, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } + { + uint64_t ha, hb, hc; + if (db_sync_last_chain_hash8(si_a, &ha) || db_sync_last_chain_hash8(si_b, &hb) + || db_sync_last_chain_hash8(si_c, &hc) || ha != hb || hb != hc) { test_phase = 2; goto done; } + if (db_sync_chain_verify(si_a) || db_sync_chain_verify(si_b) || db_sync_chain_verify(si_c)) { test_phase = 2; goto done; } + } + printf(" R1 PASS (%u records)\n", ca_target); + + // -- Round 2: live PUSH — 20 records on random peer, verify convergence -- + { + int r2_peer = rand() % 3; + struct DB_SYNC_INSTANCE* pick = (r2_peer == 0) ? si_a : (r2_peer == 1) ? si_b : si_c; + struct UTUN_INSTANCE* pick_inst = (r2_peer == 0) ? inst_a : (r2_peer == 1) ? inst_b : inst_c; + const char* pick_name = (r2_peer == 0) ? "A" : (r2_peer == 1) ? "B" : "C"; + printf(" R2: +5 on peer %s\n", pick_name); + if (insert_many(pick, pick_inst, 300, 5) != 0) { test_phase = 2; goto done; } + ca_target = cb_target = cc_target = (uint32_t)(total + 5); + if (!wait_for("R2 converge", _cond_all, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } + uint64_t ha, hb, hc; + if (db_sync_last_chain_hash8(si_a, &ha) || db_sync_last_chain_hash8(si_b, &hb) + || db_sync_last_chain_hash8(si_c, &hc) || ha != hb || hb != hc) { test_phase = 2; goto done; } + printf(" R2 PASS (%u records)\n", ca_target); + total += 5; + } + + // -- Round 3: simulate disconnect (reset sync state), add random records, reinitiate, sync -- + printf(" R3: disconnect/reset sync state, insert random, reinitiate\n"); + db_sync_peer_set_state(si_a, inst_b->node_id, 0); db_sync_peer_set_state(si_a, inst_c->node_id, 0); + db_sync_peer_set_state(si_b, inst_a->node_id, 0); + db_sync_peer_set_state(si_c, inst_a->node_id, 0); + + int r3_a = rand() % 2, r3_b = rand() % 2, r3_c = rand() % 2; + total += r3_a + r3_b + r3_c; + printf(" R3: A=%d B=%d C=%d (total=%d)\n", r3_a, r3_b, r3_c, total); + if ((r3_a > 0 && insert_many(si_a, inst_a, 400, r3_a) != 0) + || (r3_b > 0 && insert_many(si_b, inst_b, 500, r3_b) != 0) + || (r3_c > 0 && insert_many(si_c, inst_c, 600, r3_c) != 0)) { test_phase = 2; goto done; } + + db_sync_reinitiate(si_b, inst_a->node_id); + db_sync_reinitiate(si_c, inst_a->node_id); + ca_target = cb_target = cc_target = (uint32_t)total; + if (!wait_for("R3 converge", _cond_all, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } + { + uint64_t ha, hb, hc; + if (db_sync_last_chain_hash8(si_a, &ha) || db_sync_last_chain_hash8(si_b, &hb) + || db_sync_last_chain_hash8(si_c, &hc) || ha != hb || hb != hc) { test_phase = 2; goto done; } + } + printf(" R3 PASS (%u records)\n", ca_target); + done: remove_si(&si_a); remove_si(&si_b); remove_si(&si_c); if (timeout_id && ua) { uasync_cancel_timeout(ua, timeout_id); timeout_id = NULL; }