diff --git a/doc/db_sync_v2.md b/doc/db_sync_v2.md new file mode 100644 index 00000000..7b488422 --- /dev/null +++ b/doc/db_sync_v2.md @@ -0,0 +1,292 @@ +# DB Sync Protocol v2 — Техническое задание + +## 1. Проблема + +chain_hash = SHA256(prev_chain_hash || id || ts || datahash) — хеш-цепочка, где хеш каждой записи зависит от предыдущей. + +`db_record_insert` (db_sync.c:341-463) при вставке записи в середину сортированного порядка делает каскадный пересчёт chain_hash всех последующих записей — O(N). При двустороннем обмене PUSH'ами получается O(N²). + +## 2. Решение + +**Разделить два концерна:** +- **Сравнение при синхронизации** — перейти на `datahash` (первые 8 байт SHA256(data), позиционно-независимый). Не требует каскада. +- **Целостность цепочки** — `chain_hash` остаётся в схеме БД, но используется только для однократной стартап-проверки и диагностики. + +**Синхронизация:** per-peer позиционная (`synced_pos`). При PUSH каскад не делается. При SEND_DATA батч-вставка + один каскад после батча. + +## 3. Модель данных + +### 3.1. SI_PEER + +```c +struct SI_PEER { + uint64_t node_id; + uint32_t synced_pos; // последняя подтверждённо общая позиция (0-based index) + uint8_t sync_state; // 0=not_synced, 1=syncing, 2=synced +}; +``` + +- `synced_pos` = индекс последней записи в глобальном порядке `ORDER BY timestamp, datahash`, для которой chain_hash (а после v2 — datahash) подтверждённо совпадает у обоих пиров. +- При PUSH-вставке записи на позицию P, для всех peer'ов у которых `synced_pos >= P`: `synced_pos = P - 1`. +- При SEND_DATA: `synced_pos = from + received_count - 1`. + +### 3.2. Схема БД (без изменений) + +```sql +CREATE TABLE "db_sync__" ( + timestamp INTEGER NOT NULL, + datahash INTEGER NOT NULL, + id INTEGER NOT NULL, + chain_hash BLOB NOT NULL, + author INTEGER NOT NULL, + flags INTEGER NOT NULL DEFAULT 0, + data BLOB, + author_signature BLOB, + delivered_peers INTEGER NOT NULL DEFAULT 0, + delivery_chain TEXT NOT NULL DEFAULT '', + PRIMARY KEY (timestamp, datahash) +); +``` + +## 4. Wire-формат + +### 4.1. Сообщения, которые меняются + +#### INIT_SYNC + +``` +Было: [type:1][count:4][last_chain_hash:32] = 37 байт +Стало: [type:1][count:4] = 5 байт +``` + +#### INIT_RESP + +``` +Было: [type:1][tp:4][chain_at_tp:32][sc:1][(pos:4,chain_hash:32)*sc] +Стало: [type:1][tp:4][dh_at_tp:8] [sc:1][(pos:4,dh:8)*sc] + + dh_at_tp — datahash на позиции tp (8 байт вместо 32) + sparse — datahash вместо chain_hash (8 байт на элемент вместо 32) + + При sc=0 (exact match): длина 1+4+8+1 = 14 байт (было 37) + При sc=16 sparse: длина 14 + 16*12 = 206 байт (было 37 + 16*36 = 613) +``` + +#### SYNC_DONE + +``` +Было: [type:1][count:4][last_chain_hash:32] = 37 байт +Стало: [type:1][count:4][last_dh:8] = 13 байт +``` + +### 4.2. Сообщения, которые НЕ меняются + +- **PUSH**: `[type:1][id:8][ts:8][dh:8][dlen:4][data][sig_len:1][sig]` — без изменений. `prev_ch` НЕ добавляется. +- **ACK_PUSH**: `[type:1][dh:8][ts:8]` — без изменений. +- **REFINE**: `[type:1][from:4][to:4][hc:1][(pos:4,dh:8)*hc]` — уже использует datahash, без изменений. +- **SEND_DATA**: `[type:1][from:4][count:2][(id:8,ts:8,dh:8,dlen:4,data,sig_len:1,sig)*count]` — без изменений. +- **ERROR**: `[type:1][code:1]` — без изменений. + +## 5. Новые/изменённые функции + +### 5.1. `db_cascade_from` + +```c +static void db_cascade_from(struct DB_SYNC_INSTANCE* si, uint32_t from_pos); +``` + +Выполняет каскадный пересчёт `chain_hash` для записей с позиции `from_pos` до конца таблицы. + +Алгоритм: +1. BEGIN IMMEDIATE +2. SELECT chain_hash записи на позиции `from_pos - 1` (или zero если from_pos = 0) как prev_ch +3. SELECT id, timestamp, datahash от позиции from_pos до конца +4. Для каждой: chain_hash = SHA256(prev_ch || id || ts || dh), UPDATE в БД, prev_ch = новый chain_hash +5. COMMIT + +Выделяется из текущего кода `db_record_insert` (строки 412-451) в отдельную функцию. + +### 5.2. `db_record_insert` — новый параметр + +```c +static int db_record_insert(struct DB_SYNC_INSTANCE* si, + uint64_t id, uint64_t ts, uint64_t dh, + const char* json, size_t jlen, + const uint8_t* sig, size_t sig_len, + int do_cascade); +``` + +- `do_cascade=1` — после INSERT выполняется cascade (строки 412-451). Используется для локальных вставок. +- `do_cascade=0` — без cascade. Используется для PUSH и SEND_DATA (каскад делается отдельно). + +### 5.3. `db_verify_chain` (стартап-проверка) + +```c +static void db_verify_chain(struct DB_SYNC_INSTANCE* si); +``` + +Вызывается один раз после `db_sync_instance_add` (перед инициацией sync с пирами). + +Алгоритм: +1. Проход записей 0..N-1 в порядке `ORDER BY timestamp, datahash` +2. Для каждой: chain_hash = SHA256(prev_ch || id || ts || dh) +3. Сравнить с хранимым в БД +4. При первом расхождении на позиции P: пересчитать цепочку от P до конца через `db_cascade_from(si, P)`. Завершить. + +Сложность: O(N) SHA256, однократно при старте. + +### 5.4. `si_find_pos` (новая) + +```c +static uint32_t si_find_pos(struct DB_SYNC_INSTANCE* si, + uint64_t ts, uint64_t dh); +``` + +Возвращает позицию (0-based index) записи с ключом `(ts, dh)` в глобальном порядке. Используется после вставки PUSH для корректировки `synced_pos`. + +### 5.5. `db_datahash_at` — уже существует (стр. 295) + +Возвращает `datahash` на заданной позиции. Используется для сравнения в INIT_SYNC/INIT_RESP/SYNC_DONE. + +## 6. Обработчики сообщений + +### 6.1. PUSH (`db_handle_push`) + +``` +Было: db_record_insert(si, ..., /* cascade встроен */) +Стало: db_record_insert(si, ..., /* do_cascade= */ 0) + uint32_t pos = si_find_pos(si, ts, dh); + for each peer: if peer->synced_pos >= pos: peer->synced_pos = pos - 1; +``` + +### 6.2. INIT_SYNC (`db_handle_init_sync`) + +``` +Изменения: + - Принимает [count:4] вместо [count:4][chain_hash:32] (проверка len >= 5 вместо 36) + - Сравнение: datahash вместо chain_hash (8 байт вместо 32) + - sparse: datahash вместо chain_hash (8 байт на элемент вместо 32) + - Размер sparse-элемента: 4(pos) + 8(dh) = 12 байт (было 4+32=36) + - resp буфер: 4096 → достаточно (14 + 16*12 = 206 байт при max sparse) +``` + +### 6.3. INIT_RESP (`db_handle_init_resp`) + +``` +Изменения: + - Принимает [tp:4][dh:8][sc:1][...] вместо [tp:4][ch:32][sc:1][...] + - Проверка len >= 13 вместо 37 + - Сравнение: datahash вместо chain_hash + - Пустая БД пира: zero-хеш 8 байт вместо 32 + - Sparse элементы: 12 байт вместо 36 + - При совпадении: synced_pos = tp (в дополнение к sync_state = 2) +``` + +### 6.4. SEND_DATA (`db_handle_send_data`) + +``` +Стало: + struct SI_PEER* sp = si_peer_find(si, src); + uint32_t fix_from = sp ? sp->synced_pos : 0; + + for each record: + db_record_insert(si, ..., /* do_cascade= */ 0); + + db_cascade_from(si, fix_from); + + sp->synced_pos = from + received_count - 1; + sp->sync_state = 2; +``` + +Каскад от `fix_from` (позиция, которая была synced ДО этого батча) гарантирует, что все chain_hash после этой точки пересчитаны — независимо от того, какие PUSH'и испортили их между синхронизациями. + +### 6.5. SYNC_DONE (`db_handle_sync_done`) + +``` +Изменения: + - Принимает [count:4][last_dh:8] вместо [count:4][last_chain_hash:32] (len >= 13 вместо 36) + - Сравнение последнего datahash вместо chain_hash +``` + +## 7. Механика многопировой синхронизации + +При наличии нескольких пиров с одинаковым `sync_state` (например, syncing=1) выбирается пир с минимальным `synced_pos` — у него самый старый общий префикс. После завершения его синхронизации выбор повторяется. + +Это гарантирует: двигаемся от самого старого несинхронизированного участка, не прыгая. Новые записи (в хвосте) синхронизируются последними. + +## 8. Полный сценарий: три участника A, B, C + +``` +Начальное состояние: все пусты. + +=== Шаг 1. A создаёт R1 (ts=1000) === +A: PUSH R1 → B, C. Все вставляют в конец, cascade 0. +synced_pos = 0 у всех. + +=== Шаг 2. B создаёт R2 (ts=2000), C создаёт R3 (ts=1500) === +B: [R1(1000), R2(2000)], PUSH R2 → A, C. + A: R2 в конец. C: R2 в конец. +C: [R1(1000), R2(2000)]? Нет — C создал R3(1500). + C: [R1(1000), R3(1500)] локально, потом PUSH R2 → R2.ts=2000 → в конец. + C: [R1(1000), R3(1500), R2(2000)] +C: PUSH R3 → A, B. + A: [R1(1000), R2(2000)] + R3(1500) → между R1 и R2: + A: [R1, R3, R2] ← R3 в середину, cascade ПРОПУЩЕН + B: [R1, R3, R2] ← аналогично + + synced_pos после PUSH R3 на A относительно B: + A вставил R3 на поз.1 → synced_pos_B: было 1 → стало 0 + (A считает, что совпадает с B только на позиции 0) + +=== Шаг 3. A инициирует sync с B (чемпион: min synced_pos) === +A→B: INIT_SYNC [count=3] +B→A: INIT_RESP [tp=min(3,3)-1=2, dh_at_tp=H2, sc=2, sparse:{(1,H3),(0,H1)}] +A проверяет: A@2=R2→dh=H2✓, A@1=R3→dh=H3✓, A@0=R1→dh=H1✓ +→ sync complete, synced_pos_B=2, sync_state=2 + +=== Шаг 4. Рестарт A === +db_verify_chain: + pos 0: R1.ch = SHA256(0||R1) ✓ + pos 1: R3.ch = SHA256(R1.ch||R3) ✓ (посчитан при PUSH без cascade — корректен) + pos 2: R2.ch = SHA256(R3.ch||R2) ✗ (STALE: был посчитан как SHA256(R1.ch||R2) до PUSH R3) + → расхождение на pos=2 → db_cascade_from(2): пересчёт R2.ch = SHA256(R3.ch||R2) ✓ + +=== Шаг 5. A, B, C активно обмениваются === +Каждый PUSH: вставка без cascade. +Периодические sync (по min synced_pos): SEND_DATA + один cascade от synced_pos. +synced_pos растёт. +``` + +## 9. Что НЕ трогать + +- `db_sync_insert_signed` — локальная вставка. ts всегда монотонный, запись в конец, `do_cascade=1`, каскад всегда 0 строк. +- TTL cleanup (`db_sync_instance_ttl_cb`) — удаление старых записей. Нужен cascade после удаления (уже существующая логика, не меняется). +- ACK_PUSH, ERROR, REFINE — без изменений. +- Схема БД — без изменений. +- `etcp_bind`/эпилог — без изменений. + +## 10. Сложность операций (итого) + +| Операция | Каскад | Сложность | +|----------|--------|-----------| +| Локальная вставка | 0 строк (конец) | O(1) | +| PUSH (приём) | нет | O(1) | +| SEND_DATA (батч N записей) | 1 каскад после батча | O(N + tail) | +| Стартап-проверка | 1 раз | O(total) | +| Sync-протокол (INIT_SYNC/INIT_RESP/REFINE) | нет cascade | O(log N) сообщений | + +## 11. Порядок реализации + +1. **Выделить `db_cascade_from(si, from_pos)`** из тела `db_record_insert`. +2. **Добавить параметр `do_cascade`** в `db_record_insert`. При `0` — только INSERT, cascade вызывается отдельно. +3. **Добавить `synced_pos`** в `struct SI_PEER`. Инициализировать в 0. +4. **Добавить `si_find_pos(si, ts, dh)`** — поиск позиции записи в глобальном порядке. +5. **Изменить wire-формат** INIT_SYNC, INIT_RESP, SYNC_DONE (chain_hash → datahash, 32→8 байт). +6. **Изменить обработчики**: + - `db_handle_init_sync`: новые размеры, datahash вместо chain_hash + - `db_handle_init_resp`: новые размеры, datahash вместо chain_hash, обновлять synced_pos + - `db_handle_send_data`: `do_cascade=0`, `db_cascade_from(fix_from)`, обновить synced_pos + - `db_handle_push`: `do_cascade=0`, скорректировать все synced_pos + - `db_handle_sync_done`: обновить synced_pos +7. **Добавить `db_verify_chain(si)`** — вызвать при старте после `db_sync_instance_add`. +8. **Многопировая синхронизация**: в `db_sync_peer_check_cb` выбирать пира с min synced_pos вместо первого попавшегося. diff --git a/src/db_sync.c b/src/db_sync.c index 681c96e4..f780a053 100644 --- a/src/db_sync.c +++ b/src/db_sync.c @@ -39,7 +39,8 @@ static void db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si, struct SI_PEER { uint64_t node_id; - uint8_t synced; + uint32_t synced_pos; + uint8_t sync_state; }; struct DB_SYNC_INSTANCE { @@ -108,6 +109,8 @@ static int db_sqlite_open(struct DB_SYNC* db, const char* path) DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync: WAL pragma: %s", err); sqlite3_free(err); } + sqlite3_exec(db->db, "PRAGMA synchronous=NORMAL", NULL, NULL, NULL); + sqlite3_exec(db->db, "PRAGMA wal_autocheckpoint=10000", NULL, NULL, NULL); DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: SQLite opened at %s", path); return 0; } @@ -243,7 +246,8 @@ static struct SI_PEER* si_peer_add(struct DB_SYNC_INSTANCE* si, } p = &si->peers[si->peer_count++]; p->node_id = node_id; - p->synced = 0; + p->synced_pos = 0; + p->sync_state = 0; return p; } @@ -311,6 +315,23 @@ static int db_datahash_at(struct DB_SYNC_INSTANCE* si, return 0; } +static uint32_t si_find_pos(struct DB_SYNC_INSTANCE* si, + uint64_t ts, uint64_t dh) +{ + sqlite3_stmt* stmt; + if (si_prep(si, &stmt, + "SELECT COUNT(*) FROM \"%s\"" + " WHERE timestamp 0) { + sqlite3_stmt* stmt; + if (si_prep(si, &stmt, + "SELECT chain_hash FROM \"%s\"" + " ORDER BY timestamp, datahash" + " LIMIT 1 OFFSET ?") == SQLITE_OK) + { + sqlite3_bind_int64(stmt, 1, (sqlite3_int64)(from_pos - 1)); + if (sqlite3_step(stmt) == SQLITE_ROW) { + const void* b = sqlite3_column_blob(stmt, 0); + if (b && sqlite3_column_bytes(stmt, 0) >= 32) + memcpy(prev_ch, b, 32); + } + sqlite3_finalize(stmt); + } + } + + sqlite3_stmt* sel, *upd; + if (si_prep(si, &sel, + "SELECT id,timestamp,datahash FROM \"%s\"" + " ORDER BY timestamp,datahash" + " LIMIT -1 OFFSET ?") != SQLITE_OK) + { + sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); + return; + } + sqlite3_bind_int64(sel, 1, (sqlite3_int64)from_pos); + + if (si_prep(si, &upd, + "UPDATE \"%s\" SET chain_hash=?" + " WHERE timestamp=? AND datahash=?") != SQLITE_OK) + { + sqlite3_finalize(sel); + sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); + return; + } + + uint8_t ch[32]; + while (sqlite3_step(sel) == SQLITE_ROW) { + sqlite3_int64 n_id = sqlite3_column_int64(sel, 0); + sqlite3_int64 n_ts = sqlite3_column_int64(sel, 1); + sqlite3_int64 n_dh = sqlite3_column_int64(sel, 2); + db_chain_hash_compute(prev_ch, (uint64_t)n_id, + (uint64_t)n_ts, (uint64_t)n_dh, ch); + sqlite3_reset(upd); + sqlite3_bind_blob(upd, 1, ch, 32, SQLITE_STATIC); + sqlite3_bind_int64(upd, 2, n_ts); + sqlite3_bind_int64(upd, 3, n_dh); + sqlite3_step(upd); + memcpy(prev_ch, ch, 32); + } + sqlite3_finalize(upd); + sqlite3_finalize(sel); + + rc = sqlite3_exec(db, "COMMIT", NULL, NULL, NULL); + if (rc != SQLITE_OK) + DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, + "cascade_from COMMIT: %s", sqlite3_errmsg(db)); +} + static int db_record_insert(struct DB_SYNC_INSTANCE* si, uint64_t id, uint64_t ts, uint64_t dh, const char* json, size_t jlen, - const uint8_t* sig, size_t sig_len) + const uint8_t* sig, size_t sig_len, + int do_cascade) { sqlite3* db = SI_DB(si); sqlite3_stmt* stmt; @@ -410,6 +504,7 @@ static int db_record_insert(struct DB_SYNC_INSTANCE* si, } // Recompute chain hashes for all subsequent records + if (do_cascade) { sqlite3_stmt* sel; if (si_prep(si, &sel, "SELECT timestamp,datahash,id FROM \"%s\"" @@ -451,6 +546,7 @@ static int db_record_insert(struct DB_SYNC_INSTANCE* si, } sqlite3_finalize(upd); sqlite3_finalize(sel); + } si->next_id = id + 1; rc = sqlite3_exec(db, "COMMIT", NULL, NULL, NULL); @@ -575,14 +671,13 @@ static void db_handle_init_sync(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len) { - if (len < 36) { + if (len < 4) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "INIT_SYNC too short %zu from %016llx", len, (unsigned long long)src); return; } uint32_t pc = *(uint32_t*)p; - const uint8_t* plc = p + 4; uint32_t mc = db_count(si); uint32_t tp = (pc < mc ? pc : mc); if (tp > 0) @@ -593,24 +688,11 @@ static void db_handle_init_sync(struct DB_SYNC_INSTANCE* si, resp[off++] = DB_MSG_INIT_RESP; memcpy(resp + off, &tp, 4); off += 4; - uint8_t my_ch[32]; - if (mc == 0) - memset(my_ch, 0, 32); - else if (db_chain_hash_at(si, tp, my_ch) != 0) - memset(my_ch, 0, 32); - memcpy(resp + off, my_ch, 32); off += 32; - - // Exact match — sync complete - if (memcmp(my_ch, plc, 32) == 0) { - resp[off++] = 0; - db_sync_send(si, src, resp, off); - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, - "INIT_RESP complete to %016llx tp=%u", - (unsigned long long)src, tp); - return; - } + uint64_t my_dh = 0; + if (mc > 0) + db_datahash_at(si, tp, &my_dh); + memcpy(resp + off, &my_dh, 8); off += 8; - // Sparse hashes — binary search uint32_t sp = off; off++; int scnt = 0; @@ -619,12 +701,13 @@ static void db_handle_init_sync(struct DB_SYNC_INSTANCE* si, if (tp < step) break; uint32_t pos = tp - step; - if (db_chain_hash_at(si, pos, my_ch) != 0) + uint64_t sdh; + if (db_datahash_at(si, pos, &sdh) != 0) break; - if (off + 36 > sizeof(resp)) + if (off + 12 > sizeof(resp)) break; - memcpy(resp + off, &pos, 4); off += 4; - memcpy(resp + off, my_ch, 32); off += 32; + memcpy(resp + off, &pos, 4); off += 4; + memcpy(resp + off, &sdh, 8); off += 8; scnt++; } resp[sp] = (uint8_t)scnt; @@ -638,122 +721,119 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len) { - if (len < 37) { + if (len < 13) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "INIT_RESP too short %zu from %016llx", len, (unsigned long long)src); return; } uint32_t tp = *(uint32_t*)p; - const uint8_t* pc = p + 4; - uint8_t sc = p[36]; + uint64_t peer_dh; + memcpy(&peer_dh, p + 4, 8); + uint8_t sc = p[12]; - uint8_t my_ch[32]; - if (db_chain_hash_at(si, tp, my_ch) != 0) - memset(my_ch, 0, 32); + uint64_t my_dh = 0; + db_datahash_at(si, tp, &my_dh); - // Exact match at tp — sync complete - if (sc == 0 && memcmp(my_ch, pc, 32) == 0) { + if (sc == 0 && my_dh == peer_dh) { struct SI_PEER* sp = si_peer_find(si, src); - if (sp) - sp->synced = 2; + if (sp) { + sp->synced_pos = tp; + sp->sync_state = 2; + } DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync complete with %016llx", (unsigned long long)src); return; } - // Peer has empty DB — send all records - { - uint8_t z[32]; - memset(z, 0, 32); - if (memcmp(pc, z, 32) == 0 && sc == 0) { - uint32_t mc = db_count(si); - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, - "peer %016llx empty, sending all %u", - (unsigned long long)src, mc); - uint32_t sent = 0; - while (sent < mc) { - uint32_t b = mc - sent; - if (b > DB_SEND_DATA_MAX) - b = DB_SEND_DATA_MAX; - - uint8_t sdbuf[8192]; - uint32_t off = 0; - sdbuf[off++] = DB_MSG_SEND_DATA; - memcpy(sdbuf + off, &sent, 4); off += 4; - uint16_t rc = 0; - uint16_t* rcp = (uint16_t*)(sdbuf + off); off += 2; - - sqlite3_stmt* stmt; - si_prep(si, &stmt, - "SELECT id,timestamp,datahash,data,author_signature" - " FROM \"%s\" ORDER BY timestamp,datahash" - " LIMIT ? OFFSET ?"); - sqlite3_bind_int64(stmt, 1, (sqlite3_int64)b); - sqlite3_bind_int64(stmt, 2, (sqlite3_int64)sent); - 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 rdh = (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; - - int rec_sz = 28 + rdl + 1 + rsl; - if (off + rec_sz > (int)sizeof(sdbuf)) - break; - - memcpy(sdbuf + off, &rid, 8); off += 8; - memcpy(sdbuf + off, &rts, 8); off += 8; - memcpy(sdbuf + off, &rdh, 8); off += 8; - memcpy(sdbuf + off, &rdl, 4); off += 4; - if (rdl > 0) { - memcpy(sdbuf + off, rd, rdl); off += rdl; - } - sdbuf[off++] = (uint8_t)rsl; - if (rsl > 0) { - memcpy(sdbuf + off, rsig, rsl); off += rsl; - } - rc++; - } - sqlite3_finalize(stmt); - - *rcp = rc; - sent += rc; - db_sync_send(si, src, sdbuf, off); - if (rc == 0) + if (peer_dh == 0 && sc == 0) { + uint32_t mc = db_count(si); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, + "peer %016llx empty, sending all %u", + (unsigned long long)src, mc); + uint32_t sent = 0; + while (sent < mc) { + uint32_t b = mc - sent; + if (b > DB_SEND_DATA_MAX) + b = DB_SEND_DATA_MAX; + + uint8_t sdbuf[8192]; + uint32_t off = 0; + sdbuf[off++] = DB_MSG_SEND_DATA; + memcpy(sdbuf + off, &sent, 4); off += 4; + uint16_t rc = 0; + uint16_t* rcp = (uint16_t*)(sdbuf + off); off += 2; + + sqlite3_stmt* stmt; + si_prep(si, &stmt, + "SELECT id,timestamp,datahash,data,author_signature" + " FROM \"%s\" ORDER BY timestamp,datahash" + " LIMIT ? OFFSET ?"); + sqlite3_bind_int64(stmt, 1, (sqlite3_int64)b); + sqlite3_bind_int64(stmt, 2, (sqlite3_int64)sent); + 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 rdh = (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; + + int rec_sz = 28 + rdl + 1 + rsl; + if (off + rec_sz > (int)sizeof(sdbuf)) break; + + memcpy(sdbuf + off, &rid, 8); off += 8; + memcpy(sdbuf + off, &rts, 8); off += 8; + memcpy(sdbuf + off, &rdh, 8); off += 8; + memcpy(sdbuf + off, &rdl, 4); off += 4; + if (rdl > 0) { + memcpy(sdbuf + off, rd, rdl); off += rdl; + } + sdbuf[off++] = (uint8_t)rsl; + if (rsl > 0) { + memcpy(sdbuf + off, rsig, rsl); off += rsl; + } + rc++; } - return; + sqlite3_finalize(stmt); + + *rcp = rc; + sent += rc; + db_sync_send(si, src, sdbuf, off); + if (rc == 0) + break; } + return; } - // Prefix match — sync confirmed - if (memcmp(my_ch, pc, 32) == 0) { + if (my_dh == peer_dh) { struct SI_PEER* sp = si_peer_find(si, src); - if (sp) - sp->synced = 2; + if (sp) { + sp->synced_pos = tp; + sp->sync_state = 2; + } DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, - "chain_hash match with %016llx at %u", + "dh match with %016llx at %u", (unsigned long long)src, tp); return; } - // Binary search for divergence point uint32_t ds = 0, de = tp; - const uint8_t* spr = p + 37; - for (int i = 0; i < sc && spr + 36 <= p + len; i++) { + const uint8_t* spr = p + 13; + for (int i = 0; i < sc && spr + 12 <= p + len; i++) { uint32_t pos = *(uint32_t*)spr; - const uint8_t* pch = spr + 4; - uint8_t mch[32]; - if (db_chain_hash_at(si, pos, mch) == 0) { - if (memcmp(mch, pch, 32) == 0) { + uint64_t pdh; + memcpy(&pdh, spr + 4, 8); + uint64_t mdh; + if (db_datahash_at(si, pos, &mdh) == 0) { + if (mdh == pdh) { if (pos + 1 > ds) ds = pos + 1; } @@ -762,14 +842,13 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, de = pos; } } - spr += 36; + spr += 12; } DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "divergence with %016llx range [%u,%u]", (unsigned long long)src, ds, de); if (de - ds <= 1) { - // Small divergence — REFINE with 0 hashes uint8_t ref[512]; uint32_t roff = 0; ref[roff++] = DB_MSG_REFINE; @@ -779,7 +858,6 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, db_sync_send(si, src, ref, roff); } else { - // Sparse REFINE uint8_t ref[512]; uint32_t roff = 0; ref[roff++] = DB_MSG_REFINE; @@ -952,6 +1030,10 @@ static void db_handle_send_data(struct DB_SYNC_INSTANCE* si, uint16_t count = *(uint16_t*)(p + 4); const uint8_t* ptr = p + 6; + struct SI_PEER* sp = si_peer_find(si, src); + uint32_t fix_from = sp ? sp->synced_pos : 0; + uint16_t received = 0; + for (uint16_t i = 0; i < count; i++) { uint64_t rid, rts, rdh; uint32_t rdlen; @@ -963,56 +1045,59 @@ static void db_handle_send_data(struct DB_SYNC_INSTANCE* si, break; int ret = db_record_insert(si, rid, rts, rdh, (const char*)rdata, rdlen, - rsig, rsiglen); + rsig, rsiglen, 0); + if (ret >= 0) + received++; if (ret == 0 && si->on_insert) si->on_insert(si, (const char*)rdata, rdlen, src, si->on_insert_arg); } + db_cascade_from(si, fix_from); + + if (sp && received > 0) + sp->synced_pos = from + received - 1; + uint32_t mc = db_count(si); - struct SI_PEER* sp = si_peer_find(si, src); - uint32_t pk = from + count; + uint32_t pk = from + received; - // Request next batch if more records expected - if (mc > pk && sp && sp->synced == 1) { + 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); } - // Send SYNC_DONE when we've received all uint32_t nc = db_count(si); - uint8_t lch[32]; - if (nc > 0 && db_chain_hash_at(si, nc - 1, lch) == 0) { - uint8_t sd[37]; - sd[0] = DB_MSG_SYNC_DONE; - memcpy(sd + 1, &nc, 4); - memcpy(sd + 5, lch, 32); - db_sync_send(si, src, sd, 37); - } + uint64_t ldh = 0; + if (nc > 0) + db_datahash_at(si, nc - 1, &ldh); + uint8_t sd[13]; + sd[0] = DB_MSG_SYNC_DONE; + memcpy(sd + 1, &nc, 4); + memcpy(sd + 5, &ldh, 8); + db_sync_send(si, src, sd, 13); if (sp) - sp->synced = 2; + sp->sync_state = 2; DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "received %u records from %016llx", - count, (unsigned long long)src); + received, (unsigned long long)src); } static void db_handle_sync_done(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint8_t* p, size_t len) { - if (len < 36) + if (len < 12) return; uint32_t pc = *(uint32_t*)p; - const uint8_t* pch = p + 4; + uint64_t pdh; + memcpy(&pdh, p + 4, 8); uint32_t mc = db_count(si); - uint8_t mch[32]; + uint64_t mdh = 0; if (mc > 0) - db_chain_hash_at(si, mc - 1, mch); - else - memset(mch, 0, 32); + db_datahash_at(si, mc - 1, &mdh); - if (mc != pc || memcmp(mch, pch, 32) != 0) { + if (mc != pc || mdh != pdh) { DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "SYNC_DONE mismatch %016llx my=%u peer=%u", (unsigned long long)src, mc, pc); @@ -1020,8 +1105,10 @@ static void db_handle_sync_done(struct DB_SYNC_INSTANCE* si, return; } struct SI_PEER* sp = si_peer_find(si, src); - if (sp) - sp->synced = 2; + if (sp) { + sp->synced_pos = (mc < pc ? mc : pc) > 0 ? (mc < pc ? mc : pc) - 1 : 0; + sp->sync_state = 2; + } DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "sync confirmed with %016llx count=%u", (unsigned long long)src, mc); @@ -1075,8 +1162,13 @@ static void db_handle_push(struct DB_SYNC_INSTANCE* si, uint64_t src, int ret = db_record_insert(si, rid, rts, rdh, (const char*)rdata, rdlen, - rsig, rsiglen); + rsig, rsiglen, 0); if (ret == 0) { + uint32_t ins_pos = si_find_pos(si, rts, rdh); + 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; + } if (si->on_insert) si->on_insert(si, (const char*)rdata, rdlen, src, si->on_insert_arg); @@ -1206,8 +1298,8 @@ static void db_sync_on_conn_up(struct ETCP_CONN* conn, void* arg) if (!si->enabled) continue; struct SI_PEER* p = si_peer_add(si, pid); - if (p && p->synced == 0 && conn->initialized && conn->links_up) { - p->synced = 1; + if (p && p->sync_state == 0 && conn->initialized && conn->links_up) { + p->sync_state = 1; db_sync_initiate_sync(si, pid); } } @@ -1231,13 +1323,13 @@ static void db_sync_on_conn_down(struct ETCP_CONN* conn, void* arg) for (int i = 0; i < db->instance_count; i++) { struct SI_PEER* p = si_peer_find(&db->instances[i], pid); if (p) - p->synced = 0; + p->sync_state = 0; } int any = 0; for (int i = 0; i < db->instance_count; i++) { for (int j = 0; j < db->instances[i].peer_count; j++) { - if (db->instances[i].peers[j].synced >= 1) { + if (db->instances[i].peers[j].sync_state >= 1) { any = 1; goto cd_done; } @@ -1262,22 +1354,57 @@ static void db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si, "initiate_sync to %016llx my=%u", (unsigned long long)pid, mc); - uint8_t msg[37]; + uint8_t msg[5]; msg[0] = DB_MSG_INIT_SYNC; memcpy(msg + 1, &mc, 4); - if (mc > 0) { - if (db_chain_hash_at(si, mc - 1, msg + 5) != 0) - memset(msg + 5, 0, 32); - } - else { - memset(msg + 5, 0, 32); - } - db_sync_send(si, pid, msg, 37); + db_sync_send(si, pid, msg, 5); DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "INIT_SYNC -> %016llx my=%u", (unsigned long long)pid, mc); } +static void db_verify_chain(struct DB_SYNC_INSTANCE* si) +{ + uint32_t mc = db_count(si); + if (mc == 0) + return; + + uint8_t prev_ch[32]; memset(prev_ch, 0, 32); + uint8_t exp_ch[32], stored_ch[32]; + int bad_pos = -1; + + sqlite3_stmt* stmt; + if (si_prep(si, &stmt, + "SELECT id,timestamp,datahash,chain_hash FROM \"%s\"" + " ORDER BY timestamp,datahash") != SQLITE_OK) + return; + 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 rdh = (uint64_t)sqlite3_column_int64(stmt, 2); + const void* b = sqlite3_column_blob(stmt, 3); + if (b && sqlite3_column_bytes(stmt, 3) >= 32) + memcpy(stored_ch, b, 32); + else + memset(stored_ch, 0, 32); + + db_chain_hash_compute(prev_ch, rid, rts, rdh, exp_ch); + if (memcmp(exp_ch, stored_ch, 32) != 0) { + bad_pos = (int)pos; + break; + } + memcpy(prev_ch, exp_ch, 32); + } + sqlite3_finalize(stmt); + + if (bad_pos >= 0) { + DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, + "chain_hash mismatch at pos %d/%u in %s, recalculating", + bad_pos, mc, SI_TBL(si)); + db_cascade_from(si, (uint32_t)bad_pos); + } +} + // ============================================================ // Timers // ============================================================ @@ -1312,16 +1439,24 @@ static void db_sync_peer_check_cb(void* arg) && item->conn->links_up && item->conn->initialized) { uint64_t pid = item->conn->peer_node_id; - if (pid != db->inst->node_id) { - struct SI_PEER* p = si_peer_add(si, pid); - if (p && p->synced == 0) { - p->synced = 1; - db_sync_initiate_sync(si, pid); - } - } + if (pid != db->inst->node_id) + si_peer_add(si, pid); } e = e->next; } + + struct SI_PEER* best = NULL; + uint32_t min_pos = UINT32_MAX; + for (int j = 0; j < si->peer_count; j++) { + if (si->peers[j].sync_state == 0 && si->peers[j].synced_pos < min_pos) { + min_pos = si->peers[j].synced_pos; + best = &si->peers[j]; + } + } + if (best) { + best->sync_state = 1; + db_sync_initiate_sync(si, best->node_id); + } } db->peer_check_timer = @@ -1581,6 +1716,8 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, sqlite3_finalize(stmt); } + db_verify_chain(si); + si->ttl_timer = uasync_set_timeout(inst->ua, DB_SYNC_TTL_INTERVAL * 10000u, @@ -1597,8 +1734,8 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, && ce->conn->links_up > 0 && ce->conn->initialized) { struct SI_PEER* p = si_peer_add(si, pid); - if (p && p->synced == 0) { - p->synced = 1; + if (p && p->sync_state == 0) { + p->sync_state = 1; db_sync_initiate_sync(si, pid); } } @@ -1680,7 +1817,7 @@ int db_sync_insert_signed(struct DB_SYNC_INSTANCE* si, uint64_t id = si->next_id; int ret = db_record_insert(si, id, ts, dh, - json_data, len, sig, sig_len); + json_data, len, sig, sig_len, 1); if (ret != 0) return ret; if (si->on_insert) @@ -1707,7 +1844,7 @@ int db_sync_insert_signed(struct DB_SYNC_INSTANCE* si, // Push to all synced peers for (int i = 0; i < si->peer_count; i++) { - if (si->peers[i].synced >= 1 + 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) diff --git a/src/utun_instance.c b/src/utun_instance.c index 2cd4f851..a01a93a5 100644 --- a/src/utun_instance.c +++ b/src/utun_instance.c @@ -412,7 +412,10 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) { } instance->etcp_sockets = NULL; DEBUG_INFO(DEBUG_CATEGORY_MEMORY, "[INSTANCE_DESTROY] ETCP sockets cleanup complete"); - + + // Cleanup db_sync before connections cleanup (db_sync_destroy iterates connections) + db_sync_destroy(instance); + // Cleanup ETCP connections (phase 1 detach + deferred phase 2 via call_soon) { struct ll_entry* entry = instance->connections->head; @@ -469,9 +472,6 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) { nat_transport_destroy(instance); } - // Cleanup db_sync (before etcp_router_destroy) - db_sync_destroy(instance); - // Cleanup etcp_router etcp_router_destroy(instance); diff --git a/tests/Makefile.am b/tests/Makefile.am index 7ac20371..f664ba3f 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -50,6 +50,7 @@ check_PROGRAMS = \ test_conn_mgr \ test_etcp_connect \ test_db_sync \ + test_chat_sync_stress \ test_stcp_traffic \ test_bbr_integration \ test_intensive_memory_pool \ @@ -274,6 +275,10 @@ test_db_sync_SOURCES = test_db_sync.c test_db_sync_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_db_sync_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_chat_sync_stress_SOURCES = test_chat_sync_stress.c +test_chat_sync_stress_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib +test_chat_sync_stress_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) + test_stcp_traffic_SOURCES = test_stcp_traffic.c test_stcp_traffic_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_stcp_traffic_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) diff --git a/tests/test_chat_sync_stress.c b/tests/test_chat_sync_stress.c new file mode 100644 index 00000000..45603f4a --- /dev/null +++ b/tests/test_chat_sync_stress.c @@ -0,0 +1,246 @@ +// test_chat_sync_stress.c — 10 раундов: 100 батчей вставок + 1 синхронизация за раунд +#include +#include +#include +#include +#include +#include "../lib/platform_compat.h" +#include "test_utils.h" +#ifdef _WIN32 +#include +#include +#else +#include +#endif + +#include "../src/utun_instance.h" +#include "../src/config_parser.h" +#include "../src/config_updater.h" +#include "../src/db_sync.h" +#include "../lib/u_async.h" +#include "../lib/debug_config.h" + +#define ROUNDS 10 +#define BATCHES_PER_ROUND 100 +#define MAX_PER_BATCH 100 +#define SYNC_TIMEOUT_TB 300000 // 30s +#define TOTAL_TIMEOUT_TB (ROUNDS * SYNC_TIMEOUT_TB + 600000) +#define POLL_INTERVAL_MS 5 + +#define NODE_ID_A 0xAAAAAAAAAAAAAAAAULL +#define NODE_ID_B 0xBBBBBBBBBBBBBBBBULL +#define INSTANCE_NAME "chat" +#define INSTANCE_ID 1 + +static struct UTUN_INSTANCE* inst_a = NULL; +static struct UTUN_INSTANCE* inst_b = NULL; +static struct DB_SYNC_INSTANCE* si_a = NULL; +static struct DB_SYNC_INSTANCE* si_b = NULL; +static struct UASYNC* ua = NULL; +static int test_phase = 0; +static void* timeout_id = NULL; +static uint32_t expected_total = 0; +static char log_path[320]; +static unsigned int g_seed; +static char temp_dir[] = "/tmp/utun_chatsync_XXXXXX"; +static char config_a[256], config_b[256]; +static int port_a_srv, port_b_srv; + +static int write_file(const char* path, const char* fmt, ...) { + va_list ap; FILE* f = fopen(path, "w"); if (!f) return -1; + va_start(ap, fmt); vfprintf(f, fmt, ap); va_end(ap); fclose(f); return 0; +} +static char* get_pubkey(const char* path) { + struct utun_config* cfg = parse_config(path); if (!cfg) return NULL; + char* pub = strdup(cfg->global.my_public_key_hex); free_config(cfg); return pub; +} +static int create_temp_configs(void) { + if (test_mkdtemp(temp_dir) != 0) { return -1; } + int base = 42000 + (getpid() % 15000); + port_a_srv = base; port_b_srv = base + 1; + snprintf(config_a, sizeof(config_a), "%s/a.conf", temp_dir); + snprintf(config_b, sizeof(config_b), "%s/b.conf", temp_dir); + char db_path[320]; + snprintf(db_path, sizeof(db_path), "%s/db_a", temp_dir); utun_mkdir(db_path, 0755); + snprintf(db_path, sizeof(db_path), "%s/db_b", temp_dir); utun_mkdir(db_path, 0755); + snprintf(log_path, sizeof(log_path), "%s/test.log", temp_dir); + if (write_file(config_a, + "[global]\nmy_node_id=0xAAAAAAAAAAAAAAAA\ntun_ip=10.200.0.1/24\ntun_ifname=tun200\n" + "db_path=%s/db_a\ndb_sync_enabled=1\n\n" + "[server: srv_a]\naddr=127.0.0.1:%d\ntype=public\n\n[allowed_keys]\nallow_all=1\n", + temp_dir, port_a_srv) != 0) return -1; + if (config_ensure_keys_and_node_id(config_a) != 0) return -1; + char* pub_a = get_pubkey(config_a); if (!pub_a) return -1; + if (write_file(config_b, + "[global]\nmy_node_id=0xBBBBBBBBBBBBBBBB\ntun_ip=10.200.0.2/24\ntun_ifname=tun201\n" + "db_path=%s/db_b\ndb_sync_enabled=1\n\n" + "[server: srv_b]\naddr=127.0.0.1:%d\ntype=public\n\n" + "[client: to_a]\nkeepalive=1\npeer_public_key=%s\nlink=srv_b:127.0.0.1:%d\n\n" + "[allowed_keys]\nallow_all=1\n", + temp_dir, port_b_srv, pub_a, port_a_srv) != 0) { free(pub_a); return -1; } + free(pub_a); + if (config_ensure_keys_and_node_id(config_b) != 0) return -1; + return 0; +} +static void cleanup_temp_configs(void) { + unlink(config_a); unlink(config_b); unlink(log_path); + char pa[320]; + snprintf(pa, sizeof(pa), "%s/db_a/sync", temp_dir); unlink(pa); + snprintf(pa, sizeof(pa), "%s/db_a/sync-wal", temp_dir); unlink(pa); + snprintf(pa, sizeof(pa), "%s/db_a/sync-shm", temp_dir); unlink(pa); + snprintf(pa, sizeof(pa), "%s/db_b/sync", temp_dir); unlink(pa); + snprintf(pa, sizeof(pa), "%s/db_b/sync-wal", temp_dir); unlink(pa); + snprintf(pa, sizeof(pa), "%s/db_b/sync-shm", temp_dir); unlink(pa); + snprintf(pa, sizeof(pa), "%s/db_a", temp_dir); test_rmdir(pa); + snprintf(pa, sizeof(pa), "%s/db_b", temp_dir); test_rmdir(pa); + test_rmdir(temp_dir); +} + +static void test_timeout(void* arg) { (void)arg; test_phase = 2; } + +static int wait_for(const char* desc, int timeout_tb) { + uint64_t start = get_time_tb(); + uint32_t ca0 = si_a ? db_sync_count(si_a) : 0, cb0 = si_b ? db_sync_count(si_b) : 0; + DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "SYNC_START %s A=%u B=%u expect=%u", desc, ca0, cb0, expected_total); + while ((db_sync_count(si_a) != expected_total || db_sync_count(si_b) != expected_total) + && (get_time_tb() - start) < (uint64_t)timeout_tb && test_phase == 0) + uasync_poll(ua, POLL_INTERVAL_MS); + uint32_t ca = si_a ? db_sync_count(si_a) : 0, cb = si_b ? db_sync_count(si_b) : 0; + uint64_t elapsed = (get_time_tb() - start) / 10; + int ok = (ca == expected_total && cb == expected_total); + if (ok) + DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "SYNC_DONE %s A=%u B=%u %llums", desc, ca, cb, (unsigned long long)elapsed); + else { + DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "SYNC_FAIL %s A=%u B=%u expect=%u %llums", + desc, ca, cb, expected_total, (unsigned long long)elapsed); + if (test_phase == 0) test_phase = 2; + } + return ok; +} + +static int insert_batch(struct DB_SYNC_INSTANCE* si, int seq, int count, uint64_t node_id) { + if (count == 0) return 0; + uint64_t t0 = get_time_tb(); + char buf[384]; + for (int i = 0; i < count && test_phase == 0; i++) { + snprintf(buf, sizeof(buf), + "{\"seq\":%d,\"n\":%llu,\"ch\":\"test\",\"ct\":\"text\",\"d\":\"msg_%d\"," + "\"pad\":\"xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx\"}", + seq + i, (unsigned long long)node_id, seq + i); + if (db_sync_insert(si, buf) < 0) { + DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "INSERT_FAIL seq=%d", seq + i); + return -1; + } + } + uint64_t dt = (get_time_tb() - t0) / 10; + if (dt > 50) DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "INSERT_BATCH %d recs %llums", count, (unsigned long long)dt); + return 0; +} + +int main(void) { + g_seed = (unsigned int)time(NULL); + srand(g_seed); + + debug_config_init(); + debug_set_category_level(DEBUG_CATEGORY_DEBUG, DEBUG_LEVEL_WARN); + + if (create_temp_configs() != 0) { fprintf(stderr, "FAIL: config creation\n"); return 1; } + debug_enable_file_output(log_path, 1); + DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "START seed=%u temp=%s", g_seed, temp_dir); + + utun_instance_set_tun_init_enabled(0); + ua = uasync_create(); + if (!ua) { DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "uasync_create"); cleanup_temp_configs(); return 1; } + inst_a = utun_instance_create(ua, config_a); + inst_b = utun_instance_create(ua, config_b); + if (!inst_a || !inst_b) { DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "instance_create"); cleanup_temp_configs(); return 1; } + if (utun_instance_init(inst_a) != 0 || utun_instance_init(inst_b) != 0) { + DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "instance_init"); cleanup_temp_configs(); return 1; + } + si_a = db_sync_instance_add(inst_a, INSTANCE_NAME, INSTANCE_ID); + si_b = db_sync_instance_add(inst_b, INSTANCE_NAME, INSTANCE_ID); + if (!si_a || !si_b) { DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "db_sync_instance_add"); return 1; } + timeout_id = uasync_set_timeout(ua, TOTAL_TIMEOUT_TB, NULL, test_timeout, "total_timeout"); + + uint64_t t0 = get_time_tb(); + + // Phase 1: local + if (db_sync_insert(si_a, "{\"test\":1}") != 0 + || db_sync_insert(si_a, "{\"test\":2}") != 0 + || db_sync_count(si_a) != 2) { + DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "PHASE1_FAIL A=%u", db_sync_count(si_a)); + test_phase = 2; goto cleanup; + } + DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "PHASE1_OK A=%u", db_sync_count(si_a)); + + // Phase 2: initial sync + expected_total = 2; + if (!wait_for("phase2", SYNC_TIMEOUT_TB)) { test_phase = 2; goto cleanup; } + DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "PHASE2_OK A=%u B=%u", db_sync_count(si_a), db_sync_count(si_b)); + + // Stress rounds + int global_seq = 3; + int rounds_ok = 0; + + for (int r = 0; r < ROUNDS && test_phase == 0; r++) { + uint32_t added = 0, cnt_a = 0, cnt_b = 0; + + DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "ROUND%d_INSERT A=%u B=%u expect=%u", + r, db_sync_count(si_a), db_sync_count(si_b), expected_total); + + for (int b = 0; b < BATCHES_PER_ROUND && test_phase == 0; b++) { + int n = rand() % (MAX_PER_BATCH + 1); + if (n == 0) continue; + int side = rand() & 1; + struct DB_SYNC_INSTANCE* si = side ? si_b : si_a; + uint64_t node = side ? NODE_ID_B : NODE_ID_A; + + if (insert_batch(si, global_seq, n, node) < 0) { test_phase = 2; break; } + global_seq += n; + added += n; + expected_total += n; + if (side) cnt_b += n; else cnt_a += n; + if ((added & 127) == 127) uasync_poll(ua, POLL_INTERVAL_MS); + } + if (test_phase != 0) break; + + DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "ROUND%d_ADDED added=%u A=%u B=%u total=%u", + r, added, cnt_a, cnt_b, expected_total); + + uint64_t ts = get_time_tb(); + int synced = wait_for("round_sync", SYNC_TIMEOUT_TB); + uint64_t elapsed = (get_time_tb() - ts) / 10; + + uint32_t ca = db_sync_count(si_a), cb = db_sync_count(si_b); + if (!synced || ca != expected_total || cb != expected_total) { + DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "ROUND%d_FAIL A=%u B=%u expect=%u sync=%llums", + r, ca, cb, expected_total, (unsigned long long)elapsed); + fprintf(stderr, "FAIL round %d: A=%u B=%u expected=%u\n", r, ca, cb, expected_total); + test_phase = 2; break; + } + rounds_ok++; + DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "ROUND%d_OK sync=%llums A=%u B=%u", + r, (unsigned long long)elapsed, ca, cb); + } + + uint64_t total_ms = (get_time_tb() - t0) / 10; + if (test_phase == 0) + DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "ALL_OK rounds=%d A=%u B=%u total=%llums", + rounds_ok, db_sync_count(si_a), db_sync_count(si_b), (unsigned long long)total_ms); + +cleanup: + DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "CLEANUP A=%u B=%u phase=%d", + si_a ? db_sync_count(si_a) : 0, si_b ? db_sync_count(si_b) : 0, test_phase); + if (timeout_id && ua) { uasync_cancel_timeout(ua, timeout_id); timeout_id = NULL; } + if (si_a) { db_sync_instance_remove(si_a); si_a = NULL; } + if (si_b) { db_sync_instance_remove(si_b); si_b = NULL; } + if (inst_a) { inst_a->running = 0; utun_instance_destroy(inst_a); inst_a = NULL; } + if (inst_b) { inst_b->running = 0; utun_instance_destroy(inst_b); inst_b = NULL; } + if (ua) { uasync_destroy(ua, 0); ua = NULL; } + debug_disable_file_output(); + + int result = (test_phase == 0) ? 0 : 1; + fprintf(stderr, "=== %s seed=%u (log: %s) ===\n", result ? "FAIL" : "PASS", g_seed, log_path); + cleanup_temp_configs(); + return result; +}