diff --git a/src/db_sync.c b/src/db_sync.c index 4b16bd12..f261d1f3 100644 --- a/src/db_sync.c +++ b/src/db_sync.c @@ -828,8 +828,9 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, uint64_t src, const } // ---- SEND_DATA batch helper ---- -// Wire format: [id:8][ts:8][author:8][dlen:4][data][sig_len:1=64][sig:64] -static void si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32_t from, uint32_t count, int allocated_buf) +// Wire format: [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) { uint8_t sbuf_stack[8192]; uint8_t* buf = allocated_buf ? u_malloc(8192) : NULL; @@ -842,10 +843,16 @@ static void si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32 uint16_t* rcp = (uint16_t*)(buf + off); off += 2; sqlite3_stmt* stmt; - si_prep(si, &stmt, - "SELECT id,timestamp,author,data,author_signature" + int prep_rc = si_prep(si, &stmt, + "SELECT id,timestamp,node_id,data,author_signature" " FROM \"%s\" ORDER BY timestamp, author_signature" " LIMIT ? OFFSET ?"); + if (prep_rc != SQLITE_OK || !stmt) { + DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] SEND_DATA batch: si_prep failed rc=%d", + SI_SHRT(si), (unsigned long long)(dst >> 16), prep_rc); + if (allocated_buf) u_free(buf); + return -1; + } sqlite3_bind_int64(stmt, 1, (sqlite3_int64)count); sqlite3_bind_int64(stmt, 2, (sqlite3_int64)from); while (sqlite3_step(stmt) == SQLITE_ROW) { @@ -873,7 +880,12 @@ static void si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32 db_sync_send(si, dst, buf, off); DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → SEND DATA: %u records from pos=%u", SI_SHRT(si), (unsigned long long)(dst >> 16), rc, from); + if (rc == 0 && count > 0) { + DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] SEND_DATA batch: loaded 0 records from pos=%u count=%u — SQL error or empty range", + SI_SHRT(si), (unsigned long long)(dst >> 16), from, count); + } if (allocated_buf) u_free(buf); + return (int)rc; } static void db_handle_refine(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len) @@ -892,9 +904,16 @@ static void db_handle_refine(struct DB_SYNC_INSTANCE* si, uint64_t src, const ui SI_SHRT(si), (unsigned long long)(src >> 16), from, mc); return; } - si_send_data_batch(si, src, from, scnt, 1); - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← REFINE: range narrowed to single pos → sending %u records from pos=%u", - SI_SHRT(si), (unsigned long long)(src >> 16), scnt, from); + 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; } @@ -916,10 +935,21 @@ static void db_handle_refine(struct DB_SYNC_INSTANCE* si, uint64_t src, const ui } uint32_t mc = db_count(si); uint32_t scnt = 4; - if (fm > to || fm >= mc) scnt = 0; - si_send_data_batch(si, src, fm, scnt, 0); - 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 → sending %u records from pos=%u", - SI_SHRT(si), (unsigned long long)(src >> 16), from, to, hc, last_match, first_diff, scnt, fm); + 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 ---- @@ -1010,9 +1040,15 @@ static void db_handle_send_data(struct DB_SYNC_INSTANCE* si, uint64_t src, const if (mc > pk && sp && sp->sync_state == 1) { uint32_t scnt = mc - pk; if (scnt > DB_SEND_DATA_MAX) scnt = DB_SEND_DATA_MAX; - si_send_data_batch(si, src, pk, scnt, 1); - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → DATA: requesting next %u records from pos=%u (%u remaining behind)", - SI_SHRT(si), (unsigned long long)(src >> 16), scnt, pk, mc - pk); + 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; + } + 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, @@ -1113,6 +1149,8 @@ static void db_handle_push(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint int ret = db_record_insert(si, rid, rts, rauthor, (const char*)rdata, rdlen, rsig, 0); 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++; } @@ -1123,10 +1161,10 @@ 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), ACK sent", + 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", SI_SHRT(si), (unsigned long long)(src >> 16), - (unsigned long long)rid, (unsigned long long)rauthor, (unsigned long long)rts, ins_pos, db_count(si)); - if (adj_count > 0 && ins_pos < db_count(si) - 1) { + (unsigned long long)rid, (unsigned long long)rauthor, (unsigned long long)rts, ins_pos, total, ins_pos); + 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); } @@ -1217,17 +1255,6 @@ static void db_sync_on_conn_status(struct ETCP_CONN* conn, int status, void* arg } } -#if 0 /* replaced by db_sync_on_conn_status */ -static void db_sync_on_new_conn(struct ETCP_CONN* conn, void* arg) -{ - (void)arg; - if (!conn || !conn->instance || !conn->instance->db_sync) return; - DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "conn=%p node=%016llx", (void*)conn, (unsigned long long)conn->peer_node_id); - etcp_conn_add_up_cbk(conn, db_sync_on_conn_up, NULL); - etcp_conn_add_down_cbk(conn, db_sync_on_conn_down, NULL); -} -#endif - static void db_sync_on_conn_up(struct ETCP_CONN* conn, void* arg) { (void)arg; @@ -1344,6 +1371,39 @@ static void db_verify_chain(struct DB_SYNC_INSTANCE* si) } } +int db_sync_chain_verify(struct DB_SYNC_INSTANCE* si) +{ + if (!si || !si->enabled) return 0; + uint32_t mc = db_count(si); + if (mc == 0) return 0; + + uint8_t prev_ch[32]; memset(prev_ch, 0, 32); + uint8_t exp_ch[32], stored_ch[32]; + + sqlite3_stmt* stmt; + if (si_prep(si, &stmt, + "SELECT id,timestamp,node_id,author_signature,chain_hash FROM \"%s\"" + " ORDER BY timestamp, author_signature") != SQLITE_OK) + return -1; + for (uint32_t pos = 0; sqlite3_step(stmt) == SQLITE_ROW; pos++) { + 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 void* sig_blob = sqlite3_column_blob(stmt, 3); + uint8_t sig[DB_SIG_SIZE]; + if (sig_blob) memcpy(sig, sig_blob, DB_SIG_SIZE); else memset(sig, 0, DB_SIG_SIZE); + const void* b = sqlite3_column_blob(stmt, 4); + if (b && sqlite3_column_bytes(stmt, 4) >= 32) memcpy(stored_ch, b, 32); + else memset(stored_ch, 0, 32); + + db_chain_hash_compute(prev_ch, rid, rts, rauth, sig, exp_ch); + if (memcmp(exp_ch, stored_ch, 32) != 0) { sqlite3_finalize(stmt); return 1; } + memcpy(prev_ch, exp_ch, 32); + } + sqlite3_finalize(stmt); + return 0; +} + // ============================================================ // Timers // ============================================================ diff --git a/src/db_sync.h b/src/db_sync.h index 09fd05a2..85437591 100644 --- a/src/db_sync.h +++ b/src/db_sync.h @@ -120,6 +120,9 @@ typedef void (*db_sync_insert_cb)(struct DB_SYNC_INSTANCE* si, uint64_t author_node_id, void* arg); void db_sync_set_insert_cb(struct DB_SYNC_INSTANCE* si, db_sync_insert_cb cb, void* arg); +// Verify chain integrity: returns 0 if all chain_hashes are correct, 1 if any mismatch found +int db_sync_chain_verify(struct DB_SYNC_INSTANCE* si); + #ifdef __cplusplus } #endif diff --git a/tests/test_db_sync.c b/tests/test_db_sync.c index 1e2255ae..90938684 100644 --- a/tests/test_db_sync.c +++ b/tests/test_db_sync.c @@ -123,6 +123,8 @@ static int cond_links_init(void) { static uint32_t ca_target, cb_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_both(void) { return _cond_ca() && _cond_cb(); } +static int _cond_chain_ok(void) { return db_sync_chain_verify(si_a) == 0 && db_sync_chain_verify(si_b) == 0; } static int insert_many(struct DB_SYNC_INSTANCE* si, struct UTUN_INSTANCE* inst, int start, int count) { char buf[128]; @@ -220,6 +222,65 @@ int main(void) { if (!wait_for("B count=30 (peer empty)", _cond_cb, PHASE_TIMEOUT_TB)) { test_phase=2; goto done; } printf(" Phase 3 PASS: B=%u A=%u\n", db_sync_count(si_b), db_sync_count(si_a)); + // =================================================================== + // Phase 4: Divergence — A=10 B=5 DIFFERENT records in same table. + // Insert via different hash si (no PUSH between them), then re-add + // with same hash → both have existing records → chains differ → + // REFINE → si_send_data_batch → converge to 15. + // =================================================================== + printf("Phase 4: divergence — A=10 B=5 independent records, REFINE/SEND_DATA merge\n"); + + remove_si(&si_a); + remove_si(&si_b); + + { // Insert with different hashes → no PUSH between A and B + struct DB_SYNC_INSTANCE* tmp_a = db_sync_instance_add(inst_a, "test_div", 10); + struct DB_SYNC_INSTANCE* tmp_b = db_sync_instance_add(inst_b, "test_div", 20); + if (!tmp_a || !tmp_b) { test_phase = 2; goto done; } + if (insert_many(tmp_a, inst_a, 0, 10) != 0 || db_sync_count(tmp_a) != 10) { test_phase = 2; goto done; } + if (insert_many(tmp_b, inst_b, 100, 5) != 0 || db_sync_count(tmp_b) != 5) { test_phase = 2; goto done; } + printf(" A has 10 records (indices 0-9), B has 5 records (indices 100-104)\n"); + remove_si(&tmp_a); + remove_si(&tmp_b); + } + + si_a = db_sync_instance_add(inst_a, "test_div", 30); + si_b = db_sync_instance_add(inst_b, "test_div", 30); + if (!si_a || !si_b) { test_phase = 2; goto done; } + + ca_target = 15; cb_target = 15; // 10 + 5 = 15 after merge + if (!wait_for("both count=15 (divergence resolved)", _cond_both, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } + if (db_sync_chain_verify(si_a) != 0 || db_sync_chain_verify(si_b) != 0) { + fprintf(stderr, "FAIL: chain_hash mismatch after divergence merge\n"); test_phase = 2; goto done; + } + printf(" Phase 4 PASS: B=%u A=%u chains OK\n", db_sync_count(si_b), db_sync_count(si_a)); + + // =================================================================== + // Phase 5: PUSH not at tail — insert record with early timestamp + // on A → PUSH to B → B must cascade_from to fix chain_hashes + // =================================================================== + printf("Phase 5: PUSH not at tail — early timestamp, cascade verification\n"); + { + char ebuf[128]; + 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; + size_t jl = strlen(ebuf); memcpy(msg + moff, ebuf, jl); moff += jl; + uint8_t sig[64]; + if (sc_ed25519_sign(inst_a->my_ed25519_privkey, msg, moff, sig) != SC_OK) + { fprintf(stderr, "sign fail\n"); test_phase = 2; goto done; } + if (db_sync_insert_signed(si_a, ebuf, jl, sig, 64, early_ts) != 0) + { fprintf(stderr, "insert_signed fail\n"); test_phase = 2; goto done; } + } + + ca_target = 16; cb_target = 16; + if (!wait_for("both count=16 after PUSH", _cond_both, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; } + if (db_sync_chain_verify(si_a) != 0 || db_sync_chain_verify(si_b) != 0) { + fprintf(stderr, "FAIL: chain_hash mismatch after PUSH cascade\n"); test_phase = 2; goto done; + } + printf(" Phase 5 PASS: B=%u A=%u chains OK\n", db_sync_count(si_b), db_sync_count(si_a)); + done: remove_si(&si_a); remove_si(&si_b); diff --git a/tools/chatgui/src/messagedelegate.cpp b/tools/chatgui/src/messagedelegate.cpp index e408148f..11b93e82 100644 --- a/tools/chatgui/src/messagedelegate.cpp +++ b/tools/chatgui/src/messagedelegate.cpp @@ -582,7 +582,11 @@ void MessageDelegate::drawBubble(QPainter *p, const Layout &L, const QStyleOptio QColor base = opt.palette.window().color(); bubbleFill = base.lighter(135); } else { - bubbleFill = QColor(0xE1, 0xF5, 0xE1); + QColor base = opt.palette.window().color(); + bubbleFill = base.lighter(135); + bubbleFill.setRed(qMax(0, bubbleFill.red() - 10)); + bubbleFill.setGreen(qMin(255, bubbleFill.green() + 5)); + bubbleFill.setBlue(qMax(0, bubbleFill.blue() - 10)); } p->setPen(Qt::NoPen); p->setBrush(bubbleFill);