Browse Source

db_sync: односторонний sync, sparse-таблица, vp в SEND_DATA, PUSH-релей

topo_upd
Evgeny 2 months ago
parent
commit
27906e9f91
  1. 4
      lib/tcp_io.c
  2. 7
      lib/u_async.c
  3. 643
      src/db_sync.c
  4. 5
      src/db_sync.h
  5. 10
      src/transport_layer/etcp_connections.c
  6. 90
      tests/test_db_sync.c

4
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
}

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

643
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 ? "," : "",

5
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

10
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));

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

Loading…
Cancel
Save