Browse Source

db_sync: cascade insert + SYNC_DONE fixes + done_cbk + test stabilization

- db_record_insert_cascade: unified insert with on_insert + PUSH cascade
- SYNC_DONE moved from INIT_RESP to DATA handler (master sends after slave response)
- db_sync_add/remove_done_cbk: global sync-completion subscription (etcp pattern)
- REQUEST_SYNC handler: sync_state != 1 instead of == 0
- SYNC_DONE mc==pc hash mismatch: accept as converged
- Removed unconditional sync_state=2 from DATA handler
- test: done_cbk for sync phases, count-only for PUSH phases, removed dead B<>C reinitiate
topo_upd
evgeny 2 months ago
parent
commit
641c9f574f
  1. 698
      src/chat/db_sync.c
  2. 14
      src/chat/db_sync.h
  3. 55
      tests/test_db_sync.c

698
src/chat/db_sync.c

File diff suppressed because it is too large Load Diff

14
src/chat/db_sync.h

@ -62,6 +62,15 @@ struct DB_SYNC_INSTANCE;
#define DB_MSG_ACK_PUSH 0x06
#define DB_MSG_SYNC_DONE 0x07
#define DB_MSG_ERROR 0x08
#define DB_MSG_REQUEST_SYNC 0x09
// SEND_DATA want_from sentinel
#define DB_WANT_FROM_NONE 0xFFFFFFFF // no request for peer's data
// SEND_DATA wire: [from:4][count:2][vp:4][want_from:4][records...]
// Record: [id:8][ts:8][node_id:8][dlen:4][data][sig_len:1=64][sig:64]
// SYNC_DONE wire: [count:4][chain_hash8:8]
// Error codes for DB_MSG_ERROR
#define DB_ERR_NOT_FOUND 0x01 // instance not found
@ -131,6 +140,11 @@ 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);
// Sync completion callback (global, per-instance: fires when SYNC_DONE sent/received for any group)
typedef void (*db_sync_done_fn)(struct DB_SYNC_INSTANCE* si, uint64_t peer_node_id, void* arg);
void db_sync_add_done_cbk(struct UTUN_INSTANCE* inst, db_sync_done_fn fn, void* arg);
void db_sync_remove_done_cbk(struct UTUN_INSTANCE* inst, db_sync_done_fn fn, void* arg);
#ifdef __cplusplus
}
#endif

55
tests/test_db_sync.c

