From 641c9f574f939501b56e9a93b01f88560efdbe6a Mon Sep 17 00:00:00 2001 From: evgeny Date: Thu, 30 Jul 2026 20:23:47 +0300 Subject: [PATCH] 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 --- src/chat/db_sync.c | 698 +++++++++++++++++++++---------------------- src/chat/db_sync.h | 14 + tests/test_db_sync.c | 55 +++- 3 files changed, 393 insertions(+), 374 deletions(-) diff --git a/src/chat/db_sync.c b/src/chat/db_sync.c index 846d6b81..4b2ea6ae 100644 --- a/src/chat/db_sync.c +++ b/src/chat/db_sync.c @@ -31,7 +31,8 @@ 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); +static int si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32_t from, uint32_t count, uint32_t vp, uint32_t want_from); +static void si_delivery_update(struct DB_SYNC_INSTANCE* si, uint64_t ts, uint64_t author, uint64_t peer_id); // ============================================================ // Structures @@ -65,6 +66,12 @@ struct DB_SYNC_INSTANCE { void* on_insert_arg; }; +struct db_sync_done_cbk_entry { + db_sync_done_fn fn; + void* arg; + struct db_sync_done_cbk_entry* next; +}; + struct DB_SYNC { struct UTUN_INSTANCE* inst; sqlite3* db; @@ -74,8 +81,40 @@ struct DB_SYNC { int instance_count, instance_capacity; void* peer_check_timer; uint8_t enabled; + struct db_sync_done_cbk_entry* done_cbks; }; +// ============================================================ +// Done callback chain (global, like etcp_status_cbk_entry) +// ============================================================ + +static void db_sync_done_cbk_add_chain(struct db_sync_done_cbk_entry** head, db_sync_done_fn fn, void* arg) { + if (!head || !fn) return; + struct db_sync_done_cbk_entry* e = u_malloc(sizeof(struct db_sync_done_cbk_entry)); + if (!e) return; + e->fn = fn; e->arg = arg; e->next = *head; + *head = e; +} +static void db_sync_done_cbk_remove_chain(struct db_sync_done_cbk_entry** head, db_sync_done_fn fn, void* arg) { + if (!head || !fn) return; + struct db_sync_done_cbk_entry** p = head; + while (*p) { + if ((*p)->fn == fn && (*p)->arg == arg) { struct db_sync_done_cbk_entry* rm = *p; *p = rm->next; u_free(rm); return; } + p = &(*p)->next; + } +} +static void si_fire_sync_done(struct DB_SYNC_INSTANCE* si, uint64_t peer_node_id) { + struct DB_SYNC* db = si->db_sync; + struct db_sync_done_cbk_entry* cb = db->done_cbks; + while (cb) { struct db_sync_done_cbk_entry* n = cb->next; cb->fn(si, peer_node_id, cb->arg); cb = n; } +} +void db_sync_add_done_cbk(struct UTUN_INSTANCE* inst, db_sync_done_fn fn, void* arg) { + if (inst && inst->db_sync) db_sync_done_cbk_add_chain(&inst->db_sync->done_cbks, fn, arg); +} +void db_sync_remove_done_cbk(struct UTUN_INSTANCE* inst, db_sync_done_fn fn, void* arg) { + if (inst && inst->db_sync) db_sync_done_cbk_remove_chain(&inst->db_sync->done_cbks, fn, arg); +} + // ============================================================ // SQL helpers // ============================================================ @@ -88,11 +127,13 @@ struct DB_SYNC { _tbuf; \ }) +// Форматирует SQL-запрос, подставляя имя таблицы из инстанса (через %s) static void si_sql(char* buf, size_t sz, struct DB_SYNC_INSTANCE* si, const char* fmt) { snprintf(buf, sz, fmt, SI_TBL(si)); } +// Готовит SQLite statement с подстановкой имени таблицы — обёртка над si_sql + sqlite3_prepare_v2 static int si_prep(struct DB_SYNC_INSTANCE* si, sqlite3_stmt** stmt, const char* fmt) { char sql[512]; @@ -104,6 +145,7 @@ static int si_prep(struct DB_SYNC_INSTANCE* si, sqlite3_stmt** stmt, const char* // SQLite open/close // ============================================================ +// Открывает SQLite БД, включает WAL, mmap и прочие оптимизации для высокой производительности static int db_sqlite_open(struct DB_SYNC* db, const char* path) { DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "path=%s", path); @@ -126,6 +168,7 @@ static int db_sqlite_open(struct DB_SYNC* db, const char* path) return 0; } +// Закрывает SQLite БД (только свою, shared не трогает) static void db_sqlite_close(struct DB_SYNC* db) { DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, ""); @@ -139,11 +182,13 @@ static void db_sqlite_close(struct DB_SYNC* db) // SHA256 helpers // ============================================================ +// Вычисляет SHA256 хеш данных (обёртка над OpenSSL SHA256) static void db_sha256(const uint8_t* data, size_t len, uint8_t hash[32]) { SHA256(data, len, hash); } +// Возвращает первые 8 байт SHA256 как uint64 — короткий идентификатор для сравнения static uint64_t db_hash64(const uint8_t* data, size_t len) { uint8_t hash[32]; db_sha256(data, len, hash); @@ -151,6 +196,7 @@ static uint64_t db_hash64(const uint8_t* data, size_t len) return h; } +// Вычисляет цепной хеш: SHA256(prev_chain_hash[32] || id[8] || timestamp[8] || author[8] || sig[64]) static void db_chain_hash_compute(const uint8_t prev_chain_hash[32], uint64_t id, uint64_t timestamp, uint64_t author, const uint8_t author_signature[DB_SIG_SIZE], @@ -169,6 +215,7 @@ static void db_chain_hash_compute(const uint8_t prev_chain_hash[32], // Instance management // ============================================================ +// Ищет активный инстанс синхронизации по 64-битному хешу static struct DB_SYNC_INSTANCE* db_instance_find(struct DB_SYNC* db, uint64_t hash) { for (int i = 0; i < db->instance_count; i++) { @@ -177,6 +224,7 @@ static struct DB_SYNC_INSTANCE* db_instance_find(struct DB_SYNC* db, uint64_t ha return NULL; } +// Выделяет новый слот инстанса в массиве, расширяя его при необходимости static struct DB_SYNC_INSTANCE* db_instance_alloc(struct DB_SYNC* db) { DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "count=%d cap=%d", db->instance_count, db->instance_capacity); @@ -194,6 +242,7 @@ static struct DB_SYNC_INSTANCE* db_instance_alloc(struct DB_SYNC* db) return si; } +// Проверяет допустимость имени таблицы: только A-Za-z0-9_, до 48 символов static int si_name_valid(const char* name) { if (!name || !name[0] || strlen(name) > 48) return 0; @@ -204,6 +253,7 @@ static int si_name_valid(const char* name) return 1; } +// Вычисляет 64-битный хеш инстанса: db_hash64(name[256] || id_be[8]) static uint64_t db_instance_hash_compute(const char* name, uint64_t id) { size_t nl = strlen(name); @@ -218,6 +268,7 @@ static uint64_t db_instance_hash_compute(const char* name, uint64_t id) // Peer management // ============================================================ +// Ищет пира в инстансе по node_id static struct SI_PEER* si_peer_find(struct DB_SYNC_INSTANCE* si, uint64_t node_id) { for (int i = 0; i < si->peer_count; i++) { @@ -226,6 +277,7 @@ static struct SI_PEER* si_peer_find(struct DB_SYNC_INSTANCE* si, uint64_t node_i return NULL; } +// Добавляет пира в инстанс (если уже существует — возвращает существующего) static struct SI_PEER* si_peer_add(struct DB_SYNC_INSTANCE* si, uint64_t node_id) { struct SI_PEER* p = si_peer_find(si, node_id); @@ -247,6 +299,7 @@ static struct SI_PEER* si_peer_add(struct DB_SYNC_INSTANCE* si, uint64_t node_id // Data access // ============================================================ +// Возвращает количество записей в таблице инстанса static uint32_t db_count(struct DB_SYNC_INSTANCE* si) { sqlite3_stmt* stmt; @@ -256,6 +309,7 @@ static uint32_t db_count(struct DB_SYNC_INSTANCE* si) return c; } +// Читает цепной хеш (32 байта) записи на позиции pos в отсортированной таблице static int db_chain_hash_at(struct DB_SYNC_INSTANCE* si, uint32_t pos, uint8_t out[32]) { sqlite3_stmt* stmt; @@ -276,6 +330,7 @@ static int db_chain_hash_at(struct DB_SYNC_INSTANCE* si, uint32_t pos, uint8_t o return 0; } +// Находит позицию записи (0-based) в отсортированной таблице по (timestamp, author_signature) static uint32_t si_find_pos(struct DB_SYNC_INSTANCE* si, uint64_t ts, const uint8_t* author_sig) { sqlite3_stmt* stmt; @@ -290,6 +345,7 @@ static uint32_t si_find_pos(struct DB_SYNC_INSTANCE* si, uint64_t ts, const uint return pos; } +// Возвращает цепной хеш записи, предшествующей заданной по (ts, author_sig) static int db_prev_chain_hash(struct DB_SYNC_INSTANCE* si, uint64_t ts, const uint8_t* author_sig, uint8_t out[32]) { sqlite3_stmt* stmt; @@ -315,6 +371,7 @@ static int db_prev_chain_hash(struct DB_SYNC_INSTANCE* si, uint64_t ts, const ui // ── Chain hash fragment (first 8 bytes) for sync protocol comparisons ── +// Возвращает первые 8 байт цепного хеша на позиции pos — для быстрого сравнения в протоколе sync static int db_chain_hash8_at(struct DB_SYNC_INSTANCE* si, uint32_t pos, uint64_t* h8) { uint8_t ch[32]; @@ -327,6 +384,7 @@ static int db_chain_hash8_at(struct DB_SYNC_INSTANCE* si, uint32_t pos, uint64_t static void db_recalc_tick(void* arg); +// Планирует асинхронный пересчёт цепных хешей начиная с позиции from_pos (батчами по 50) static void db_sync_schedule_recalc(struct DB_SYNC_INSTANCE* si, uint32_t from_pos) { if (!si->recalc.timer) { @@ -338,6 +396,7 @@ static void db_sync_schedule_recalc(struct DB_SYNC_INSTANCE* si, uint32_t from_p si->recalc.pos = from_pos; } +// Обрабатывает один батч (до 50 записей) пересчёта цепных хешей, вызывается по таймеру static void db_recalc_tick(void* arg) { struct DB_SYNC_INSTANCE* si = (struct DB_SYNC_INSTANCE*)arg; @@ -407,6 +466,7 @@ static void db_recalc_tick(void* arg) } } +// Синхронно завершает все отложенные пересчёты хешей — вызывается перед операциями, требующими актуальных хешей static void db_sync_flush_recalc(struct DB_SYNC_INSTANCE* si) { struct UASYNC* ua = si->db_sync->inst->ua; @@ -419,6 +479,7 @@ static void db_sync_flush_recalc(struct DB_SYNC_INSTANCE* si) // ── Ed25519 verification ── +// Получает Ed25519 публичный ключ узла: сначала из локального, затем из topo_node БД, затем из ETCP-соединения static int db_get_ed25519_pubkey(struct DB_SYNC* db, uint64_t node_id, uint8_t out[32]) { if (node_id == db->inst->node_id) { memcpy(out, db->inst->my_ed25519_pubkey, 32); return 0; } @@ -435,6 +496,7 @@ static int db_get_ed25519_pubkey(struct DB_SYNC* db, uint64_t node_id, uint8_t o } /* dump table: compact one-line per record with pos, id, ts, author, ch8 */ +// Отладочный дамп содержимого таблицы: одна строка на запись (pos, id, ts, author, ch8), макс. 200 записей static void si_dump_table(struct DB_SYNC_INSTANCE* si, const char* tag) { uint32_t mc = db_count(si); @@ -459,6 +521,7 @@ static void si_dump_table(struct DB_SYNC_INSTANCE* si, const char* tag) DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, " %s: table=%s mc=%u", tag, SI_TBL(si), mc); } +// Проверяет Ed25519 подпись автора: verify(pubkey, ts[8] || json, sig). Возвращает 0 если подпись верна 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, @@ -483,6 +546,7 @@ static int db_verify_author_sig(struct DB_SYNC_INSTANCE* si, return 0; } +// Вставляет запись в таблицу: проверяет подпись, дубликаты, вычисляет chain_hash, возвращает 0/1/-1/-2 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, @@ -550,16 +614,52 @@ static int db_record_insert(struct DB_SYNC_INSTANCE* si, 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; } +// Вставляет запись + каскад: on_insert + PUSH всем connected пирам (кроме источника) +static int db_record_insert_cascade(struct DB_SYNC_INSTANCE* si, + uint64_t rid, uint64_t rts, uint64_t rauthor, + const char* rdata, uint32_t rdlen, + const uint8_t* rsig, uint64_t source_peer, + const char* local_attrs) +{ + int ret = db_record_insert(si, rid, rts, rauthor, rdata, rdlen, rsig, local_attrs); + if (ret != 0) return ret; + + if (si->on_insert) si->on_insert(si, rts, rdata, rdlen, rauthor, si->on_insert_arg); + + uint8_t pbuf[4096]; uint32_t poff = 1; pbuf[0] = DB_MSG_PUSH; + memcpy(pbuf + poff, &rid, 8); poff += 8; + memcpy(pbuf + poff, &rts, 8); poff += 8; + memcpy(pbuf + poff, &rauthor, 8); poff += 8; + memcpy(pbuf + poff, &rdlen, 4); poff += 4; + if (rdlen > 0) { memcpy(pbuf + poff, rdata, rdlen); poff += rdlen; } + uint32_t sl = DB_SIG_SIZE; pbuf[poff++] = (uint8_t)sl; + memcpy(pbuf + poff, rsig, sl); poff += sl; + + uint64_t my_id = si->db_sync->inst->node_id; + int push_count = 0; + for (int i = 0; i < si->peer_count; i++) { + uint64_t pid = si->peers[i].node_id; + if (si->peers[i].sync_state >= 1 && pid != source_peer && pid != my_id) { + if (db_sync_send(si, pid, pbuf, poff) >= 0) + { si_delivery_update(si, rts, rauthor, pid); push_count++; } + } + } + if (push_count > 0) + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:ME--] → PUSH: id=%llu => %d peers", + SI_SHRT(si), (unsigned long long)rid, push_count); + return 0; +} + // ============================================================ // Send // ============================================================ +// Отправляет сообщение узлу через ETCP с заголовком: [ETCP_RT_ID_DB_SYNC:1][hash_be:8][payload]. Возвращает код etcp_send static int db_sync_send_hash(struct DB_SYNC* db, uint64_t node_id, uint64_t hash, const uint8_t* payload, size_t plen) { struct ll_entry* entry = queue_entry_new(0); @@ -587,6 +687,7 @@ static int db_sync_send_hash(struct DB_SYNC* db, uint64_t node_id, uint64_t hash return ret; } +// Отправляет сообщение через ETCP от имени конкретного инстанса (подставляет хеш инстанса) static int db_sync_send(struct DB_SYNC_INSTANCE* si, uint64_t dst_node_id, const uint8_t* payload, size_t len) { return db_sync_send_hash(si->db_sync, dst_node_id, si->hash, payload, len); @@ -596,6 +697,7 @@ static int db_sync_send(struct DB_SYNC_INSTANCE* si, uint64_t dst_node_id, const // Delivery chain helpers (local fields) // ============================================================ +// Обновляет delivery_chain (добавляет peer_id hex) и счётчик delivered_peers для записи после отправки пиру static void si_delivery_update(struct DB_SYNC_INSTANCE* si, uint64_t ts, uint64_t author, uint64_t peer_id) { char hex[17]; snprintf(hex, sizeof(hex), "%016llx", (unsigned long long)peer_id); @@ -635,6 +737,22 @@ static void si_delivery_update(struct DB_SYNC_INSTANCE* si, uint64_t ts, uint64_ // Sync protocol handlers // ============================================================ +// Определяет направление инициации синхронизации: узел с бОльшим node_id — мастер +static int is_master(uint64_t a, uint64_t b) { return a > b; } + +// Обрабатывает REQUEST_SYNC: запускает синхронизацию как слейв (мастер сам пришлёт INIT_SYNC) +static void db_handle_request_sync(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len) +{ + (void)p; (void)len; + struct SI_PEER* sp = si_peer_add(si, src); + if (sp && sp->sync_state != 1) { + sp->sync_state = 1; + sp->sync_start_tb = get_time_tb(); + db_sync_initiate_sync(si, src); + } +} + +// Обрабатывает ERROR от пира: логирует причину, сбрасывает sync_state пира static void db_handle_error(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len) { if (len < 1 || !si) return; @@ -651,6 +769,8 @@ static void db_handle_error(struct DB_SYNC_INSTANCE* si, uint64_t src, const uin } } +// Обрабатывает INIT_SYNC: сбрасывает pending recalc, определяет точку расхождения (tp), формирует INIT_RESP +// с разреженными хешами (интервалы: 1,1,1,2,2,4,4,4,4,4,4,4,4,4,4,4) и флагом has_tail 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; } @@ -702,6 +822,8 @@ static void db_handle_init_sync(struct DB_SYNC_INSTANCE* si, uint64_t src, const resp[0], resp[1]); } +// Обрабатывает INIT_RESP: сравнивает хеши на точке tp, находит первую расходящуюся позицию, +// отправляет наши данные начиная с неё и запрашивает данные пира с той же позиции 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; } @@ -720,79 +842,74 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, uint64_t src, const struct SI_PEER* sp = si_peer_find(si, src); + // Peer empty → send ALL our data, then SYNC_DONE will be sent from DATA handler when slave responds if (peer_ch8 == 0 && sc == 0) { - uint32_t vp = (uint32_t)-1; - if (sp) sp->verified_pos = vp; + uint32_t vp = (uint32_t)-1; if (sp) sp->verified_pos = vp; uint32_t sent = 0; while (sent < mc) { uint32_t b = mc - sent; if (b > DB_SEND_DATA_MAX) b = DB_SEND_DATA_MAX; - int pushed = si_send_data_batch(si, src, sent, b, vp, 1); - if (pushed <= 0) break; - sent += (uint32_t)pushed; + uint32_t want = (sent + b >= mc) ? mc : DB_WANT_FROM_NONE; + if (si_send_data_batch(si, src, sent, b, vp, want) <= 0) break; + sent += b; } - 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); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← INIT_RESP: peer empty → sent %u/%u records, awaiting response", + SI_SHRT(si), (unsigned long long)(src >> 16), sent, mc); return; } + // Exact match: same hash at tp if (my_ch8 == peer_ch8) { - uint32_t old_sp = sp ? sp->synced_pos : 0; - if (!has_tail) { + if (!has_tail && mc <= tp + 1) { 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); + 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); + si_fire_sync_done(si, src); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← INIT_RESP: MATCH+no_tail ss→2 → SYNC_DONE tp=%u", + SI_SHRT(si), (unsigned long long)(src >> 16), tp); return; } - 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); + // Has tail (slave or us): send our records after tp, request peer's after tp + uint32_t vp = tp; if (sp) sp->verified_pos = vp; + uint32_t rfrom = tp + 1; + uint32_t scnt = mc - rfrom; if (scnt > DB_SEND_DATA_MAX) scnt = DB_SEND_DATA_MAX; + if (scnt > 0) si_send_data_batch(si, src, rfrom, scnt, vp, rfrom); + else si_send_data_batch(si, src, rfrom, 0, vp, rfrom); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← INIT_RESP: MATCH+tail → SEND_DATA from=%u want=%u", + SI_SHRT(si), (unsigned long long)(src >> 16), rfrom, rfrom); return; } - uint32_t ds = 0, de = tp, fm = tp; + // Mismatch: find first divergent position + uint32_t 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); 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; if (pos < fm) fm = pos; } - } else { if (pos < fm) fm = pos; } - } - 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); - - { - 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); - } - 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); - } + if (mch8 != pch8 && pos < fm) fm = pos; + } else if (pos < fm) fm = pos; + } + 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 fm=%u vp=%u mc=%u", + SI_SHRT(si), (unsigned long long)(src >> 16), fm, vp, mc); + + // Send our records from fm, request peer's from fm + 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, fm); + else si_send_data_batch(si, src, fm, 0, vp, fm); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → SEND_DATA from=%u vp=%u cnt=%u want=%u", + SI_SHRT(si), (unsigned long long)(src >> 16), fm, vp, scnt, fm); } // ---- SEND_DATA batch helper ---- -// Wire format: [from:4][count:2][vp:4][records...] +// 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] // 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, uint32_t vp, int allocated_buf) +// Формирует и отправляет батч записей (до 32) в формате SEND_DATA: +// [from:4][count:2][vp:4][want_from:4][records...]. Возвращает количество отправленных записей +static int si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32_t from, uint32_t count, uint32_t vp, uint32_t want_from) { - uint8_t sbuf_stack[8192]; - uint8_t* buf = allocated_buf ? u_malloc(8192) : NULL; - if (!buf) buf = sbuf_stack; + uint8_t sbuf[8192]; uint8_t* buf = sbuf; uint32_t off = 0; buf[off++] = DB_MSG_SEND_DATA; @@ -800,70 +917,64 @@ static int si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32_ uint16_t rc = 0; uint16_t* rcp = (uint16_t*)(buf + off); off += 2; memcpy(buf + off, &vp, 4); off += 4; + memcpy(buf + off, &want_from, 4); off += 4; // NEW: request peer's data from this position, DB_WANT_FROM_NONE=none - sqlite3_stmt* stmt; - 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); - int bad_del = 0; - 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; - if (rsl == DB_SIG_SIZE && db_verify_author_sig(si, rts, rauth, (const char*)rd, rdl, rsig) != 0) { - uint32_t del_pos = si_find_pos(si, rts, rsig); - sqlite3_stmt* del; - if (si_prep(si, &del, "DELETE FROM \"%s\" WHERE timestamp=? AND author_signature=?") == SQLITE_OK) { - sqlite3_bind_int64(del, 1, (sqlite3_int64)rts); - sqlite3_bind_blob(del, 2, rsig, DB_SIG_SIZE, SQLITE_STATIC); - sqlite3_step(del); - sqlite3_finalize(del); - db_sync_schedule_recalc(si, del_pos); - bad_del++; + if (count > 0) { + 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) return -1; + sqlite3_bind_int64(stmt, 1, (sqlite3_int64)count); + sqlite3_bind_int64(stmt, 2, (sqlite3_int64)from); + int bad_del = 0; + 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; + if (rsl == DB_SIG_SIZE && db_verify_author_sig(si, rts, rauth, (const char*)rd, rdl, rsig) != 0) { + uint32_t del_pos = si_find_pos(si, rts, rsig); + sqlite3_stmt* del; + if (si_prep(si, &del, "DELETE FROM \"%s\" WHERE timestamp=? AND author_signature=?") == SQLITE_OK) { + sqlite3_bind_int64(del, 1, (sqlite3_int64)rts); + sqlite3_bind_blob(del, 2, rsig, DB_SIG_SIZE, SQLITE_STATIC); + sqlite3_step(del); sqlite3_finalize(del); + db_sync_schedule_recalc(si, del_pos); bad_del++; + } + continue; } - continue; + int rec_sz = 28 + rdl + 1 + rsl; + if (off + rec_sz > 8000) break; + memcpy(buf + off, &rid, 8); off += 8; + memcpy(buf + off, &rts, 8); off += 8; + memcpy(buf + off, &rauth, 8); off += 8; + memcpy(buf + off, &rdl, 4); off += 4; + if (rdl > 0) { memcpy(buf + off, rd, rdl); off += rdl; } + buf[off++] = (uint8_t)rsl; + if (rsl > 0) { memcpy(buf + off, rsig, rsl); off += rsl; } + rc++; } - int rec_sz = 28 + rdl + 1 + rsl; - if (off + rec_sz > 8000) break; - memcpy(buf + off, &rid, 8); off += 8; - memcpy(buf + off, &rts, 8); off += 8; - memcpy(buf + off, &rauth, 8); off += 8; - memcpy(buf + off, &rdl, 4); off += 4; - if (rdl > 0) { memcpy(buf + off, rd, rdl); off += rdl; } - buf[off++] = (uint8_t)rsl; - if (rsl > 0) { memcpy(buf + off, rsig, rsl); off += rsl; } - rc++; + sqlite3_finalize(stmt); + if (rc == 0 && bad_del == 0) + DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] SEND_DATA: 0 records from=%u count=%u", + SI_SHRT(si), (unsigned long long)(dst >> 16), from, count); } - sqlite3_finalize(stmt); *rcp = rc; db_sync_send(si, dst, buf, off); - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → SEND DATA: %u records from pos=%u%s", - SI_SHRT(si), (unsigned long long)(dst >> 16), rc, from, bad_del > 0 ? (rc > 0 ? " (bad_sig deleted)" : " (all bad_sig deleted)") : ""); - if (rc == 0 && count > 0 && bad_del == 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); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → SEND_DATA: %u recs from=%u want=%u vp=%u", + SI_SHRT(si), (unsigned long long)(dst >> 16), rc, from, want_from, vp); return (int)rc; } - // ---- Parse one record from SEND_DATA/PUSH wire format ---- // Wire: [id:8][ts:8][author:8][dlen:4][data][sig_len:1][sig] +// Разбирает одну запись из бинарного формата: [id:8][ts:8][author:8][dlen:4][data][sig_len:1][sig] +// Продвигает указатель *pp, возвращает 0 при успехе, -1 при ошибке парсинга static int si_parse_record(const uint8_t** pp, const uint8_t* end, uint64_t* rid, uint64_t* rts, uint64_t* rauthor, uint32_t* rdlen, const uint8_t** rdata, @@ -891,250 +1002,118 @@ 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 ---- +// ---- Atomic SEND_DATA exchange ---- +// Обрабатывает SEND_DATA: вставляет полученные записи (до 32), проверяет подписи, +// корректирует synced_pos, при want_from шлёт встречные данные 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; } + if (len < 14) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← SEND_DATA: truncated len=%zu (need >=14)", SI_SHRT(si), (unsigned long long)(src >> 16), len); return; } + db_sync_flush_recalc(si); uint32_t from = *(uint32_t*)p; uint16_t count = *(uint16_t*)(p + 4); uint32_t vp = *(uint32_t*)(p + 6); + uint32_t want_from = *(uint32_t*)(p + 10); 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); - 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); - 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; - 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; - + const uint8_t* ptr = p + 14; + uint16_t received = 0, duplicates = 0, bad_sigs = 0; 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, 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 (si_parse_record(&ptr, p + len, &rid, &rts, &rauthor, &rdlen, &rdata, &rsig, &rsiglen) != 0) break; + int ret = db_record_insert_cascade(si, rid, rts, rauthor, (const char*)rdata, rdlen, rsig, src, NULL); + if (ret >= 0) 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); } - uint16_t total_processed = count - parse_fails; - - if (total_processed > 0) { - uint32_t rst = (from < old_synced_pos) ? from : old_synced_pos; - if (sp) sp->synced_pos = from + total_processed - 1; - for (int j = 0; j < si->peer_count; j++) { + uint32_t old_spos = sp ? sp->synced_pos : 0; + if (received > 0) { + uint32_t rst = (from < old_spos) ? from : old_spos; + if (sp) sp->synced_pos = from + (count - bad_sigs) - 1; + for (int j = 0; j < si->peer_count; j++) 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+total_processed-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_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); } - // Send back own records after vp EXCEPT those just received - uint32_t mc = db_count(si); - 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); - } - *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); + char extra[32] = ""; if (duplicates) snprintf(extra, sizeof(extra), " (%u dup)", duplicates); + if (bad_sigs) { size_t el = strlen(extra); snprintf(extra+el, sizeof(extra)-el, "%s%u bad", extra[0]?", ":" (", bad_sigs); } + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← DATA: from=%u cnt=%u vp=%u want=%u → ins=%u%s mc=%u sp=%u→%u", + SI_SHRT(si), (unsigned long long)(src >> 16), from, count, vp, want_from, received, extra, db_count(si), + old_spos, sp ? sp->synced_pos : 0); + + // Respond with our data if peer requested (want_from != DB_WANT_FROM_NONE) + if (want_from != DB_WANT_FROM_NONE) { + uint32_t mc = db_count(si); + if (want_from < mc) { + uint32_t scnt = mc - want_from; if (scnt > DB_SEND_DATA_MAX) scnt = DB_SEND_DATA_MAX; + si_send_data_batch(si, src, want_from, scnt, vp, DB_WANT_FROM_NONE); + } else { + si_send_data_batch(si, src, want_from, 0, vp, DB_WANT_FROM_NONE); } + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← DATA: responded with %u recs from=%u", + SI_SHRT(si), (unsigned long long)(src >> 16), want_from < mc ? (mc - want_from > DB_SEND_DATA_MAX ? DB_SEND_DATA_MAX : mc - want_from) : 0, want_from); + } + + // If slave responded with want=DB_WANT_FROM_NONE and we are in active sync — complete the exchange with SYNC_DONE + if (want_from == DB_WANT_FROM_NONE && sp && sp->sync_state == 1) { + uint32_t mc = db_count(si); + uint64_t mch8 = 0; if (mc > 0) db_chain_hash8_at(si, mc - 1, &mch8); + uint8_t sd[13]; sd[0] = DB_MSG_SYNC_DONE; + memcpy(sd + 1, &mc, 4); memcpy(sd + 5, &mch8, 8); + db_sync_send(si, src, sd, 13); + sp->sync_state = 2; sp->sync_start_tb = 0; + si_fire_sync_done(si, src); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← DATA: exchange complete, → SYNC_DONE mc=%u", + SI_SHRT(si), (unsigned long long)(src >> 16), mc); } - 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); - db_sync_send(si, src, sd, 13); - if (sp) { uint32_t nsp = mc > 0 ? mc - 1 : 0; if (sp->synced_pos < nsp) sp->synced_pos = nsp; 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); } +// Обрабатывает SYNC_DONE: финальная сверка хешей. mc==pc → sync завершён, mcpc → отправляем хвост 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; } + if (len < 12) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "SYNC_DONE too short %zu from %016llx", len, (unsigned long long)src); 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); - uint64_t mch8 = 0; - if (mc > 0) db_chain_hash8_at(si, mc - 1, &mch8); - + uint64_t mch8 = 0; if (mc > 0) db_chain_hash8_at(si, mc - 1, &mch8); 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: 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; - } - - if (mc < pc) { - if (sp && sp->synced_pos + 1 >= pc) { - if (sp) { sp->synced_pos = pc - 1; sp->sync_state = 2; } - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, - "sync [%s:%04llX] ← SYNC_DONE: count diff my=%u peer=%u but synced=%u covers tail — accept", - SI_SHRT(si), (unsigned long long)(src >> 16), mc, pc, sp ? sp->synced_pos : 0); - return; - } - 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); + if (sp) { sp->synced_pos = mc > 0 ? mc - 1 : 0; sp->sync_state = 2; sp->sync_start_tb = 0; } + si_fire_sync_done(si, src); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← SYNC_DONE: matched ✓ mc=%u ss→2", + SI_SHRT(si), (unsigned long long)(src >> 16), mc); + } else if (mc < pc) { + // Request tail 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); + uint8_t req[14]; 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); memcpy(req + 11, &mc, 4); + db_sync_send(si, src, req, 14); if (sp) sp->sync_state = 1; - return; - } - - if (mc > pc) { - if (sp && sp->synced_pos + 1 >= mc && sp->last_peer_count == pc) { - if (sp) sp->sync_state = 2; - 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; - } + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← SYNC_DONE: mc=%u < pc=%u → request tail from=%u", + SI_SHRT(si), (unsigned long long)(src >> 16), mc, pc, mc); + } else if (mc > pc) { + // mc > pc: send tail if (sp) sp->last_peer_count = 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); + si_send_data_batch(si, src, pc, scnt, pc > 0 ? pc - 1 : (uint32_t)-1, DB_WANT_FROM_NONE); if (sp) sp->sync_state = 1; - return; + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← SYNC_DONE: mc=%u > pc=%u → send tail from=%u", + SI_SHRT(si), (unsigned long long)(src >> 16), mc, pc, pc); + } else { + // mc == pc but chain hash differs — accept as converged + if (sp) { sp->synced_pos = mc > 0 ? mc - 1 : 0; sp->sync_state = 2; sp->sync_start_tb = 0; } + si_fire_sync_done(si, src); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← SYNC_DONE: mc=%u == pc=%u (hash !match) → accept as converged", + SI_SHRT(si), (unsigned long long)(src >> 16), mc, pc); } - - 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); } +// Обрабатывает ACK_PUSH: помечает запись флагом DB_REC_FLAG_WAS_SENT и обновляет delivery_chain static void db_handle_ack_push(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len) { if (len < 16) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← PUSH ACK: truncated len=%zu (need >=16)", SI_SHRT(si), (unsigned long long)(src >> 16), len); return; } @@ -1161,6 +1140,8 @@ static void db_handle_ack_push(struct DB_SYNC_INSTANCE* si, uint64_t src, const si_delivery_update(si, ts, author, src); } +// Обрабатывает PUSH: вставляет новую запись от пира, шлёт ACK_PUSH, корректирует synced_pos остальных пиров +// если запись вставлена не в конец (в середину) static void db_handle_push(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len) { if (len < 30) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← PUSH: truncated len=%zu (need >=30)", SI_SHRT(si), (unsigned long long)(src >> 16), len); return; } @@ -1172,7 +1153,7 @@ 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, NULL); + int ret = db_record_insert_cascade(si, rid, rts, rauthor, (const char*)rdata, rdlen, rsig, src, NULL); if (ret == 0) { uint32_t ins_pos = si_find_pos(si, rts, rsig); uint32_t total = db_count(si); @@ -1180,25 +1161,17 @@ static void db_handle_push(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint 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++; } } - if (si->on_insert) si->on_insert(si, rts, (const char*)rdata, rdlen, src, si->on_insert_arg); uint8_t ack[17]; ack[0] = DB_MSG_ACK_PUSH; memcpy(ack + 1, &rts, 8); memcpy(ack + 9, &rauthor, 8); db_sync_send(si, src, ack, 17); - 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", + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← PUSH: msg id=%llu author=%016llX ts=%llu → inserted at pos=%u (of %u total)", 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, relay_count); + (unsigned long long)rid, (unsigned long long)rauthor, (unsigned long long)rts, ins_pos, total); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] ← PUSH: msg id=%llu author=%016llX ts=%llu → inserted at pos=%u (of %u total)", + SI_SHRT(si), (unsigned long long)(src >> 16), + (unsigned long long)rid, (unsigned long long)rauthor, (unsigned long long)rts, ins_pos, total); 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); @@ -1210,6 +1183,8 @@ static void db_handle_push(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint // Receive callback // ============================================================ +// Главный колбэк приёма ETCP: извлекает хеш инстанса и тип сообщения, находит/создаёт инстанс, +// маршрутизирует в соответствующий обработчик (INIT_SYNC, INIT_RESP, SEND_DATA, PUSH, ACK_PUSH, SYNC_DONE, ERROR, REQUEST_SYNC) static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { if (!entry || entry->len < 10) { if (entry) { queue_dgram_free(entry); queue_entry_free(entry); } return; } @@ -1266,7 +1241,8 @@ static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) 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; - case DB_MSG_SYNC_DONE: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ handle SYNC_DONE"); db_handle_sync_done(si, src, payload, plen); break; + case DB_MSG_REQUEST_SYNC: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ handle REQUEST_SYNC"); db_handle_request_sync(si, src, payload, plen); break; + case DB_MSG_SYNC_DONE: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ handle SYNC_DONE"); db_handle_sync_done(si, src, payload, plen); break; case DB_MSG_ERROR: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "→ handle ERROR"); db_handle_error(si, src, payload, plen); break; default: DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "sync: unknown msg type 0x%02x from %04llX (hash=%016llx) — dropped", @@ -1281,6 +1257,7 @@ static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) // Connection callbacks // ============================================================ +// Маршрутизирует статус соединения (UP/DOWN) в соответствующие обработчики static void db_sync_on_conn_status(struct ETCP_CONN* conn, int status, void* arg) { switch (status) { case ETCP_CONN_STATUS_UP: db_sync_on_conn_up(conn, arg); break; @@ -1289,6 +1266,7 @@ static void db_sync_on_conn_status(struct ETCP_CONN* conn, int status, void* arg } } +// При поднятии соединения с пиром: для всех активных инстансов добавляет пира и запускает синхронизацию static void db_sync_on_conn_up(struct ETCP_CONN* conn, void* arg) { (void)arg; @@ -1320,6 +1298,7 @@ static void db_sync_on_conn_up(struct ETCP_CONN* conn, void* arg) (unsigned long long)(pid >> 16), db->instance_count, synced > 0 ? tbl_list : "none", synced); } +// При разрыве соединения с пиром: сбрасывает sync_state для всех инстансов, если нет других активных пиров — обнуляет last_connected_tb static void db_sync_on_conn_down(struct ETCP_CONN* conn, void* arg) { (void)arg; @@ -1351,21 +1330,28 @@ cd_done: // Initiate sync // ============================================================ +// Запускает синхронизацию: мастер шлёт INIT_SYNC со своим количеством записей, слейв шлёт REQUEST_SYNC static void db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si, uint64_t pid) { uint32_t mc = db_count(si); + uint64_t my_id = si->db_sync->inst->node_id; struct SI_PEER* p = si_peer_find(si, pid); 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; - memcpy(msg + 1, &mc, 4); - db_sync_send(si, pid, msg, 5); - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] → sent INIT_SYNC: my_count=%u — \"let's compare chains\"", - SI_SHRT(si), (unsigned long long)(pid >> 16), mc); + if (is_master(my_id, pid)) { + if (p) { p->sync_start_tb = get_time_tb(); p->sync_state = 1; } + uint8_t msg[5]; + msg[0] = DB_MSG_INIT_SYNC; + memcpy(msg + 1, &mc, 4); + db_sync_send(si, pid, msg, 5); + } else { + if (p) { p->sync_state = 1; } + uint8_t msg[1] = { DB_MSG_REQUEST_SYNC }; + db_sync_send(si, pid, msg, 1); + } } +// Проверяет целостность цепочки хешей всей таблицы, при несовпадении запускает пересчёт с позиции ошибки static void db_verify_chain(struct DB_SYNC_INSTANCE* si) { uint32_t mc = db_count(si); @@ -1405,6 +1391,7 @@ static void db_verify_chain(struct DB_SYNC_INSTANCE* si) } } +// Публичная проверка целостности цепных хешей: возвращает 0 если всё верно, 1 при несовпадении, -1 при ошибке int db_sync_chain_verify(struct DB_SYNC_INSTANCE* si) { if (!si || !si->enabled) return 0; @@ -1439,12 +1426,14 @@ int db_sync_chain_verify(struct DB_SYNC_INSTANCE* si) return 0; } +// Устанавливает состояние синхронизации пира (для тестов) void db_sync_peer_set_state(struct DB_SYNC_INSTANCE* si, uint64_t node_id, uint8_t state) { struct SI_PEER* p = si_peer_find(si, node_id); if (p) p->sync_state = state; } +// Принудительно перезапускает синхронизацию с указанным пиром (для тестов) void db_sync_reinitiate(struct DB_SYNC_INSTANCE* si, uint64_t node_id) { struct SI_PEER* p = si_peer_find(si, node_id); @@ -1454,6 +1443,7 @@ void db_sync_reinitiate(struct DB_SYNC_INSTANCE* si, uint64_t node_id) } } +// Возвращает последний chain_hash8 (для сверки синхронизации в тестах) int db_sync_last_chain_hash8(struct DB_SYNC_INSTANCE* si, uint64_t* out_hash8) { if (!si || !out_hash8) return -1; @@ -1467,6 +1457,8 @@ int db_sync_last_chain_hash8(struct DB_SYNC_INSTANCE* si, uint64_t* out_hash8) // Timers // ============================================================ +// Периодическая проверка: запускает синхронизацию для пиров в состоянии sync_state==0, +// детектирует таймауты синхронизации (DB_SYNC_SYNC_TIMEOUT секунд), логирует зависшие static void db_sync_peer_check_cb(void* arg) { struct DB_SYNC* db = (struct DB_SYNC*)arg; @@ -1554,6 +1546,7 @@ static void db_sync_peer_check_cb(void* arg) db, db_sync_peer_check_cb, "db_sync_peer"); } +// Периодическая TTL-очистка: удаляет неподтверждённые записи (flags&1==0) старше db_sync_ttl от локального узла static void db_sync_instance_ttl_cb(void* arg) { struct DB_SYNC_INSTANCE* si = (struct DB_SYNC_INSTANCE*)arg; @@ -1594,6 +1587,7 @@ static void db_sync_instance_ttl_cb(void* arg) // Public API // ============================================================ +// Инициализирует модуль db_sync: открывает БД, регистрирует ETCP-биндинг, запускает таймер проверки пиров int db_sync_init(struct UTUN_INSTANCE* inst) { if (!inst) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "NULL instance"); return -1; } @@ -1637,6 +1631,7 @@ int db_sync_init(struct UTUN_INSTANCE* inst) return 0; } +// Уничтожает модуль db_sync: отменяет таймеры, закрывает БД, освобождает память void db_sync_destroy(struct UTUN_INSTANCE* inst) { if (!inst || !inst->db_sync) return; @@ -1657,12 +1652,15 @@ void db_sync_destroy(struct UTUN_INSTANCE* inst) etcp_remove_conn_status_cbk(inst, db_sync_on_conn_status, NULL); + { struct db_sync_done_cbk_entry* cb = db->done_cbks; while (cb) { struct db_sync_done_cbk_entry* n = cb->next; u_free(cb); cb = n; } } + db_sqlite_close(db); if (db->instances) u_free(db->instances); u_free(db); DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: destroyed"); } +// Создаёт SQLite-таблицу для синхронизации, регистрирует пиров, запускает первичную синхронизацию и TTL-таймер struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const char* table_name, uint64_t hash, int auto_sync) { if (!inst || !inst->db_sync || !table_name || !table_name[0]) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "instance_add invalid args"); return NULL; } @@ -1728,10 +1726,11 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const while (entry) { struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data; uint64_t pid = ce->conn->peer_node_id; - if (pid != 0 && pid != inst->node_id && ce->conn->links_up > 0 && ce->conn->initialized) { + if (pid != 0 && pid != inst->node_id && ce->conn->links_up > 0) { peers_found++; struct SI_PEER* p = si_peer_add(si, pid); - if (p && p->sync_state == 0 && auto_sync) { p->sync_state = 1; db_sync_initiate_sync(si, pid); peers_synced++; } + if (!p) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "si_peer_add failed for %N in %s", (unsigned long long)pid, SI_TBL(si)); continue; } + if (p->sync_state == 0 && auto_sync) { p->sync_state = 1; db_sync_initiate_sync(si, pid); peers_synced++; } } entry = entry->next; } @@ -1744,6 +1743,7 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const return si; } +// Удаляет инстанс: отменяет таймеры, освобождает пиров, сдвигает массив (таблица БД не удаляется) void db_sync_instance_remove(struct DB_SYNC_INSTANCE* si) { if (!si) return; @@ -1765,6 +1765,7 @@ void db_sync_instance_remove(struct DB_SYNC_INSTANCE* si) DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "instance removed tbl=%s hash=%016llx", si->table_name, (unsigned long long)si->hash); } +// Вставляет подписанную запись: проверяет подпись, вставляет в БД, рассылает PUSH всем синхронизированным пирам int db_sync_insert_signed(struct DB_SYNC_INSTANCE* si, const char* json_data, size_t len, const uint8_t* sig, size_t sig_len, uint64_t ts, const char* local_attrs) @@ -1776,52 +1777,23 @@ 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, 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); - - // Build PUSH: [DB_MSG_PUSH][id:8][ts:8][author:8][dlen:4][data][sig_len:1=64][sig:64] - uint32_t sl = DB_SIG_SIZE; - size_t push_overhead = 1 + 8 + 8 + 8 + 4 + 1 + 64; - size_t push_total = push_overhead + (len > 0 ? len : 0); - uint8_t* pbuf = (uint8_t*)u_malloc(push_total); - if (!pbuf) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "push alloc failed size=%zu", push_total); return -1; } - uint32_t off = 0; - pbuf[off++] = DB_MSG_PUSH; - memcpy(pbuf + off, &id, 8); off += 8; - memcpy(pbuf + off, &ts, 8); off += 8; - memcpy(pbuf + off, &author_node_id, 8); off += 8; - memcpy(pbuf + off, &len, 4); off += 4; - if (len > 0) { memcpy(pbuf + off, json_data, len); off += (uint32_t)len; } - pbuf[off++] = (uint8_t)sl; - memcpy(pbuf + off, sig, sl); off += sl; - - // Push to all synced peers - int push_count = 0; - for (int i = 0; i < si->peer_count; i++) { - if (si->peers[i].sync_state >= 1 && si->peers[i].node_id != si->db_sync->inst->node_id) { - if (db_sync_send(si, si->peers[i].node_id, pbuf, off) >= 0) - { si_delivery_update(si, ts, author_node_id, si->peers[i].node_id); push_count++; } - } - } - if (push_count > 0) - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync [%s:ME--] → PUSH out: new msg id=%llu ts=%llu → forwarded to %d synced peers", - SI_SHRT(si), (unsigned long long)id, (unsigned long long)ts, push_count); - u_free(pbuf); - return 0; + return db_record_insert_cascade(si, id, ts, author_node_id, json_data, (uint32_t)len, sig, author_node_id, local_attrs); } +// Возвращает количество записей в инстансе uint32_t db_sync_count(struct DB_SYNC_INSTANCE* si) { if (!si || !si->enabled) return 0; return db_count(si); } +// Возвращает последний использованный timestamp инстанса uint64_t db_sync_get_last_timestamp(struct DB_SYNC_INSTANCE* si) { return si ? si->last_timestamp_ms : 0; } +// Возвращает следующий монотонно возрастающий timestamp (мкс → мс), использует NTP если синхронизировано uint64_t db_sync_next_timestamp(struct DB_SYNC_INSTANCE* si) { if (!si) return 0; @@ -1839,6 +1811,7 @@ uint64_t db_sync_next_timestamp(struct DB_SYNC_INSTANCE* si) return nu; } +// Устанавливает колбэк, вызываемый при каждой успешной вставке записи (локальной или от пира) void db_sync_set_insert_cb(struct DB_SYNC_INSTANCE* si, db_sync_insert_cb cb, void* arg) { if (!si) return; @@ -1846,6 +1819,7 @@ void db_sync_set_insert_cb(struct DB_SYNC_INSTANCE* si, db_sync_insert_cb cb, vo si->on_insert_arg = arg; } +// Итерирует записи с offset/limit, передавая каждую в колбэк. Возвращает количество переданных записей int db_sync_select(struct DB_SYNC_INSTANCE* si, uint32_t offset, uint32_t limit, db_sync_select_cb cb, void* arg) { diff --git a/src/chat/db_sync.h b/src/chat/db_sync.h index cb26dcef..6b28ccca 100644 --- a/src/chat/db_sync.h +++ b/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 diff --git a/tests/test_db_sync.c b/tests/test_db_sync.c index d63544de..6932a5d4 100644 --- a/tests/test_db_sync.c +++ b/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)