From c093630d4d478fdd0eaa819e4ee9aed3af81dd78 Mon Sep 17 00:00:00 2001 From: evgeny Date: Thu, 30 Jul 2026 15:53:23 +0300 Subject: [PATCH] =?UTF-8?q?db=5Fsync:=20async=20batch=20chain=5Fhash=20rec?= =?UTF-8?q?alc=20(batch=2050,=20timeout=200=E2=86=92100tb),=20flush=20befo?= =?UTF-8?q?re=20protocol=20chain=5Fhash=20reads,=20remove=20db=5Fcascade?= =?UTF-8?q?=5Ffrom/range=20+=20do=5Fcascade=20param?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- src/chat/db_sync.c | 246 +++++++++++++++++++++++-------------------- tests/test_db_sync.c | 58 ++++++---- 2 files changed, 165 insertions(+), 139 deletions(-) diff --git a/src/chat/db_sync.c b/src/chat/db_sync.c index 70e5e662..846d6b81 100644 --- a/src/chat/db_sync.c +++ b/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); diff --git a/tests/test_db_sync.c b/tests/test_db_sync.c index 32cb8177..d63544de 100644 --- a/tests/test_db_sync.c +++ b/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; }