Browse Source

db_sync: async batch chain_hash recalc (batch 50, timeout 0→100tb), flush before protocol chain_hash reads, remove db_cascade_from/range + do_cascade param

topo_upd
evgeny 2 months ago
parent
commit
c093630d4d
  1. 246
      src/chat/db_sync.c
  2. 58
      tests/test_db_sync.c

246
src/chat/db_sync.c

@ -44,6 +44,7 @@ struct SI_PEER {
uint32_t last_peer_count;
uint8_t sync_state;
uint64_t sync_start_tb;
uint32_t retry_count;
};
struct DB_SYNC_INSTANCE {
@ -54,6 +55,10 @@ struct DB_SYNC_INSTANCE {
uint64_t last_timestamp_ms;
uint8_t enabled;
void* ttl_timer;
struct {
uint32_t pos;
void* timer;
} recalc;
struct SI_PEER* peers;
int peer_count, peer_capacity;
db_sync_insert_cb on_insert;
@ -77,7 +82,11 @@ struct DB_SYNC {
#define SI_DB(si) ((si)->db_sync->db)
#define SI_TBL(si) ((si)->table_name)
#define SI_SHRT(si) ({ const char* _t = SI_TBL(si); const char* _u = strrchr(_t, '_'); _u ? _u + 1 : _t; })
#define SI_SHRT(si) ({ \
static char _tbuf[28]; \
snprintf(_tbuf, sizeof(_tbuf), "%s.%04llX", SI_TBL(si), (unsigned long long)((si)->hash & 0xFFFF)); \
_tbuf; \
})
static void si_sql(char* buf, size_t sz, struct DB_SYNC_INSTANCE* si, const char* fmt)
{
@ -304,45 +313,67 @@ static int db_prev_chain_hash(struct DB_SYNC_INSTANCE* si, uint64_t ts, const ui
return 0;
}
static void db_cascade_from(struct DB_SYNC_INSTANCE* si, uint32_t from_pos)
// ── Chain hash fragment (first 8 bytes) for sync protocol comparisons ──
static int db_chain_hash8_at(struct DB_SYNC_INSTANCE* si, uint32_t pos, uint64_t* h8)
{
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "tbl=%s from=%u", SI_TBL(si), from_pos);
uint8_t ch[32];
if (db_chain_hash_at(si, pos, ch) != 0) return -1;
memcpy(h8, ch, 8);
return 0;
}
// ── Recalc scheduler (batch of 50, async via timer) ──
static void db_recalc_tick(void* arg);
static void db_sync_schedule_recalc(struct DB_SYNC_INSTANCE* si, uint32_t from_pos)
{
if (!si->recalc.timer) {
si->recalc.pos = from_pos;
si->recalc.timer = uasync_set_timeout(si->db_sync->inst->ua, 0, si, db_recalc_tick, "db_recalc");
return;
}
if (from_pos < si->recalc.pos)
si->recalc.pos = from_pos;
}
static void db_recalc_tick(void* arg)
{
struct DB_SYNC_INSTANCE* si = (struct DB_SYNC_INSTANCE*)arg;
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_from BEGIN: %s", sqlite3_errmsg(db)); return; }
uint32_t total = db_count(si);
uint8_t prev_ch[32]; memset(prev_ch, 0, 32);
if (from_pos > 0) {
sqlite3_stmt* stmt;
if (si_prep(si, &stmt,
"SELECT chain_hash FROM \"%s\""
" ORDER BY timestamp, author_signature"
" LIMIT 1 OFFSET ?") == SQLITE_OK)
{
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)(from_pos - 1));
if (sqlite3_step(stmt) == SQLITE_ROW) {
const void* b = sqlite3_column_blob(stmt, 0);
if (b && sqlite3_column_bytes(stmt, 0) >= 32) memcpy(prev_ch, b, 32);
}
sqlite3_finalize(stmt);
}
if (si->recalc.pos >= total) { si->recalc.timer = NULL; return; }
uint8_t prev_ch[32];
if (si->recalc.pos > 0) {
if (db_chain_hash_at(si, si->recalc.pos - 1, prev_ch) != 0)
memset(prev_ch, 0, 32);
} else {
memset(prev_ch, 0, 32);
}
uint32_t batch = (total - si->recalc.pos > 50) ? 50 : total - si->recalc.pos;
int rc = sqlite3_exec(db, "BEGIN IMMEDIATE", NULL, NULL, NULL);
if (rc != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "recalc BEGIN: %s", sqlite3_errmsg(db)); return; }
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)
{ DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_cascade_from SELECT prep: %s", sqlite3_errmsg(db)); sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return; }
sqlite3_bind_int64(sel, 1, (sqlite3_int64)from_pos);
{ DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "recalc SELECT prep: %s", sqlite3_errmsg(db)); sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return; }
sqlite3_bind_int64(sel, 1, (sqlite3_int64)si->recalc.pos);
if (si_prep(si, &upd,
"UPDATE \"%s\" SET chain_hash=?"
" WHERE timestamp=? AND author_signature=?") != SQLITE_OK)
{ DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_cascade_from UPDATE prep: %s", sqlite3_errmsg(db)); sqlite3_finalize(sel); sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return; }
{ DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "recalc UPDATE prep: %s", sqlite3_errmsg(db)); sqlite3_finalize(sel); sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return; }
uint8_t ch[32];
while (sqlite3_step(sel) == SQLITE_ROW) {
uint8_t ch[32]; uint32_t processed = 0;
while (sqlite3_step(sel) == SQLITE_ROW && processed < batch) {
sqlite3_int64 n_id = sqlite3_column_int64(sel, 0);
sqlite3_int64 n_ts = sqlite3_column_int64(sel, 1);
sqlite3_int64 n_auth = sqlite3_column_int64(sel, 2);
@ -356,22 +387,34 @@ static void db_cascade_from(struct DB_SYNC_INSTANCE* si, uint32_t from_pos)
sqlite3_bind_blob(upd, 3, sig, DB_SIG_SIZE, SQLITE_STATIC);
sqlite3_step(upd);
memcpy(prev_ch, ch, 32);
processed++;
}
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_from COMMIT: %s", sqlite3_errmsg(db));
}
if (rc != SQLITE_OK) DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "recalc COMMIT: %s", sqlite3_errmsg(db));
// ── Chain hash fragment (first 8 bytes) for sync protocol comparisons ──
si->recalc.pos += processed;
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "recalc tick tbl=%s pos=%u batch=%u/%u total=%u",
SI_TBL(si), si->recalc.pos - processed, processed, batch, total);
static int db_chain_hash8_at(struct DB_SYNC_INSTANCE* si, uint32_t pos, uint64_t* h8)
if (si->recalc.pos >= total) {
si->recalc.timer = NULL;
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "recalc complete tbl=%s pos=%u(total)", SI_TBL(si), si->recalc.pos);
} else {
si->recalc.timer = uasync_set_timeout(si->db_sync->inst->ua, 100, si, db_recalc_tick, "db_recalc");
}
}
static void db_sync_flush_recalc(struct DB_SYNC_INSTANCE* si)
{
uint8_t ch[32];
if (db_chain_hash_at(si, pos, ch) != 0) return -1;
memcpy(h8, ch, 8);
return 0;
struct UASYNC* ua = si->db_sync->inst->ua;
while (si->recalc.timer) {
uasync_cancel_timeout(ua, si->recalc.timer);
si->recalc.timer = NULL;
db_recalc_tick(si);
}
}
// ── Ed25519 verification ──
@ -391,6 +434,31 @@ static int db_get_ed25519_pubkey(struct DB_SYNC* db, uint64_t node_id, uint8_t o
return -1;
}
/* dump table: compact one-line per record with pos, id, ts, author, ch8 */
static void si_dump_table(struct DB_SYNC_INSTANCE* si, const char* tag)
{
uint32_t mc = db_count(si);
if (mc > 200) return; /* skip very large tables */
sqlite3_stmt* stmt;
if (si_prep(si, &stmt,
"SELECT id,timestamp,node_id,chain_hash FROM \"%s\" ORDER BY timestamp,author_signature") != SQLITE_OK) return;
char buf[256]; int total = 0;
while (sqlite3_step(stmt) == SQLITE_ROW) {
uint64_t id = (uint64_t)sqlite3_column_int64(stmt, 0);
uint64_t ts = (uint64_t)sqlite3_column_int64(stmt, 1);
uint64_t auth = (uint64_t)sqlite3_column_int64(stmt, 2);
const uint8_t* ch = (const uint8_t*)sqlite3_column_blob(stmt, 3);
uint64_t ch8 = 0; if (ch) memcpy(&ch8, ch, 8);
snprintf(buf, sizeof(buf), "[%2u] id=%-2llu ts=%-10llu auth=%04llX ch8=%016llX",
total, (unsigned long long)id, (unsigned long long)ts,
(unsigned long long)(auth >> 16), ch8);
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, " %s %s", tag, buf);
total++;
}
sqlite3_finalize(stmt);
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, " %s: table=%s mc=%u", tag, SI_TBL(si), mc);
}
static int db_verify_author_sig(struct DB_SYNC_INSTANCE* si,
uint64_t ts, uint64_t author_node_id,
const char* json, size_t jlen,
@ -418,10 +486,10 @@ static int db_verify_author_sig(struct DB_SYNC_INSTANCE* si,
static int db_record_insert(struct DB_SYNC_INSTANCE* si,
uint64_t id, uint64_t ts, uint64_t author_node_id,
const char* json, size_t jlen,
const uint8_t* author_sig, int do_cascade,
const uint8_t* author_sig,
const char* local_attrs)
{
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "id=%llu ts=%llu author=%N len=%zu cascade=%d", (unsigned long long)id, (unsigned long long)ts, (unsigned long long)author_node_id, jlen, do_cascade);
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "id=%llu ts=%llu author=%N len=%zu", (unsigned long long)id, (unsigned long long)ts, (unsigned long long)author_node_id, jlen);
sqlite3* db = SI_DB(si);
sqlite3_stmt* stmt;
int rc;
@ -435,7 +503,7 @@ static int db_record_insert(struct DB_SYNC_INSTANCE* si,
sqlite3_step(del);
int chg = sqlite3_changes(SI_DB(si));
sqlite3_finalize(del);
if (chg > 0) db_cascade_from(si, si_find_pos(si, ts, author_sig));
if (chg > 0) db_sync_schedule_recalc(si, si_find_pos(si, ts, author_sig));
}
return -2;
}
@ -478,48 +546,13 @@ static int db_record_insert(struct DB_SYNC_INSTANCE* si,
sqlite3_finalize(stmt);
if (rc != SQLITE_DONE) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "INSERT: %s", sqlite3_errmsg(db)); sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return -1; }
// Recompute chain hashes for all subsequent records
if (do_cascade) {
sqlite3_stmt* sel;
if (si_prep(si, &sel,
"SELECT timestamp,id,node_id,author_signature FROM \"%s\""
" WHERE timestamp>?1"
" OR (timestamp=?1 AND author_signature>?2)"
" ORDER BY timestamp, author_signature") != SQLITE_OK)
{ DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_record_insert cascade SELECT prep: %s", sqlite3_errmsg(db)); sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return -1; }
sqlite3_bind_int64(sel, 1, (sqlite3_int64)ts);
sqlite3_bind_blob(sel, 2, author_sig, DB_SIG_SIZE, SQLITE_STATIC);
sqlite3_stmt* upd;
if (si_prep(si, &upd,
"UPDATE \"%s\" SET chain_hash=?"
" WHERE timestamp=? AND author_signature=?") != SQLITE_OK)
{ DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_record_insert cascade UPDATE prep: %s", sqlite3_errmsg(db)); sqlite3_finalize(sel); sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return -1; }
uint8_t rch[32]; memcpy(rch, ch, 32);
while (sqlite3_step(sel) == SQLITE_ROW) {
sqlite3_int64 n_ts = sqlite3_column_int64(sel, 0);
sqlite3_int64 n_id = sqlite3_column_int64(sel, 1);
sqlite3_int64 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);
uint8_t nch[32];
db_chain_hash_compute(rch, (uint64_t)n_id, (uint64_t)n_ts, (uint64_t)n_auth, sig, nch);
sqlite3_reset(upd);
sqlite3_bind_blob(upd, 1, nch, 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(rch, nch, 32);
}
sqlite3_finalize(upd);
sqlite3_finalize(sel);
}
si->next_id = id + 1;
rc = sqlite3_exec(db, "COMMIT", NULL, NULL, NULL);
if (rc != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "insert COMMIT: %s", sqlite3_errmsg(db)); return -1; }
// Trigger async chain_hash recalc from this record onward
uint32_t ins_pos = si_find_pos(si, ts, author_sig);
db_sync_schedule_recalc(si, ins_pos);
return 0;
}
@ -621,6 +654,7 @@ static void db_handle_error(struct DB_SYNC_INSTANCE* si, uint64_t src, const uin
static void db_handle_init_sync(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len)
{
if (len < 4) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "INIT_SYNC too short %zu from %016llx", len, (unsigned long long)src); return; }
db_sync_flush_recalc(si);
uint32_t pc = *(uint32_t*)p;
uint32_t mc = db_count(si);
uint32_t tp = (pc < mc ? pc : mc);
@ -659,6 +693,10 @@ static void db_handle_init_sync(struct DB_SYNC_INSTANCE* si, uint64_t src, const
}
resp[sp] = (uint8_t)scnt;
int sret = db_sync_send(si, src, resp, off);
{
struct SI_PEER* sp = si_peer_add(si, src);
if (sp && sp->sync_state == 0) { sp->sync_state = 1; sp->sync_start_tb = get_time_tb(); }
}
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → INIT_RESP: sent %u bytes to peer=%016llx ret=%d flow=%u,%u",
SI_SHRT(si), (unsigned long long)(src >> 16), off, (unsigned long long)src, sret,
resp[0], resp[1]);
@ -667,6 +705,7 @@ 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 < 14) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "INIT_RESP too short %zu from %016llx", len, (unsigned long long)src); return; }
db_sync_flush_recalc(si);
uint32_t tp = *(uint32_t*)p;
uint64_t peer_ch8; memcpy(&peer_ch8, p + 4, 8);
uint8_t has_tail = p[12];
@ -792,7 +831,7 @@ static int si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32_
sqlite3_bind_blob(del, 2, rsig, DB_SIG_SIZE, SQLITE_STATIC);
sqlite3_step(del);
sqlite3_finalize(del);
db_cascade_from(si, del_pos);
db_sync_schedule_recalc(si, del_pos);
bad_del++;
}
continue;
@ -853,39 +892,6 @@ static int si_parse_record(const uint8_t** pp, const uint8_t* end,
}
// ---- 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 < 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; }
@ -903,6 +909,7 @@ static void db_handle_send_data(struct DB_SYNC_INSTANCE* si, uint64_t src, const
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);
db_sync_flush_recalc(si);
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);
@ -925,7 +932,7 @@ static void db_handle_send_data(struct DB_SYNC_INSTANCE* si, uint64_t src, const
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, NULL);
int ret = db_record_insert(si, rid, rts, rauthor, (const char*)rdata, rdlen, rsig, NULL);
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++;
@ -937,7 +944,6 @@ static void db_handle_send_data(struct DB_SYNC_INSTANCE* si, uint64_t src, const
if (total_processed > 0) {
uint32_t rst = (from < old_synced_pos) ? from : old_synced_pos;
db_cascade_range(si, rst, from + total_processed);
if (sp) sp->synced_pos = from + total_processed - 1;
for (int j = 0; j < si->peer_count; j++) {
if (si->peers[j].synced_pos >= rst && &si->peers[j] != sp)
@ -993,7 +999,7 @@ static void db_handle_send_data(struct DB_SYNC_INSTANCE* si, uint64_t src, const
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);
if (old_synced_pos < nc) db_sync_schedule_recalc(si, old_synced_pos);
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);
}
@ -1051,6 +1057,7 @@ static void db_handle_send_data(struct DB_SYNC_INSTANCE* si, uint64_t src, const
}
}
db_sync_flush_recalc(si);
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, &mc, 4); memcpy(sd + 5, &lch8, 8);
@ -1063,6 +1070,7 @@ static void db_handle_send_data(struct DB_SYNC_INSTANCE* si, uint64_t src, const
static void db_handle_sync_done(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len)
{
if (len < 12) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← SYNC_DONE: truncated len=%zu (need >=12)", SI_SHRT(si), (unsigned long long)(src >> 16), len); return; }
db_sync_flush_recalc(si);
uint32_t pc = *(uint32_t*)p;
uint64_t pch8; memcpy(&pch8, p + 4, 8);
uint32_t mc = db_count(si);
@ -1105,6 +1113,7 @@ static void db_handle_sync_done(struct DB_SYNC_INSTANCE* si, uint64_t src, const
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"sync [%s:%04llX] ← SYNC_DONE: count diff my=%u peer=%u synced=%u — same pc, accept",
SI_SHRT(si), (unsigned long long)(src >> 16), mc, pc, sp ? sp->synced_pos : 0);
si_dump_table(si, "SAME_PC_ACCEPT");
return;
}
if (sp) sp->last_peer_count = pc;
@ -1121,6 +1130,7 @@ static void db_handle_sync_done(struct DB_SYNC_INSTANCE* si, uint64_t src, const
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);
si_dump_table(si, "HASH_MISMATCH");
if (sp) { sp->sync_state = 1; sp->sync_start_tb = get_time_tb(); }
db_sync_initiate_sync(si, src);
}
@ -1162,11 +1172,10 @@ static void db_handle_push(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint
if (si_parse_record(&ptr, p + len, &rid, &rts, &rauthor, &rdlen, &rdata, &rsig, &rsiglen) != 0) return;
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "from=%016llx id=%llu author=%016llx ts=%llu len=%u", (unsigned long long)src, (unsigned long long)rid, (unsigned long long)rauthor, (unsigned long long)rts, rdlen);
int ret = db_record_insert(si, rid, rts, rauthor, (const char*)rdata, rdlen, rsig, 0, NULL);
int ret = db_record_insert(si, rid, rts, rauthor, (const char*)rdata, rdlen, rsig, NULL);
if (ret == 0) {
uint32_t ins_pos = si_find_pos(si, rts, rsig);
uint32_t total = db_count(si);
db_cascade_from(si, ins_pos);
int adj_count = 0;
for (int j = 0; j < si->peer_count; j++) {
if (si->peers[j].synced_pos >= ins_pos) { si->peers[j].synced_pos = ins_pos > 0 ? ins_pos - 1 : 0; adj_count++; }
@ -1392,13 +1401,14 @@ static void db_verify_chain(struct DB_SYNC_INSTANCE* si)
if (bad_pos >= 0) {
DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "chain_hash mismatch at pos %d/%u in %s, recalculating",
bad_pos, mc, SI_TBL(si));
db_cascade_from(si, (uint32_t)bad_pos);
db_sync_schedule_recalc(si, (uint32_t)bad_pos);
}
}
int db_sync_chain_verify(struct DB_SYNC_INSTANCE* si)
{
if (!si || !si->enabled) return 0;
db_sync_flush_recalc(si);
uint32_t mc = db_count(si);
if (mc == 0) return 0;
@ -1447,6 +1457,7 @@ void db_sync_reinitiate(struct DB_SYNC_INSTANCE* si, uint64_t node_id)
int db_sync_last_chain_hash8(struct DB_SYNC_INSTANCE* si, uint64_t* out_hash8)
{
if (!si || !out_hash8) return -1;
db_sync_flush_recalc(si);
uint32_t mc = db_count(si);
if (mc == 0) { *out_hash8 = 0; return 0; }
return db_chain_hash8_at(si, mc - 1, out_hash8);
@ -1524,6 +1535,7 @@ static void db_sync_peer_check_cb(void* arg)
"sync [%s:%04llX] TIMEOUT: no response for %llu ms, sync_state 1→0 (synced_pos was %u)",
SI_SHRT(si), (unsigned long long)(p->node_id >> 16),
(unsigned long long)((now - p->sync_start_tb) / 10), p->synced_pos);
si_dump_table(si, "STUCK_SYNC_TIMEOUT");
p->sync_state = 0;
static int dump_ctr = 0;
if (dump_ctr++ == 0 || (dump_ctr % 10) == 0)
@ -1639,6 +1651,7 @@ void db_sync_destroy(struct UTUN_INSTANCE* inst)
for (int i = 0; i < db->instance_count; i++) {
struct DB_SYNC_INSTANCE* si = &db->instances[i];
if (si->ttl_timer) { uasync_cancel_timeout(inst->ua, si->ttl_timer); si->ttl_timer = NULL; }
if (si->recalc.timer) { uasync_cancel_timeout(inst->ua, si->recalc.timer); si->recalc.timer = NULL; }
if (si->peers) u_free(si->peers);
}
@ -1740,6 +1753,7 @@ void db_sync_instance_remove(struct DB_SYNC_INSTANCE* si)
si->enabled = 0;
if (si->ttl_timer) { uasync_cancel_timeout(inst->ua, si->ttl_timer); si->ttl_timer = NULL; }
if (si->recalc.timer) { uasync_cancel_timeout(inst->ua, si->recalc.timer); si->recalc.timer = NULL; }
if (si->peers) { u_free(si->peers); si->peers = NULL; si->peer_count = si->peer_capacity = 0; }
int idx = (int)(si - db->instances);
@ -1762,7 +1776,7 @@ int db_sync_insert_signed(struct DB_SYNC_INSTANCE* si, const char* json_data, si
uint64_t id = si->next_id;
uint64_t author_node_id = si->db_sync->inst->node_id;
int ret = db_record_insert(si, id, ts, author_node_id, json_data, len, sig, 1, local_attrs);
int ret = db_record_insert(si, id, ts, author_node_id, json_data, len, sig, local_attrs);
if (ret != 0) return ret;
if (si->on_insert) si->on_insert(si, ts, json_data, len, author_node_id, si->on_insert_arg);

58
tests/test_db_sync.c

@ -45,6 +45,7 @@ static void* timeout_id = NULL;
static char temp_dir[] = "/tmp/utun_dbsync_XXXXXX";
static char config_a[256], config_b[256], config_c[256];
static int port_a_srv, port_b_srv, port_c_srv;
static uint32_t ca_target, cb_target, cc_target;
static int write_file(const char* path, const char* fmt, ...) {
va_list ap; FILE* f = fopen(path, "w");
@ -123,7 +124,13 @@ static int wait_for(const char* desc, int (*cond)(void), int timeout_tb) {
while (!cond() && (get_time_tb() - start) < (uint64_t)timeout_tb && test_phase == 0)
uasync_poll(ua, POLL_INTERVAL_MS);
if (cond()) return 1;
if (test_phase == 0) { fprintf(stderr, "TIMEOUT: %s\n", desc); test_phase = 2; }
if (test_phase == 0) {
fprintf(stderr, "TIMEOUT: %s — ca=%u/%u cb=%u/%u cc=%u/%u\n", desc,
si_a ? db_sync_count(si_a) : 0, ca_target,
si_b ? db_sync_count(si_b) : 0, cb_target,
si_c ? db_sync_count(si_c) : 0, cc_target);
test_phase = 2;
}
return 0;
}
@ -136,7 +143,6 @@ static int cond_links_init(void) {
e = e->next; }
return links >= (inst_c ? 2 : 1);
}
static uint32_t ca_target, cb_target, cc_target;
static int _cond_ca(void) { return si_a && db_sync_count(si_a) == ca_target; }
static int _cond_cb(void) { return si_b && db_sync_count(si_b) == cb_target; }
static int _cond_cc(void) { return si_c && db_sync_count(si_c) == cc_target; }
@ -367,29 +373,34 @@ int main(void) {
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; }
// -- Round 3: multiple disconnect/reset/sync cycles (stress test the race condition) --
{
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; }
int rnd_rounds = 5; // run 5 sub-rounds to catch rare race conditions
for (int r = 0; r < rnd_rounds && test_phase == 0; r++) {
printf(" R3.%d: disconnect/reset, insert random, reinitiate\n", r);
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_b, inst_c->node_id, 0);
db_sync_peer_set_state(si_c, inst_a->node_id, 0); db_sync_peer_set_state(si_c, inst_b->node_id, 0);
int ra = rand() % 3, rb = rand() % 3, rc = rand() % 3;
total += ra + rb + rc;
printf(" R3.%d: A=%d B=%d C=%d (total=%d)\n", r, ra, rb, rc, total);
if ((ra > 0 && insert_many(si_a, inst_a, 400 + r*100, ra) != 0)
|| (rb > 0 && insert_many(si_b, inst_b, 500 + r*100, rb) != 0)
|| (rc > 0 && insert_many(si_c, inst_c, 600 + r*100, rc) != 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);
if (test_phase == 0) printf(" R3 PASS (%u records)\n", ca_target);
done:
remove_si(&si_a); remove_si(&si_b); remove_si(&si_c);
@ -400,6 +411,7 @@ done:
if (ua) { uasync_destroy(ua, 0); ua = NULL; }
cleanup_temp_configs();
if (test_phase == 0) test_phase = 1;
debug_disable_file_output();
printf("=== %s ===\n", test_phase == 1 ? "PASS" : "FAIL");
return test_phase == 1 ? 0 : 1;
}

Loading…
Cancel
Save