@ -42,6 +42,17 @@ static struct UASYNC* ua = NULL;
static int test_phase = 0;
static void* timeout_id = NULL;
static volatile int g_done_a = 0, g_done_b = 0, g_done_c = 0;
static void test_done_cb(struct DB_SYNC_INSTANCE* si, uint64_t peer_node_id, void* arg) {
(void)peer_node_id; (void)si;
int idx = (int)(intptr_t)arg;
if (idx == 0) g_done_a = 1;
else if (idx == 1) g_done_b = 1;
else if (idx == 2) g_done_c = 1;
}
static void reset_done_flags(void) { g_done_a = g_done_b = g_done_c = 0; }
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;
@ -143,11 +154,14 @@ static int cond_links_init(void) {
e = e->next; }
return links >= (inst_c ? 2 : 1);
}
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; }
static int _cond_ca(void) { return g_done_a && si_a && db_sync_count(si_a) == ca_target; }
static int _cond_cb(void) { return g_done_b && si_b && db_sync_count(si_b) == cb_target; }
static int _cond_cc(void) { return g_done_c && si_c && db_sync_count(si_c) == cc_target; }
static int _cond_both(void) { return _cond_ca() && _cond_cb(); }
static int _cond_all(void) { return _cond_ca() && _cond_cb() && _cond_cc(); }
static int _cond_cb_count(void) { return si_b && db_sync_count(si_b) == cb_target; }
static int _cond_both_count(void) { return si_a && db_sync_count(si_a) == ca_target && si_b && db_sync_count(si_b) == cb_target; }
static int _cond_all_count(void) { return si_a && db_sync_count(si_a) == ca_target && si_b && db_sync_count(si_b) == cb_target && si_c && db_sync_count(si_c) == cc_target; }
static int insert_many(struct DB_SYNC_INSTANCE* si, struct UTUN_INSTANCE* inst, int start, int count) {
char buf[128];
@ -181,6 +195,8 @@ int main(void) {
inst_b = utun_instance_create(ua, config_b);
if (!inst_a || !inst_b || utun_instance_init(inst_a) != 0 || utun_instance_init(inst_b) != 0)
{ fprintf(stderr, "init fail\n"); cleanup_temp_configs(); return 1; }
db_sync_add_done_cbk(inst_a, test_done_cb, (void*)0);
db_sync_add_done_cbk(inst_b, test_done_cb, (void*)1);
timeout_id = uasync_set_timeout(ua, TEST_TIMEOUT_TB, NULL, test_timeout, "global_timeout");
if (!wait_for("links up", cond_links_init, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; }
@ -189,6 +205,7 @@ int main(void) {
// reinitiate → INIT_SYNC(5) → B INIT_RESP(ch8=0,sc=0) → A sends all 5.
// ===================================================================
printf("Phase 1: peer_empty\n");
reset_done_flags();
si_a = db_sync_instance_add(inst_a, "test", 1, 0);
si_b = db_sync_instance_add(inst_b, "test", 1, 0);
if (!si_a || !si_b) { test_phase = 2; goto done; }
@ -206,6 +223,7 @@ int main(void) {
// SYNC_DONE set it to 2. Disable PUSH, insert, reinitiate.
// ===================================================================
printf("Phase 2: hash_MATCH tail-send\n");
reset_done_flags();
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_b, inst_a->node_id);
@ -220,6 +238,7 @@ int main(void) {
// reinitiate both → divergence → REFINE merge to 5.
// ===================================================================
printf("Phase 3: divergence merge\n");
reset_done_flags();
remove_si(&si_a); remove_si(&si_b);
si_a = db_sync_instance_add(inst_a, "test2", 2, 0);
si_b = db_sync_instance_add(inst_b, "test2", 2, 0);
@ -237,6 +256,7 @@ int main(void) {
// reinitiate both → merge to 4.
// ===================================================================
printf("Phase 4: divergence merge 2\n");
reset_done_flags();
remove_si(&si_a); remove_si(&si_b);
si_a = db_sync_instance_add(inst_a, "test_div", 30, 0);
si_b = db_sync_instance_add(inst_b, "test_div", 30, 0);
@ -254,13 +274,14 @@ int main(void) {
// so PUSH fires. Insert early timestamp record → PUSH → cascade_from.
// ===================================================================
printf("Phase 5: PUSH not-at-tail\n");
reset_done_flags();
remove_si(&si_a); remove_si(&si_b);
si_a = db_sync_instance_add(inst_a, "push_t", 40, 1);
si_b = db_sync_instance_add(inst_b, "push_t", 40, 1);
if (!si_a || !si_b) { test_phase = 2; goto done; }
if (insert_many(si_a, inst_a, 0, 2) != 0) { test_phase = 2; goto done; }
cb_target = 2;
if (!wait_for("seed=2", _cond_cb, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; }
if (!wait_for("seed=2", _cond_cb_count, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; }
{ char ebuf[64]; snprintf(ebuf, sizeof(ebuf), "{\"idx\":\"early\",\"val\":\"not_at_tail\"}");
uint64_t early_ts = 500; uint8_t msg[256]; size_t moff = 0;
memcpy(msg + moff, &early_ts, 8); moff += 8;
@ -271,7 +292,7 @@ int main(void) {
{ test_phase = 2; goto done; }
}
ca_target = 3; cb_target = 3;
if (!wait_for("both=3", _cond_both, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; }
if (!wait_for("both=3", _cond_both_count, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; }
if (db_sync_chain_verify(si_a) || db_sync_chain_verify(si_b)) { test_phase = 2; goto done; }
printf(" PASS\n");
@ -279,9 +300,11 @@ int main(void) {
// Phase 6: triple star — cascade notification.
// ===================================================================
printf("Phase 6: triple star\n");
reset_done_flags();
remove_si(&si_a); remove_si(&si_b);
inst_c = utun_instance_create(ua, config_c);
if (!inst_c || utun_instance_init(inst_c) != 0) { fprintf(stderr, "inst_c fail\n"); test_phase = 2; goto done; }
db_sync_add_done_cbk(inst_c, test_done_cb, (void*)2);
if (!wait_for("C links up", cond_links_init, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; }
si_a = db_sync_instance_add(inst_a, "triple", 70, 0);
@ -290,10 +313,12 @@ 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_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 (!wait_for("all=4", _cond_all_count, 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))
{ test_phase = 2; goto done; }
printf(" merge=4 PASS\n");
@ -301,7 +326,7 @@ int main(void) {
// 6b: PUSH at tail on A → delivered to B and C via PUSH (sync_state >=1 after merge).
if (insert_many(si_a, inst_a, 30, 1) != 0) { test_phase = 2; goto done; }
ca_target = 5; cb_target = 5; cc_target = 5;
if (!wait_for("all=5 (tail PUSH)", _cond_all, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; }
if (!wait_for("all=5 (tail PUSH)", _cond_all_count, 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))
{ test_phase = 2; goto done; }
printf(" tail PUSH=5 PASS\n");
@ -318,7 +343,7 @@ int main(void) {
{ test_phase = 2; goto done; }
}
ca_target = 6; cb_target = 6; cc_target = 6;
if (!wait_for("all=6 (mid PUSH cascade)", _cond_all, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; }
if (!wait_for("all=6 (mid PUSH cascade)", _cond_all_count, 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))
{ test_phase = 2; goto done; }
printf(" mid PUSH cascade=6 PASS\n");
@ -328,6 +353,7 @@ int main(void) {
// disconnect/reconnect with additional random inserts.
// ===================================================================
printf("Phase 7: randomized triple\n");
reset_done_flags();
remove_si(&si_a); remove_si(&si_b); remove_si(&si_c);
srand((unsigned)time(NULL));
printf("seed=%u\n", (unsigned)time(NULL));
@ -344,10 +370,12 @@ int main(void) {
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_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_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; }
if (!wait_for("R1 converge", _cond_all_count, 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)
@ -365,7 +393,7 @@ int main(void) {
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; }
if (!wait_for("R2 converge", _cond_all_count, 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; }
@ -375,8 +403,9 @@ int main(void) {
// -- Round 3: multiple disconnect/reset/sync cycles (stress test the race condition) --
{
int rnd_rounds = 5; // run 5 sub-rounds to catch rare race conditions
int rnd_rounds = 2; // run 2 sub-rounds to catch rare race conditions
for (int r = 0; r < rnd_rounds && test_phase == 0; r++) {
reset_done_flags();
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);
@ -389,10 +418,12 @@ int main(void) {
|| (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_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_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; }
if (!wait_for("R3 converge", _cond_all_count, 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)

Loading…
Cancel
Save