Browse Source

db_sync v2: positional sync per-peer (synced_pos), datahash instead of chain_hash, no cascade on PUSH, single cascade after SEND_DATA batch, startup chain verification. PRAGMA synchronous=NORMAL + wal_autocheckpoint=10000 for performance.

topo_upd
Evgeny 3 months ago
parent
commit
bb425c6722
  1. 292
      doc/db_sync_v2.md
  2. 471
      src/db_sync.c
  3. 6
      src/utun_instance.c
  4. 5
      tests/Makefile.am
  5. 246
      tests/test_chat_sync_stress.c

292
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_<name>_<id>" (
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 вместо первого попавшегося.

471
src/db_sync.c

@ -39,7 +39,8 @@ static void db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si,
struct SI_PEER { struct SI_PEER {
uint64_t node_id; uint64_t node_id;
uint8_t synced; uint32_t synced_pos;
uint8_t sync_state;
}; };
struct DB_SYNC_INSTANCE { 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); DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync: WAL pragma: %s", err);
sqlite3_free(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); DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: SQLite opened at %s", path);
return 0; 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 = &si->peers[si->peer_count++];
p->node_id = node_id; p->node_id = node_id;
p->synced = 0; p->synced_pos = 0;
p->sync_state = 0;
return p; return p;
} }
@ -311,6 +315,23 @@ static int db_datahash_at(struct DB_SYNC_INSTANCE* si,
return 0; 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<?1"
" OR (timestamp=?1 AND datahash<?2)") != SQLITE_OK)
return 0;
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)ts);
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)dh);
uint32_t pos = (sqlite3_step(stmt) == SQLITE_ROW)
? (uint32_t)sqlite3_column_int64(stmt, 0) : 0;
sqlite3_finalize(stmt);
return pos;
}
static int db_prev_chain_hash(struct DB_SYNC_INSTANCE* si, static int db_prev_chain_hash(struct DB_SYNC_INSTANCE* si,
uint64_t ts, uint64_t dh, uint8_t out[32]) uint64_t ts, uint64_t dh, uint8_t out[32])
{ {
@ -338,10 +359,83 @@ static int db_prev_chain_hash(struct DB_SYNC_INSTANCE* si,
return 0; return 0;
} }
static void db_cascade_from(struct DB_SYNC_INSTANCE* si, uint32_t from_pos)
{
sqlite3* db = SI_DB(si);
int rc = sqlite3_exec(db, "BEGIN IMMEDIATE", NULL, NULL, NULL);
if (rc != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC,
"cascade_from BEGIN: %s", sqlite3_errmsg(db));
return;
}
uint8_t prev_ch[32];
memset(prev_ch, 0, 32);
if (from_pos > 0) {
sqlite3_stmt* stmt;
if (si_prep(si, &stmt,
"SELECT chain_hash FROM \"%s\""
" ORDER BY timestamp, 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, static int db_record_insert(struct DB_SYNC_INSTANCE* si,
uint64_t id, uint64_t ts, uint64_t dh, uint64_t id, uint64_t ts, uint64_t dh,
const char* json, size_t jlen, 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* db = SI_DB(si);
sqlite3_stmt* stmt; sqlite3_stmt* stmt;
@ -410,6 +504,7 @@ static int db_record_insert(struct DB_SYNC_INSTANCE* si,
} }
// Recompute chain hashes for all subsequent records // Recompute chain hashes for all subsequent records
if (do_cascade) {
sqlite3_stmt* sel; sqlite3_stmt* sel;
if (si_prep(si, &sel, if (si_prep(si, &sel,
"SELECT timestamp,datahash,id FROM \"%s\"" "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(upd);
sqlite3_finalize(sel); sqlite3_finalize(sel);
}
si->next_id = id + 1; si->next_id = id + 1;
rc = sqlite3_exec(db, "COMMIT", NULL, NULL, NULL); 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, uint64_t src,
const uint8_t* p, size_t len) const uint8_t* p, size_t len)
{ {
if (len < 36) { if (len < 4) {
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC,
"INIT_SYNC too short %zu from %016llx", "INIT_SYNC too short %zu from %016llx",
len, (unsigned long long)src); len, (unsigned long long)src);
return; return;
} }
uint32_t pc = *(uint32_t*)p; uint32_t pc = *(uint32_t*)p;
const uint8_t* plc = p + 4;
uint32_t mc = db_count(si); uint32_t mc = db_count(si);
uint32_t tp = (pc < mc ? pc : mc); uint32_t tp = (pc < mc ? pc : mc);
if (tp > 0) if (tp > 0)
@ -593,24 +688,11 @@ static void db_handle_init_sync(struct DB_SYNC_INSTANCE* si,
resp[off++] = DB_MSG_INIT_RESP; resp[off++] = DB_MSG_INIT_RESP;
memcpy(resp + off, &tp, 4); off += 4; memcpy(resp + off, &tp, 4); off += 4;
uint8_t my_ch[32]; uint64_t my_dh = 0;
if (mc == 0) if (mc > 0)
memset(my_ch, 0, 32); db_datahash_at(si, tp, &my_dh);
else if (db_chain_hash_at(si, tp, my_ch) != 0) memcpy(resp + off, &my_dh, 8); off += 8;
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;
}
// Sparse hashes — binary search
uint32_t sp = off; uint32_t sp = off;
off++; off++;
int scnt = 0; int scnt = 0;
@ -619,12 +701,13 @@ static void db_handle_init_sync(struct DB_SYNC_INSTANCE* si,
if (tp < step) if (tp < step)
break; break;
uint32_t pos = tp - step; 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; break;
if (off + 36 > sizeof(resp)) if (off + 12 > sizeof(resp))
break; break;
memcpy(resp + off, &pos, 4); off += 4; memcpy(resp + off, &pos, 4); off += 4;
memcpy(resp + off, my_ch, 32); off += 32; memcpy(resp + off, &sdh, 8); off += 8;
scnt++; scnt++;
} }
resp[sp] = (uint8_t)scnt; resp[sp] = (uint8_t)scnt;
@ -638,122 +721,119 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si,
uint64_t src, uint64_t src,
const uint8_t* p, size_t len) const uint8_t* p, size_t len)
{ {
if (len < 37) { if (len < 13) {
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC,
"INIT_RESP too short %zu from %016llx", "INIT_RESP too short %zu from %016llx",
len, (unsigned long long)src); len, (unsigned long long)src);
return; return;
} }
uint32_t tp = *(uint32_t*)p; uint32_t tp = *(uint32_t*)p;
const uint8_t* pc = p + 4; uint64_t peer_dh;
uint8_t sc = p[36]; memcpy(&peer_dh, p + 4, 8);
uint8_t sc = p[12];
uint8_t my_ch[32]; uint64_t my_dh = 0;
if (db_chain_hash_at(si, tp, my_ch) != 0) db_datahash_at(si, tp, &my_dh);
memset(my_ch, 0, 32);
// Exact match at tp — sync complete if (sc == 0 && my_dh == peer_dh) {
if (sc == 0 && memcmp(my_ch, pc, 32) == 0) {
struct SI_PEER* sp = si_peer_find(si, src); struct SI_PEER* sp = si_peer_find(si, src);
if (sp) if (sp) {
sp->synced = 2; sp->synced_pos = tp;
sp->sync_state = 2;
}
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"sync complete with %016llx", "sync complete with %016llx",
(unsigned long long)src); (unsigned long long)src);
return; return;
} }
// Peer has empty DB — send all records if (peer_dh == 0 && sc == 0) {
{ uint32_t mc = db_count(si);
uint8_t z[32]; DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
memset(z, 0, 32); "peer %016llx empty, sending all %u",
if (memcmp(pc, z, 32) == 0 && sc == 0) { (unsigned long long)src, mc);
uint32_t mc = db_count(si); uint32_t sent = 0;
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, while (sent < mc) {
"peer %016llx empty, sending all %u", uint32_t b = mc - sent;
(unsigned long long)src, mc); if (b > DB_SEND_DATA_MAX)
uint32_t sent = 0; b = DB_SEND_DATA_MAX;
while (sent < mc) {
uint32_t b = mc - sent; uint8_t sdbuf[8192];
if (b > DB_SEND_DATA_MAX) uint32_t off = 0;
b = DB_SEND_DATA_MAX; sdbuf[off++] = DB_MSG_SEND_DATA;
memcpy(sdbuf + off, &sent, 4); off += 4;
uint8_t sdbuf[8192]; uint16_t rc = 0;
uint32_t off = 0; uint16_t* rcp = (uint16_t*)(sdbuf + off); off += 2;
sdbuf[off++] = DB_MSG_SEND_DATA;
memcpy(sdbuf + off, &sent, 4); off += 4; sqlite3_stmt* stmt;
uint16_t rc = 0; si_prep(si, &stmt,
uint16_t* rcp = (uint16_t*)(sdbuf + off); off += 2; "SELECT id,timestamp,datahash,data,author_signature"
" FROM \"%s\" ORDER BY timestamp,datahash"
sqlite3_stmt* stmt; " LIMIT ? OFFSET ?");
si_prep(si, &stmt, sqlite3_bind_int64(stmt, 1, (sqlite3_int64)b);
"SELECT id,timestamp,datahash,data,author_signature" sqlite3_bind_int64(stmt, 2, (sqlite3_int64)sent);
" FROM \"%s\" ORDER BY timestamp,datahash" while (sqlite3_step(stmt) == SQLITE_ROW) {
" LIMIT ? OFFSET ?"); uint64_t rid = (uint64_t)sqlite3_column_int64(stmt, 0);
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)b); uint64_t rts = (uint64_t)sqlite3_column_int64(stmt, 1);
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)sent); uint64_t rdh = (uint64_t)sqlite3_column_int64(stmt, 2);
while (sqlite3_step(stmt) == SQLITE_ROW) { const uint8_t* rd =
uint64_t rid = (uint64_t)sqlite3_column_int64(stmt, 0); (const uint8_t*)sqlite3_column_blob(stmt, 3);
uint64_t rts = (uint64_t)sqlite3_column_int64(stmt, 1); uint32_t rdl = (uint32_t)sqlite3_column_bytes(stmt, 3);
uint64_t rdh = (uint64_t)sqlite3_column_int64(stmt, 2); if (!rd) rdl = 0;
const uint8_t* rd = const uint8_t* rsig =
(const uint8_t*)sqlite3_column_blob(stmt, 3); (const uint8_t*)sqlite3_column_blob(stmt, 4);
uint32_t rdl = (uint32_t)sqlite3_column_bytes(stmt, 3); int rsl = sqlite3_column_bytes(stmt, 4);
if (!rd) rdl = 0; if (!rsig) rsl = 0;
const uint8_t* rsig =
(const uint8_t*)sqlite3_column_blob(stmt, 4); int rec_sz = 28 + rdl + 1 + rsl;
int rsl = sqlite3_column_bytes(stmt, 4); if (off + rec_sz > (int)sizeof(sdbuf))
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)
break; 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 (my_dh == peer_dh) {
if (memcmp(my_ch, pc, 32) == 0) {
struct SI_PEER* sp = si_peer_find(si, src); struct SI_PEER* sp = si_peer_find(si, src);
if (sp) if (sp) {
sp->synced = 2; sp->synced_pos = tp;
sp->sync_state = 2;
}
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"chain_hash match with %016llx at %u", "dh match with %016llx at %u",
(unsigned long long)src, tp); (unsigned long long)src, tp);
return; return;
} }
// Binary search for divergence point
uint32_t ds = 0, de = tp; uint32_t ds = 0, de = tp;
const uint8_t* spr = p + 37; const uint8_t* spr = p + 13;
for (int i = 0; i < sc && spr + 36 <= p + len; i++) { for (int i = 0; i < sc && spr + 12 <= p + len; i++) {
uint32_t pos = *(uint32_t*)spr; uint32_t pos = *(uint32_t*)spr;
const uint8_t* pch = spr + 4; uint64_t pdh;
uint8_t mch[32]; memcpy(&pdh, spr + 4, 8);
if (db_chain_hash_at(si, pos, mch) == 0) { uint64_t mdh;
if (memcmp(mch, pch, 32) == 0) { if (db_datahash_at(si, pos, &mdh) == 0) {
if (mdh == pdh) {
if (pos + 1 > ds) if (pos + 1 > ds)
ds = pos + 1; ds = pos + 1;
} }
@ -762,14 +842,13 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si,
de = pos; de = pos;
} }
} }
spr += 36; spr += 12;
} }
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"divergence with %016llx range [%u,%u]", "divergence with %016llx range [%u,%u]",
(unsigned long long)src, ds, de); (unsigned long long)src, ds, de);
if (de - ds <= 1) { if (de - ds <= 1) {
// Small divergence — REFINE with 0 hashes
uint8_t ref[512]; uint8_t ref[512];
uint32_t roff = 0; uint32_t roff = 0;
ref[roff++] = DB_MSG_REFINE; 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); db_sync_send(si, src, ref, roff);
} }
else { else {
// Sparse REFINE
uint8_t ref[512]; uint8_t ref[512];
uint32_t roff = 0; uint32_t roff = 0;
ref[roff++] = DB_MSG_REFINE; 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); uint16_t count = *(uint16_t*)(p + 4);
const uint8_t* ptr = p + 6; 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++) { for (uint16_t i = 0; i < count; i++) {
uint64_t rid, rts, rdh; uint64_t rid, rts, rdh;
uint32_t rdlen; uint32_t rdlen;
@ -963,56 +1045,59 @@ static void db_handle_send_data(struct DB_SYNC_INSTANCE* si,
break; break;
int ret = db_record_insert(si, rid, rts, rdh, int ret = db_record_insert(si, rid, rts, rdh,
(const char*)rdata, rdlen, (const char*)rdata, rdlen,
rsig, rsiglen); rsig, rsiglen, 0);
if (ret >= 0)
received++;
if (ret == 0 && si->on_insert) if (ret == 0 && si->on_insert)
si->on_insert(si, (const char*)rdata, rdlen, si->on_insert(si, (const char*)rdata, rdlen,
src, si->on_insert_arg); 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); uint32_t mc = db_count(si);
struct SI_PEER* sp = si_peer_find(si, src); uint32_t pk = from + received;
uint32_t pk = from + count;
// Request next batch if more records expected if (mc > pk && sp && sp->sync_state == 1) {
if (mc > pk && sp && sp->synced == 1) {
uint32_t scnt = mc - pk; uint32_t scnt = mc - pk;
if (scnt > DB_SEND_DATA_MAX) if (scnt > DB_SEND_DATA_MAX)
scnt = DB_SEND_DATA_MAX; scnt = DB_SEND_DATA_MAX;
si_send_data_batch(si, src, pk, scnt, 1); si_send_data_batch(si, src, pk, scnt, 1);
} }
// Send SYNC_DONE when we've received all
uint32_t nc = db_count(si); uint32_t nc = db_count(si);
uint8_t lch[32]; uint64_t ldh = 0;
if (nc > 0 && db_chain_hash_at(si, nc - 1, lch) == 0) { if (nc > 0)
uint8_t sd[37]; db_datahash_at(si, nc - 1, &ldh);
sd[0] = DB_MSG_SYNC_DONE; uint8_t sd[13];
memcpy(sd + 1, &nc, 4); sd[0] = DB_MSG_SYNC_DONE;
memcpy(sd + 5, lch, 32); memcpy(sd + 1, &nc, 4);
db_sync_send(si, src, sd, 37); memcpy(sd + 5, &ldh, 8);
} db_sync_send(si, src, sd, 13);
if (sp) if (sp)
sp->synced = 2; sp->sync_state = 2;
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"received %u records from %016llx", "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, static void db_handle_sync_done(struct DB_SYNC_INSTANCE* si,
uint64_t src, const uint8_t* p, size_t len) uint64_t src, const uint8_t* p, size_t len)
{ {
if (len < 36) if (len < 12)
return; return;
uint32_t pc = *(uint32_t*)p; 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); uint32_t mc = db_count(si);
uint8_t mch[32]; uint64_t mdh = 0;
if (mc > 0) if (mc > 0)
db_chain_hash_at(si, mc - 1, mch); db_datahash_at(si, mc - 1, &mdh);
else
memset(mch, 0, 32);
if (mc != pc || memcmp(mch, pch, 32) != 0) { if (mc != pc || mdh != pdh) {
DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC,
"SYNC_DONE mismatch %016llx my=%u peer=%u", "SYNC_DONE mismatch %016llx my=%u peer=%u",
(unsigned long long)src, mc, pc); (unsigned long long)src, mc, pc);
@ -1020,8 +1105,10 @@ static void db_handle_sync_done(struct DB_SYNC_INSTANCE* si,
return; return;
} }
struct SI_PEER* sp = si_peer_find(si, src); struct SI_PEER* sp = si_peer_find(si, src);
if (sp) if (sp) {
sp->synced = 2; sp->synced_pos = (mc < pc ? mc : pc) > 0 ? (mc < pc ? mc : pc) - 1 : 0;
sp->sync_state = 2;
}
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"sync confirmed with %016llx count=%u", "sync confirmed with %016llx count=%u",
(unsigned long long)src, mc); (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, int ret = db_record_insert(si, rid, rts, rdh,
(const char*)rdata, rdlen, (const char*)rdata, rdlen,
rsig, rsiglen); rsig, rsiglen, 0);
if (ret == 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) if (si->on_insert)
si->on_insert(si, (const char*)rdata, rdlen, si->on_insert(si, (const char*)rdata, rdlen,
src, si->on_insert_arg); 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) if (!si->enabled)
continue; continue;
struct SI_PEER* p = si_peer_add(si, pid); struct SI_PEER* p = si_peer_add(si, pid);
if (p && p->synced == 0 && conn->initialized && conn->links_up) { if (p && p->sync_state == 0 && conn->initialized && conn->links_up) {
p->synced = 1; p->sync_state = 1;
db_sync_initiate_sync(si, pid); 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++) { for (int i = 0; i < db->instance_count; i++) {
struct SI_PEER* p = si_peer_find(&db->instances[i], pid); struct SI_PEER* p = si_peer_find(&db->instances[i], pid);
if (p) if (p)
p->synced = 0; p->sync_state = 0;
} }
int any = 0; int any = 0;
for (int i = 0; i < db->instance_count; i++) { for (int i = 0; i < db->instance_count; i++) {
for (int j = 0; j < db->instances[i].peer_count; j++) { 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; any = 1;
goto cd_done; goto cd_done;
} }
@ -1262,22 +1354,57 @@ static void db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si,
"initiate_sync to %016llx my=%u", "initiate_sync to %016llx my=%u",
(unsigned long long)pid, mc); (unsigned long long)pid, mc);
uint8_t msg[37]; uint8_t msg[5];
msg[0] = DB_MSG_INIT_SYNC; msg[0] = DB_MSG_INIT_SYNC;
memcpy(msg + 1, &mc, 4); memcpy(msg + 1, &mc, 4);
if (mc > 0) { db_sync_send(si, pid, msg, 5);
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);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC,
"INIT_SYNC -> %016llx my=%u", "INIT_SYNC -> %016llx my=%u",
(unsigned long long)pid, mc); (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 // Timers
// ============================================================ // ============================================================
@ -1312,16 +1439,24 @@ static void db_sync_peer_check_cb(void* arg)
&& item->conn->links_up && item->conn->initialized) && item->conn->links_up && item->conn->initialized)
{ {
uint64_t pid = item->conn->peer_node_id; uint64_t pid = item->conn->peer_node_id;
if (pid != db->inst->node_id) { if (pid != db->inst->node_id)
struct SI_PEER* p = si_peer_add(si, pid); si_peer_add(si, pid);
if (p && p->synced == 0) {
p->synced = 1;
db_sync_initiate_sync(si, pid);
}
}
} }
e = e->next; 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 = db->peer_check_timer =
@ -1581,6 +1716,8 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst,
sqlite3_finalize(stmt); sqlite3_finalize(stmt);
} }
db_verify_chain(si);
si->ttl_timer = si->ttl_timer =
uasync_set_timeout(inst->ua, uasync_set_timeout(inst->ua,
DB_SYNC_TTL_INTERVAL * 10000u, 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) && ce->conn->links_up > 0 && ce->conn->initialized)
{ {
struct SI_PEER* p = si_peer_add(si, pid); struct SI_PEER* p = si_peer_add(si, pid);
if (p && p->synced == 0) { if (p && p->sync_state == 0) {
p->synced = 1; p->sync_state = 1;
db_sync_initiate_sync(si, pid); 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; uint64_t id = si->next_id;
int ret = db_record_insert(si, id, ts, dh, 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) if (ret != 0)
return ret; return ret;
if (si->on_insert) if (si->on_insert)
@ -1707,7 +1844,7 @@ int db_sync_insert_signed(struct DB_SYNC_INSTANCE* si,
// Push to all synced peers // Push to all synced peers
for (int i = 0; i < si->peer_count; i++) { 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) && si->peers[i].node_id != si->db_sync->inst->node_id)
{ {
if (db_sync_send(si, si->peers[i].node_id, pbuf, off) >= 0) if (db_sync_send(si, si->peers[i].node_id, pbuf, off) >= 0)

6
src/utun_instance.c

@ -413,6 +413,9 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) {
instance->etcp_sockets = NULL; instance->etcp_sockets = NULL;
DEBUG_INFO(DEBUG_CATEGORY_MEMORY, "[INSTANCE_DESTROY] ETCP sockets cleanup complete"); 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) // Cleanup ETCP connections (phase 1 detach + deferred phase 2 via call_soon)
{ {
struct ll_entry* entry = instance->connections->head; struct ll_entry* entry = instance->connections->head;
@ -469,9 +472,6 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) {
nat_transport_destroy(instance); nat_transport_destroy(instance);
} }
// Cleanup db_sync (before etcp_router_destroy)
db_sync_destroy(instance);
// Cleanup etcp_router // Cleanup etcp_router
etcp_router_destroy(instance); etcp_router_destroy(instance);

5
tests/Makefile.am

@ -50,6 +50,7 @@ check_PROGRAMS = \
test_conn_mgr \ test_conn_mgr \
test_etcp_connect \ test_etcp_connect \
test_db_sync \ test_db_sync \
test_chat_sync_stress \
test_stcp_traffic \ test_stcp_traffic \
test_bbr_integration \ test_bbr_integration \
test_intensive_memory_pool \ 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_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib
test_db_sync_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) 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_SOURCES = test_stcp_traffic.c
test_stcp_traffic_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_stcp_traffic_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib
test_stcp_traffic_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_stcp_traffic_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS)

246
tests/test_chat_sync_stress.c

@ -0,0 +1,246 @@
// test_chat_sync_stress.c — 10 раундов: 100 батчей вставок + 1 синхронизация за раунд
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <stdarg.h>
#include <time.h>
#include "../lib/platform_compat.h"
#include "test_utils.h"
#ifdef _WIN32
#include <windows.h>
#include <direct.h>
#else
#include <unistd.h>
#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;
}
Loading…
Cancel
Save