Browse Source

db_sync: fix test sync protocol isolation (add auto_sync flag, destroy order, test refactor)

- utun_instance: fix destroy order (topo_groups before conn_mgr to avoid UAF in topo_recovery)
- topo_recovery: add NULL guard for instance/ua in try_next
- db_sync_instance_add: add auto_sync flag (1=send INIT_SYNC, 0=quiet)
  replaces separate _quiet variant, cleaner API
- db_sync_peer_set_state, db_sync_reinitiate: test-only peer control utilities
- test_db_sync: full rewrite — stable, isolated sync protocol testing
  all 3 branches (peer_empty, hash_MATCH, divergence) tested without PUSH contamination
  Phase 6 extended: triple merge + tail PUSH + mid PUSH cascade on star topology
  runs in 1.4s, consistently passes 10/10
- db_sync_doc: updated protocol architecture description
- chatgui: update db_sync_instance_add call with auto_sync=1
topo_upd
Evgeny 2 months ago
parent
commit
e071a9e2f8
  1. 41
      src/db_sync.c
  2. 8
      src/db_sync.h
  3. 321
      src/db_sync_doc.md
  4. 5
      src/routing_layer/topo_recovery.c
  5. 14
      src/utun_instance.c
  6. 4
      tests/test_chat_sync_stress.c
  7. 309
      tests/test_db_sync.c
  8. 2
      tools/chatgui/transport/chat_channel.c

41
src/db_sync.c

@ -659,7 +659,8 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, uint64_t src, const
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "from=%016llx tp=%u peer_ch8=%016llx sc=%d", (unsigned long long)src, tp, (unsigned long long)peer_ch8, sc);
uint64_t my_ch8 = 0;
db_chain_hash8_at(si, tp, &my_ch8);
uint32_t mc_init = db_count(si);
if (mc_init > 0) db_chain_hash8_at(si, tp, &my_ch8);
if (peer_ch8 == 0 && sc == 0) {
uint32_t mc = db_count(si);
@ -1009,19 +1010,21 @@ static void db_handle_send_data(struct DB_SYNC_INSTANCE* si, uint64_t src, const
si->on_insert(si, rts, (const char*)rdata, rdlen, src, si->on_insert_arg);
}
db_cascade_from(si, fix_from);
if (received > 0) {
db_cascade_from(si, fix_from);
for (int j = 0; j < si->peer_count; j++) {
if (si->peers[j].synced_pos >= fix_from && &si->peers[j] != sp) {
uint32_t old = si->peers[j].synced_pos;
si->peers[j].synced_pos = fix_from > 0 ? fix_from - 1 : 0;
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] cascade changed chain from pos=%u → peer=%016llx synced_pos %u→%u",
SI_SHRT(si), (unsigned long long)(src >> 16), fix_from,
(unsigned long long)si->peers[j].node_id, old, si->peers[j].synced_pos);
for (int j = 0; j < si->peer_count; j++) {
if (si->peers[j].sync_state != 1 && si->peers[j].synced_pos >= fix_from && &si->peers[j] != sp) {
uint32_t old = si->peers[j].synced_pos;
si->peers[j].synced_pos = fix_from > 0 ? fix_from - 1 : 0;
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "sync [%s:%04llX] cascade changed chain from pos=%u → peer=%016llx synced_pos %u→%u (state=%u)",
SI_SHRT(si), (unsigned long long)(src >> 16), fix_from,
(unsigned long long)si->peers[j].node_id, old, si->peers[j].synced_pos, si->peers[j].sync_state);
}
}
}
if (sp && received > 0) sp->synced_pos = from + received - 1;
if (sp && count > 0) sp->synced_pos = from + count - 1;
uint32_t mc = db_count(si);
uint32_t pk = from + received;
@ -1211,7 +1214,7 @@ static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry)
if (!ch_id || !ch_id[0]) continue;
uint64_t h; { uint8_t sh[32]; SC_SHA256_CTX ctx; sc_sha256_init(&ctx); sc_sha256_update(&ctx, (const uint8_t*)ch_id, strlen(ch_id)); sc_sha256_final(&ctx, sh); memcpy(&h, sh, 8); }
if (h == hash) {
si = db_sync_instance_add(db->inst, tbl, hash);
si = db_sync_instance_add(db->inst, tbl, hash, 1);
if (si) { si_peer_add(si, src); }
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "lazy-register: tbl=%s ch=%s hash=%016llx", tbl, ch_id, (unsigned long long)hash);
break;
@ -1414,6 +1417,18 @@ 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);
if (p) { p->sync_state = 0; db_sync_initiate_sync(si, node_id); }
}
// ============================================================
// Timers
// ============================================================
@ -1608,7 +1623,7 @@ void db_sync_destroy(struct UTUN_INSTANCE* inst)
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: destroyed");
}
struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const char* table_name, uint64_t hash)
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; }
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "table=%s hash=%016llx", table_name, (unsigned long long)hash);
@ -1675,7 +1690,7 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const
if (pid != 0 && pid != inst->node_id && ce->conn->links_up > 0 && ce->conn->initialized) {
peers_found++;
struct SI_PEER* p = si_peer_add(si, pid);
if (p && p->sync_state == 0) { p->sync_state = 1; db_sync_initiate_sync(si, pid); peers_synced++; }
if (p && p->sync_state == 0 && auto_sync) { p->sync_state = 1; db_sync_initiate_sync(si, pid); peers_synced++; }
}
entry = entry->next;
}

8
src/db_sync.h

@ -90,7 +90,7 @@ int db_sync_init(struct UTUN_INSTANCE* inst);
void db_sync_destroy(struct UTUN_INSTANCE* inst);
// Instance management
struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const char* table_name, uint64_t hash);
struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const char* table_name, uint64_t hash, int auto_sync);
void db_sync_instance_remove(struct DB_SYNC_INSTANCE* si);
// Data operations (per-instance)
@ -123,6 +123,12 @@ void db_sync_set_insert_cb(struct DB_SYNC_INSTANCE* si, db_sync_insert_cb cb
// Verify chain integrity: returns 0 if all chain_hashes are correct, 1 if any mismatch found
int db_sync_chain_verify(struct DB_SYNC_INSTANCE* si);
// Peer state control (for testing — disable PUSH to isolate sync protocol)
void db_sync_peer_set_state(struct DB_SYNC_INSTANCE* si, uint64_t node_id, uint8_t state);
// Force re-initiate sync to a specific peer (for testing)
void db_sync_reinitiate(struct DB_SYNC_INSTANCE* si, uint64_t node_id);
#ifdef __cplusplus
}
#endif

321
src/db_sync_doc.md

@ -1,95 +1,131 @@
# db_sync — Distributed content-addressed table with SQLite + P2P sync
# db_sync — Distributed append-only table with SQLite + P2P sync
## 1. Назначение
Децентрализованная реплицируемая таблица JSON-записей, синхронизируемая между всеми узлами сети через ETCP. Каждый узел хранит полную копию данных каждого инстанса. Модуль поддерживает несколько независимых инстансов (таблиц), каждый идентифицируется хешем `SHA256(name || id_be)[0:8]`.
Реплицировать append-only таблицу JSON-записей между всеми пирами P2P-сети. Каждый пир в итоге должен иметь идентичный набор записей. Модуль поддерживает несколько независимых инстансов (таблиц), каждый идентифицируется хешем `SHA256(name || id_be)[0:8]`.
**Ключевые свойства:**
- **Content-addressed**: целостность цепочки гарантируется `chain_hash` — каскадным SHA256, где `chain_hash[n] = SHA256(chain_hash[n-1] || id || timestamp || author || author_signature)`.
- **Ed25519-подписи**: каждая запись обязательно подписана автором (`Ed25519(timestamp || json_data)`). Записи без подписи или с неверной подписью отвергаются.
- **Multi-instance**: несколько независимых таблиц внутри одного процесса (напр. `("chats", 1)` и `("chats", 2)`).
- **Ordered**: записи упорядочены по `(timestamp, author_signature)`. Первичный ключ — та же пара, определяющая уникальность (дубликаты по тому же автору в ту же миллисекунду невозможны).
- **Append-only**: записи не редактируются и не удаляются явно, только TTL-очистка собственных неотправленных записей.
- **Sync protocol**: 8 типов сообщений (INIT_SYNC → INIT_RESP → REFINE → SEND_DATA → SYNC_DONE, плюс PUSH/ACK_PUSH/ERROR). Используется бинарный поиск расхождений по chain_hash.
- **Push**: новые записи немедленно рассылаются (PUSH) всем синхронизированным пирам.
## 2. Ключевые свойства
## 2. Как пользоваться
- **Криптографическая цепь.** `chain_hash[N] = SHA256(chain_hash[N-1] || id || timestamp || author || author_signature)`. Записи упорядочены `ORDER BY timestamp, author_signature`. Первые 8 байт — `chain_hash8` — используется для быстрого сравнения в протоколе.
- **Ed25519-подписи.** Каждая запись подписана автором. Записи без подписи или с неверной подписью отвергаются.
- **Append-only.** Записи не редактируются и не удаляются явно, только TTL-очистка собственных неотправленных записей.
- **PUSH — мгновенная доставка.** При локальном `insert_signed` запись немедленно шлётся всем пирам с `sync_state >= 1` (steady state).
- **Sync — сравнение цепей.** При старте/реконнекте стороны обмениваются хешами цепей для поиска расхождений.
- **Multi-instance.** Несколько независимых таблиц внутри одного процесса (разные hash).
### Типовой сценарий
## 3. Архитектура протокола
```
1. В конфиге: db_sync_enabled = 1, db_sync_ttl = 86400
2. db_sync_init(inst) — вызывается автоматически при старте utun_instance
3. struct DB_SYNC_INSTANCE* si = db_sync_instance_add(inst, "chats", 1);
→ Создаёт таблицу SQLite, верифицирует цепочку, запускает TTL-таймер.
→ Если есть активные ETCP-соединения — автоматически запускает синхронизацию.
4. Вставка записи:
uint64_t ts = db_sync_next_timestamp(si); // монотонно возрастающий timestamp (ms)
uint8_t sig[64];
Ed25519_sign(my_privkey, ts || my_node_id || json, sig);
db_sync_insert_signed(si, json_str, json_len, sig, 64, ts);
→ Подпись ОБЯЗАТЕЛЬНА (sig=NULL → ошибка). Запись автоматически push-ится всем synced пирам.
5. Чтение:
db_sync_select(si, offset, limit, my_callback, ctx);
6. db_sync_instance_remove(si) — деактивировать инстанс (таблица БД не удаляется).
7. db_sync_destroy(inst) — вызывается автоматически при завершении.
```
### 3.1. На что оптимизирован
### Ключевые концепции
99% времени новые записи просто дописываются в конец. Самый частый сценарий — пир A добавил сообщение, пир B получил его и вставил в конец своей цепи. Протокол должен:
- Доставлять новые записи немедленно (PUSH)
- При реконнекте быстро понять "у нас всё совпадает до позиции N, добрось хвост" (hash_MATCH)
- При реальном расхождении бинарным поиском найти точку и слить (divergence + REFINE)
#### Multi-instance
Каждый экземпляр `DB_SYNC_INSTANCE` идентифицируется 64-битным хешем `hash = SHA256(name || htobe64(id))[0:8]`. Этот хеш используется как routing key в синхронизационных сообщениях. На приёмной стороне по хешу находится нужный инстанс.
### 3.2. Общие идеи реализации
#### Chain hash
Каскадный хеш цепочки записей:
**Криптографическая цепь** — гарантия целостности. Если `chain_hash8(N)` совпадает на двух пирах, цепь идентична до позиции N (вероятность коллизии 2^-64). При вставке не в конец — каскадный пересчёт `chain_hash` всех последующих записей через `db_cascade_from`.
```
chain_hash[0] = SHA256(0x00...00 || id_0 || ts_0 || author_0 || sig_0)
chain_hash[n] = SHA256(chain_hash[n-1] || id_n || ts_n || author_n || sig_n)
```
**PUSH** — рабочий механизм доставки в steady state. Запись, вставленная локально, немедленно уходит всем синхронизированным пирам. ACK_PUSH подтверждает получение и обновляет delivery_chain.
**Sync** — полное сравнение цепей при старте/реконнекте. Три ветки: peer_empty (пир пуст — отдать всё), hash_MATCH (цепи совпали до позиции N — отдать хвост), divergence (цепи разошлись — найти точку расхождения).
**Sparse checkpoints** — бинарный поиск расхождения. INIT_RESP возвращает хеши в степенях двойки от tp (tp-1, tp-2, tp-4, ..., до 16 шт). Это позволяет за O(log N) сравнений сузить диапазон, не передавая хеш каждой записи.
**REFINE** — финальное сужение. Если sparse-хешей INIT_RESP недостаточно (диапазон >1), REFINE запрашивает до 16 дополнительных хешей. Когда диапазон ≤1 — сразу SEND_DATA.
**Каскадное уведомление.** При вставке записи не в конец `synced_pos` всех остальных пиров сбрасывается до позиции вставки — им потребуется пересинхронизация.
### 3.3. Фазы протокола
**Фаза А — Инициализация.** `db_sync_instance_add` читает таблицу из SQLite, проверяет целостность цепи (`db_verify_chain` — автофикс при расхождении), и немедленно шлёт `INIT_SYNC(my_count)` всем подключённым пирам с `sync_state == 0`. Если связь появилась позже — `conn_up` делает то же самое.
**Фаза Б — Сравнение цепей.** Получатель INIT_SYNC вычисляет `tp = min(my_count, peer_count)`, tp-- если >0 (последняя гарантированно общая позиция), и возвращает INIT_RESP: `peer_ch8` на позиции tp, `sc` (количество sparse-хешей), sparse-хеши на позициях tp-2^k.
**Фаза В — Три ветки:**
| Ветка | Условие | Сценарий | Действие |
|-------|---------|----------|----------|
| peer_empty | `peer_ch8==0 && sc==0` | Пир пуст (0 записей) | Отправить все свои записи через SEND_DATA |
| hash_MATCH | `my_ch8 == peer_ch8` | Цепи идентичны до tp | Отправить хвост [tp+1..mc) через SEND_DATA |
| divergence | `my_ch8 != peer_ch8` | Разные истории | Анализ sparse-хешей → REFINE → SEND_DATA |
**Фаза Г — REFINE.** Инициатор анализирует sparse-хеши: `ds` — последняя совпавшая позиция, `de` — первая разошедшаяся. Если `de-ds ≤ 1` → сразу SEND_DATA (hc=0). Иначе → REFINE с до 16 своих хешей, равномерно распределённых в [ds..de]. Получатель сравнивает со своей цепью, находит точку совпадения, шлёт SEND_DATA от этой точки.
**Фаза Д — SEND_DATA.** Получатель вставляет записи с проверкой Ed25519-подписи, делает `db_cascade_from(fix_from)`, уведомляет остальных пиров о сдвиге цепи (сброс их synced_pos). Если у получателя после вставки записей больше чем у отправителя — proactive push-back (шлёт свой хвост). Когда все записи получены → SYNC_DONE с итоговым count и chain_hash8.
**Фаза Е — SYNC_DONE.** Сравнение итогового count и chain_hash8. Не совпало — ретрай всей процедуры с начала (до 3 раз, потом give up с partial sync). Совпало — `sync_state=2`, синхронизация завершена.
### 3.4. PUSH — отдельный от sync механизм
PUSH матчится по `hash` — если у пира нет si с таким же hash, PUSH не доставляется. Это позволяет изолировать тестирование sync-протокола от PUSH: вставлять данные через si с уникальным hash (PUSH не уходит — нет получателя), затем удалять tmp si (данные в SQLite сохраняются), создавать si с общим hash — instance_add запускает чистый sync.
## 4. Peer management
Для протокола синхронизации используются первые 8 байт (`chain_hash8`). Сравнивая эти хеши на разных позициях, узлы находят точку расхождения.
#### Sync protocol (8 message types)
| Message | Direction | Описание |
|---------|-----------|----------|
| `DB_MSG_INIT_SYNC (0x01)` | A→B | Инициатор шлёт своё количество записей (`my_count`) |
| `DB_MSG_INIT_RESP (0x02)` | B→A | Truncation point `tp = min(counts)-1`, `chain_hash8[tp]`, + до 16 sparse-хешей на позициях `tp-2^k` |
| `DB_MSG_REFINE (0x03)` | A→B, B→A | Бинарный поиск: до 16 равномерно распределённых хешей в диапазоне расхождения |
| `DB_MSG_SEND_DATA (0x04)` | A→B, B→A | Пакетная передача записей (до 32 за раз). Wire: `[id:8][ts:8][author:8][dlen:4][data][sig_len:1=64][sig:64]` |
| `DB_MSG_PUSH (0x05)` | A→B | Рассылка одной новой записи всем synced-пирам |
| `DB_MSG_ACK_PUSH (0x06)` | B→A | Подтверждение получения PUSH; обновляет delivery_chain и флаг WAS_SENT |
| `DB_MSG_SYNC_DONE (0x07)` | B→A | Финальное подтверждение: peer_count + last_chain_hash8 |
| `DB_MSG_ERROR (0x08)` | B→A | Ошибка. Коды: `0x01` — instance not found, `0x02` — instance disabled |
**Алгоритм синхронизации:**
1. `INIT_SYNC`: A → B: `my_count_a`; B вычисляет `tp = min(count_a, count_b) - 1`.
2. `INIT_RESP`: B → A: `tp`, `chain_hash8_b[tp]`, + sparse-хеши на позициях `tp-1, tp-2, tp-4, tp-8, ..., tp-32768`.
3. Если `hash8_b[tp] == hash8_a[tp]` — цепочки совпадают до `tp`. Более длинная сторона шлёт «хвост» (записи после `tp`).
4. Если sparse-хеши показывают расхождение: A определяет диапазон `[ds, de]` где хеши не совпадают. Если `de-ds ≤ 1` — пустой `REFINE` (запрос данных). Иначе — `REFINE` с до 16 хешами, равномерно распределёнными в диапазоне.
5. `REFINE`: B ищет первую позицию, где хеши разошлись, шлёт `SEND_DATA` начиная с этой позиции.
6. `SEND_DATA`: A вставляет полученные записи (с верификацией Ed25519 подписи), пересчитывает каскадный chain_hash, шлёт `SYNC_DONE`.
7. `SYNC_DONE`: если counts/hashes не совпали — реинициируется sync.
#### Защита от подделок
- Ed25519-подпись автора проверяется для КАЖДОЙ вставляемой записи (и локальной, и от пиров).
- Публичный ключ автора ищется: (1) в своих ключах (если self), (2) в `topo_node_sqlite`, (3) в `peer_ed25519_pubkey` активного ETCP-соединения.
- Если подпись невалидна — запись отвергается с логом "discarding as forgery".
#### Peer management
- При поднятии ETCP-соединения для каждого инстанса добавляется `SI_PEER` и запускается синхронизация.
- При разрыве соединения `sync_state` пира сбрасывается в 0.
- `peer_check` таймер (каждые 5с) перебирает `topo_group->senders_list` и запускает синхронизацию для пиров с `sync_state == 0`.
- `PUSH` рассылается только пирам в состоянии `sync_state >= 1`.
- `peer_check` таймер (каждые 5с) перебирает активных пиров и запускает синхронизацию для тех, у кого `sync_state == 0`. Выбирается пир с минимальным `synced_pos` — двигаемся от самого старого несинхронизированного участка.
- PUSH рассылается только пирам в состоянии `sync_state >= 1`.
## 5. Сообщения протокола
| Сообщение | Wire-формат | Описание |
|-----------|-------------|----------|
| `INIT_SYNC (0x01)` | `[type:1][count:4]` | Инициатор шлёт количество записей |
| `INIT_RESP (0x02)` | `[type:1][tp:4][ch8:8][sc:1][(pos:4,ch8:8)*sc]` | tp, chain_hash8 на tp, sc sparse-хешей |
| `REFINE (0x03)` | `[type:1][from:4][to:4][hc:1][(pos:4,ch8:8)*hc]` | hc=0 → сразу SEND_DATA; hc>0 → до 16 хешей в диапазоне |
| `SEND_DATA (0x04)` | `[type:1][from:4][count:2][(id:8,ts:8,author:8,dlen:4,data,sig_len:1,sig)*count]` | Пакет до 32 записей |
| `PUSH (0x05)` | `[type:1][id:8][ts:8][author:8][dlen:4][data][sig_len:1=64][sig:64]` | Рассылка одной записи synced-пирам |
| `ACK_PUSH (0x06)` | `[type:1][ts:8][author:8]` | Подтверждение PUSH, обновление delivery_chain |
| `SYNC_DONE (0x07)` | `[type:1][count:4][ch8:8]` | Финальный count + chain_hash8 |
| `ERROR (0x08)` | `[type:1][code:1]` | Коды: 0x01=NOT_FOUND, 0x02=DISABLED |
Все сообщения маршрутизируются через ETCP service `0x20` с префиксом `[svc:1][hash_be:8][payload]`.
## 6. Wire-формат записи
`[id:8][ts:8][author_node_id:8][dlen:4][data:json][sig_len:1=64][sig:ed25519:64]`
Дубликаты определяются по `(timestamp, author_signature)`.
## 7. Структуры данных
```c
struct DB_SYNC {
struct UTUN_INSTANCE* inst;
sqlite3* db;
uint8_t shared_db;
uint64_t last_connected_tb;
struct DB_SYNC_INSTANCE* instances;
int instance_count, instance_capacity;
void* peer_check_timer;
uint8_t enabled;
};
#### TTL cleanup
- Периодический таймер (каждые 3600с).
- Удаляет записи authored by self, где НЕ установлен флаг `DB_REC_FLAG_WAS_SENT` (0x01) и `timestamp < now - db_sync_ttl`.
- Чужие записи и свои отправленные не удаляются.
struct DB_SYNC_INSTANCE {
struct DB_SYNC* db_sync;
uint64_t hash;
char table_name[64];
uint64_t next_id;
uint64_t last_timestamp_ms;
uint8_t enabled;
void* ttl_timer;
struct SI_PEER* peers;
int peer_count, peer_capacity;
db_sync_insert_cb on_insert;
void* on_insert_arg;
};
## 3. API
struct SI_PEER {
uint64_t node_id;
uint32_t synced_pos; // последняя общая позиция (0-based)
uint8_t sync_state; // 0=idle, 1=syncing, 2=synced
uint8_t sync_retry_count; // счётчик retry SYNC_DONE mismatch
uint64_t sync_start_tb; // время начала sync (для timeout)
};
```
## 8. API
### Глобальный жизненный цикл
@ -97,17 +133,21 @@ chain_hash[n] = SHA256(chain_hash[n-1] || id_n || ts_n || author_n || sig_n)
int db_sync_init(struct UTUN_INSTANCE* inst);
void db_sync_destroy(struct UTUN_INSTANCE* inst);
```
- `db_sync_init` — инициализирует DB_SYNC, открывает SQLite по пути `<db_path>/sync`, биндит ETCP service `0x20`, вешает коллбэки на существующие и новые соединения, стартует `peer_check` таймер. Если `db_sync_enabled = 0` — создаёт структуру в disabled-режиме.
- `db_sync_destroy` — отменяет все таймеры, анбиндит ETCP service, снимает коллбэки со всех соединений, закрывает SQLite, освобождает память.
`db_sync_init` — открывает SQLite по пути `<db_path>/chats.db`, биндит ETCP service `0x20`, вешает коллбэки соединений, стартует `peer_check` таймер (5с). Если `db_sync_enabled = 0` — disabled-режим.
`db_sync_destroy` — отменяет таймеры, анбиндит сервис, снимает коллбэки, закрывает SQLite.
### Управление инстансами
```c
struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const char* name, uint64_t id);
struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const char* table_name, uint64_t hash);
void db_sync_instance_remove(struct DB_SYNC_INSTANCE* si);
```
- `db_sync_instance_add` — создаёт/регистрирует инстанс. Проверяет валидность имени (`[a-zA-Z0-9_]`, макс 48 символов). Вычисляет hash, создаёт SQLite-таблицу `db_sync_<name>_<id>`, верифицирует цепочку хешей, запускает TTL-таймер. Если уже есть активные соединения — автоматически инициирует sync.
- `db_sync_instance_remove` — деактивирует инстанс: `enabled=0`, отменяет TTL-таймер, освобождает peers, удаляет из массива. Таблица БД **не удаляется**.
`db_sync_instance_add` — создаёт/регистрирует инстанс. Создаёт SQLite-таблицу (если нет), читает цепь, `db_verify_chain` (автофикс), стартует TTL-таймер. Если есть подключённые пиры с `sync_state == 0` — немедленно шлёт `INIT_SYNC`.
`db_sync_instance_remove` — деактивирует: `enabled=0`, отменяет TTL-таймер, освобождает peers, удаляет из массива. Таблица БД не удаляется.
### Операции с данными
@ -118,12 +158,10 @@ uint32_t db_sync_count(struct DB_SYNC_INSTANCE* si);
uint64_t db_sync_get_last_timestamp(struct DB_SYNC_INSTANCE* si);
uint64_t db_sync_next_timestamp(struct DB_SYNC_INSTANCE* si);
```
- `db_sync_insert_signed` — вставляет подписанную запись. `sig` ОБЯЗАТЕЛЬНО 64 байта (Ed25519). Проверяет подпись, проверяет дубликат по `(timestamp, author_signature)`, вычисляет chain_hash, вставляет, каскадно пересчитывает хеши последующих записей, рассылает PUSH всем synced-пирам. Возвращает 0 (успех), 1 (уже существует), -1 (ошибка), -2 (неверная подпись).
- `db_sync_count` — количество записей в локальной БД.
- `db_sync_get_last_timestamp` — последний выданный `db_sync_next_timestamp`.
- `db_sync_next_timestamp` — возвращает монотонно возрастающий timestamp в миллисекундах. Гарантирует `last_timestamp_ms < returned`.
### Чтение (select)
`db_sync_insert_signed` — вставляет подписанную запись. `sig` — 64 байта Ed25519. Проверяет подпись, проверяет дубликат, вычисляет chain_hash, вставляет с каскадным пересчётом, рассылает PUSH synced-пирам. Возвращает: 0=успех, 1=дубликат, -1=ошибка, -2=неверная подпись.
### Чтение
```c
typedef void (*db_sync_select_cb)(void* arg, uint64_t id, uint64_t timestamp,
@ -133,7 +171,6 @@ typedef void (*db_sync_select_cb)(void* arg, uint64_t id, uint64_t timestamp,
int db_sync_select(struct DB_SYNC_INSTANCE* si, uint32_t offset, uint32_t limit,
db_sync_select_cb cb, void* arg);
```
- `db_sync_select` — итератор по записям, упорядоченным `ORDER BY timestamp, author_signature`. `limit=0` — без ограничения. Для каждой записи вызывает `cb` с полями: id, timestamp, data (JSON), author (node_id), author_sig (64 байта), delivered_peers (счётчик доставок), delivery_chain (hex-идентификаторы пиров через запятую). Возвращает количество переданных в callback записей.
### Callback на вставку
@ -144,96 +181,44 @@ typedef void (*db_sync_insert_cb)(struct DB_SYNC_INSTANCE* si,
uint64_t author_node_id, void* arg);
void db_sync_set_insert_cb(struct DB_SYNC_INSTANCE* si, db_sync_insert_cb cb, void* arg);
```
- `db_sync_set_insert_cb` — устанавливает callback, вызываемый после успешной вставки записи (локальной или от пира). `author_node_id = self` для локальных вставок, `= peer_node_id` для записей от пиров.
### Внутренние структуры
### Верификация цепи
```c
struct DB_SYNC {
struct UTUN_INSTANCE* inst; // обратная ссылка на инстанс
sqlite3* db; // SQLite handle
uint64_t last_connected_tb; // время последнего подключения (timebase)
struct DB_SYNC_INSTANCE* instances; // динамический массив инстансов
int instance_count, instance_capacity;
void* peer_check_timer; // handle uasync-таймера
uint8_t enabled; // глобальный флаг из конфига
};
struct DB_SYNC_INSTANCE {
struct DB_SYNC* db_sync; // обратная ссылка
uint64_t hash; // routing key = SHA256(name||id_be)[0:8]
char table_name[64]; // "db_sync_<name>_<id>"
uint64_t next_id; // монотонно возрастающий id записей
uint64_t last_timestamp_ms; // последний выданный timestamp (ms)
uint8_t enabled; // per-instance enable/disable
void* ttl_timer; // handle TTL-таймера
struct SI_PEER* peers; // динамический массив пиров
int peer_count, peer_capacity;
db_sync_insert_cb on_insert; // callback на вставку
void* on_insert_arg;
};
struct SI_PEER {
uint64_t node_id; // идентификатор пира
uint32_t synced_pos; // позиция, до которой синхронизированы
uint8_t sync_state; // 0=idle, 1=syncing, 2=synced
};
int db_sync_chain_verify(struct DB_SYNC_INSTANCE* si);
```
### Конфигурация
Возвращает 0 если все chain_hash корректны, 1 если расхождение.
| Параметр | Файл конфига | По умолчанию | Описание |
|----------|-------------|--------------|----------|
| `db_sync_enabled` | `global` | `0` | Включить модуль синхронизации |
| `db_sync_ttl` | `global` | `86400` | TTL неотправленных собственных записей (секунды) |
| `db_path` | `global` | — | Путь к БД; файл создаётся как `<db_path>/sync` |
## 9. Конфигурация
### Константы синхронизации
| Параметр | По умолчанию | Описание |
|----------|-------------|----------|
| `db_sync_enabled` | `0` | Включить модуль синхронизации |
| `db_sync_ttl` | `86400` | TTL неотправленных собственных записей (сек) |
| `db_path` | — | Путь к БД; файл `<db_path>/chats.db` |
## 10. Константы
| Константа | Значение | Описание |
|-----------|---------|----------|
| `ETCP_RT_ID_DB_SYNC` | `0x20` | ETCP service ID |
| `DB_REFINE_HASHES` | `16` | Макс. количество хешей в REFINE |
| `DB_SEND_DATA_MAX` | `32` | Макс. записей в одном SEND_DATA |
| `DB_REFINE_HASHES` | `16` | Макс. хешей в REFINE |
| `DB_SEND_DATA_MAX` | `32` | Макс. записей в SEND_DATA |
| `DB_SIG_SIZE` | `64` | Размер Ed25519 подписи |
| `DB_SYNC_PEER_CHECK_INTERVAL` | `5` | Интервал проверки пиров (секунды) |
| `DB_SYNC_TTL_INTERVAL` | `3600` | Интервал TTL-очистки (секунды) |
### Зависимости
- **etcp_api.h / etcp.h** — P2P-обмен сообщениями, коллбэки соединений
- **topo_group.h / topo_node_sqlite.h** — топология, discovery пиров, Ed25519 pubkeys
- **secure_channel.h** — SHA256 (`sc_sha256_*`), Ed25519 verify (`sc_ed25519_verify`)
- **lib/u_async.h** — асинхронный event loop, таймеры (`uasync_set_timeout`, `uasync_cancel_timeout`)
- **lib/sqlite3.h** — SQLite3 (WAL, synchronous=NORMAL, auto-checkpoint 10000)
- **lib/mem.h** — `u_malloc`, `u_calloc`, `u_realloc`, `u_free`
- **lib/debug_config.h** — `DEBUG_INFO`, `DEBUG_WARN`, `DEBUG_ERROR` (категория `DEBUG_CATEGORY_DEBUG`)
- **lib/sha256.h** — `SC_SHA256_CTX` (если не USE_OPENSSL)
- **lib/platform_compat.h** — `htobe64`, `get_time_tb`, `utun_gettimeofday`
- **openssl/evp.h** — низкоуровневый Ed25519 в `sc_ed25519_verify`
### Внутренние функции
| Функция | Назначение |
|---------|-----------|
| `db_sqlite_open/db_sqlite_close` | Открытие/закрытие SQLite с WAL-режимом |
| `db_sha256/db_hash64/db_chain_hash_compute` | SHA256 и вычисление chain_hash |
| `db_instance_find/db_instance_alloc` | Поиск/выделение инстанса по hash |
| `si_peer_find/si_peer_add` | Поиск/добавление пира в инстансе |
| `db_count` | `SELECT COUNT(*)` из таблицы инстанса |
| `db_chain_hash_at/db_chain_hash8_at` | Чтение chain_hash/chain_hash8 на позиции |
| `db_prev_chain_hash` | Поиск chain_hash записи, предшествующей данной |
| `db_cascade_from` | Пересчёт chain_hash начиная с позиции `from_pos` |
| `db_get_ed25519_pubkey` | Получение Ed25519 pubkey пира |
| `db_verify_author_sig` | Проверка Ed25519 подписи автора записи |
| `db_record_insert` | Полный цикл вставки: проверка подписи, проверка дубликата, вставка, каскад chain_hash |
| `db_sync_send/db_sync_send_hash` | Отправка сообщения пиру через ETCP |
| `si_delivery_update` | Обновление delivery_chain при доставке |
| `si_parse_record` | Разбор одной записи из wire-формата |
| `si_send_data_batch` | Формирование и отправка SEND_DATA пакета |
| `db_sync_recv_cb` | Центральный callback приёма ETCP-сообщений, маршрутизация по instance_hash |
| `db_handle_init_sync/_resp/_refine/_send_data/_push/_ack_push/_sync_done/_error` | Обработчики каждого типа сообщения |
| `db_sync_initiate_sync` | Отправка INIT_SYNC пиру |
| `db_verify_chain` | Полная верификация цепочки chain_hash (при старте инстанса) |
| `db_sync_peer_check_cb` | Таймер: периодический поиск активных пиров и запуск sync |
| `db_sync_instance_ttl_cb` | Таймер: TTL-очистка неотправленных собственных записей |
| `DB_SYNC_PEER_CHECK_INTERVAL` | `5` | Интервал проверки пиров (сек) |
| `DB_SYNC_SYNC_TIMEOUT` | `15` | Таймаут ожидания ответа при sync (сек) |
| `DB_SYNC_TTL_INTERVAL` | `3600` | Интервал TTL-очистки (сек) |
## 11. Тестирование sync-протокола
PUSH доставляет записи немедленно и независимо от sync. Чтобы тестировать чистый sync-протокол, нужно чтобы PUSH не вмешивался:
1. Вставить данные через si с **уникальным hash** — PUSH уходит, но не доставляется (нет получателя с таким hash)
2. `remove_si()` — данные сохраняются в SQLite, si деактивирован
3. Создать si с **общим hash** на обеих сторонах — `instance_add` запускает чистый sync-протокол
Этот паттерн используется для тестирования всех трёх веток INIT_RESP:
- **peer_empty**: одна сторона с данными, другая пустая
- **hash_MATCH**: обе имеют общий префикс, но у одной больше записей
- **divergence**: независимые вставки на обеих сторонах до синхронизации

5
src/routing_layer/topo_recovery.c

@ -163,6 +163,11 @@ static void topo_recovery_try_next(struct TOPO_RECOVERY_CTX* ctx) {
ctx->current_node_id = 0;
continue;
}
if (!ctx->instance || !ctx->instance->ua) {
DEBUG_WARN(DEBUG_CATEGORY_BGP, "recovery: instance/ua gone, aborting try_next for %016llx",
(unsigned long long)ctx->next_hop_id);
return;
}
ctx->connect_timer = uasync_set_timeout(ctx->instance->ua, TOPO_RECOVERY_CONNECT_TIMEOUT_MS * 10,
ctx, topo_recovery_timeout_cb, "recovery_timeout");
return;

14
src/utun_instance.c

@ -482,18 +482,18 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) {
nat_transport_destroy(instance);
}
// Cleanup conn_mgr (must be before etcp_router_destroy — unbinds from router)
if (instance->conn_mgr) { conn_mgr_destroy(instance->conn_mgr); instance->conn_mgr = NULL; }
// Cleanup etcp_router
etcp_router_destroy(instance);
// Cleanup BGP module
// Cleanup BGP module (uses conn_mgr → must be before conn_mgr_destroy)
if (instance->topo_groups) {
DEBUG_INFO(DEBUG_CATEGORY_BGP, "Destroying BGP module");
topo_groups_destroy(instance);
}
// Cleanup conn_mgr (unbinds from etcp_router → must be before etcp_router_destroy)
if (instance->conn_mgr) { conn_mgr_destroy(instance->conn_mgr); instance->conn_mgr = NULL; }
// Cleanup etcp_router
etcp_router_destroy(instance);
// Cleanup NAT detection
if (instance->nat_det) {
nat_detection_destroy(instance->nat_det);

4
tests/test_chat_sync_stress.c

@ -151,8 +151,8 @@ int main(void) {
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);
si_a = db_sync_instance_add(inst_a, INSTANCE_NAME, INSTANCE_ID, 1);
si_b = db_sync_instance_add(inst_b, INSTANCE_NAME, INSTANCE_ID, 1);
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");

309
tests/test_db_sync.c

@ -1,4 +1,4 @@
// test_db_sync.c — тестирование db_sync: dh_match tail-send и peer empty
// test_db_sync.c — sync protocol test: peer_empty, hash_MATCH, divergence merge
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
@ -28,21 +28,23 @@
#include "../lib/debug_config.h"
#include "../lib/mem.h"
#define TEST_TIMEOUT_TB 600000 // 60s
#define PHASE_TIMEOUT_TB 250000 // 25s
#define POLL_INTERVAL_MS 5
#define TEST_TIMEOUT_TB 30000 // 3s
#define PHASE_TIMEOUT_TB 25000 // 2.5s per phase
#define POLL_INTERVAL_MS 1
static struct UTUN_INSTANCE* inst_a = NULL;
static struct UTUN_INSTANCE* inst_b = NULL;
static struct UTUN_INSTANCE* inst_c = NULL;
static struct DB_SYNC_INSTANCE* si_a = NULL;
static struct DB_SYNC_INSTANCE* si_b = NULL;
static struct DB_SYNC_INSTANCE* si_c = NULL;
static struct UASYNC* ua = NULL;
static int test_phase = 0;
static void* timeout_id = NULL;
static char temp_dir[] = "/tmp/utun_dbsync_XXXXXX";
static char config_a[256], config_b[256];
static int port_a_srv, port_b_srv;
static char config_a[256], config_b[256], config_c[256];
static int port_a_srv, port_b_srv, port_c_srv;
static int write_file(const char* path, const char* fmt, ...) {
va_list ap; FILE* f = fopen(path, "w");
@ -60,44 +62,57 @@ static char* get_pubkey(const char* path) {
static int create_temp_configs(void) {
if (test_mkdtemp(temp_dir) != 0) { fprintf(stderr, "mkdtemp failed\n"); return -1; }
int base = 42000 + (getpid() % 15000);
port_a_srv = base; port_b_srv = base + 1;
port_a_srv = base; port_b_srv = base + 1; port_c_srv = base + 2;
snprintf(config_a, sizeof(config_a), "%s/a.conf", temp_dir);
snprintf(config_b, sizeof(config_b), "%s/b.conf", temp_dir);
snprintf(config_c, sizeof(config_c), "%s/c.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(db_path, sizeof(db_path), "%s/db_c", temp_dir); utun_mkdir(db_path, 0755);
if (write_file(config_a,
write_file(config_a,
"[global]\nmy_node_id=0xAAAAAAAAAAAAAAAA\ntun_ip=10.200.0.1/24\n"
"tun_ifname=tun200\nkeepalive_adaptive=0\ndb_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;
"[allowed_keys]\nallow_all=1\n", temp_dir, port_a_srv);
config_ensure_keys_and_node_id(config_a);
char* pub_a = get_pubkey(config_a); if (!pub_a) { fprintf(stderr, "pub_a fail\n"); return -1; }
if (write_file(config_b,
write_file(config_b,
"[global]\nmy_node_id=0xBBBBBBBBBBBBBBBB\ntun_ip=10.200.0.2/24\n"
"tun_ifname=tun201\nkeepalive_adaptive=0\ndb_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; }
"[allowed_keys]\nallow_all=1\n", temp_dir, port_b_srv, pub_a, port_a_srv);
write_file(config_c,
"[global]\nmy_node_id=0xCCCCCCCCCCCCCCCC\ntun_ip=10.200.0.3/24\n"
"tun_ifname=tun202\nkeepalive_adaptive=0\ndb_path=%s/db_c\ndb_sync_enabled=1\n\n"
"[server: srv_c]\naddr=127.0.0.1:%d\ntype=public\n\n"
"[client: to_a]\nkeepalive=1\npeer_public_key=%s\nlink=srv_c:127.0.0.1:%d\n\n"
"[allowed_keys]\nallow_all=1\n", temp_dir, port_c_srv, pub_a, port_a_srv);
free(pub_a);
if (config_ensure_keys_and_node_id(config_b) != 0) return -1;
config_ensure_keys_and_node_id(config_b);
config_ensure_keys_and_node_id(config_c);
return 0;
}
static void cleanup_temp_configs(void) {
unlink(config_a); unlink(config_b);
unlink(config_a); unlink(config_b); unlink(config_c);
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/chats.db", temp_dir); unlink(pa);
snprintf(pa, sizeof(pa), "%s/db_a/chats.db-wal", temp_dir); unlink(pa);
snprintf(pa, sizeof(pa), "%s/db_a/chats.db-shm", temp_dir); unlink(pa);
snprintf(pa, sizeof(pa), "%s/db_b/chats.db", temp_dir); unlink(pa);
snprintf(pa, sizeof(pa), "%s/db_b/chats.db-wal", temp_dir); unlink(pa);
snprintf(pa, sizeof(pa), "%s/db_b/chats.db-shm", temp_dir); unlink(pa);
snprintf(pa, sizeof(pa), "%s/db_c/chats.db", temp_dir); unlink(pa);
snprintf(pa, sizeof(pa), "%s/db_c/chats.db-wal", temp_dir); unlink(pa);
snprintf(pa, sizeof(pa), "%s/db_c/chats.db-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);
snprintf(pa, sizeof(pa), "%s/db_c", temp_dir); test_rmdir(pa);
test_rmdir(temp_dir);
}
@ -115,29 +130,30 @@ static int wait_for(const char* desc, int (*cond)(void), int timeout_tb) {
static int cond_links_init(void) {
if (!inst_a || !inst_b) return 0;
struct ll_entry* e = inst_a->connections->head;
int links = 0;
while (e) { struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data;
struct ETCP_LINK* l = ce->conn->links; while (l) { if (l->initialized) return 1; l = l->next; }
struct ETCP_LINK* l = ce->conn->links; while (l) { if (l->initialized) links++; l = l->next; }
e = e->next; }
return 0;
return links >= (inst_c ? 2 : 1);
}
static uint32_t ca_target, cb_target;
static uint32_t ca_target, cb_target, cc_target;
static int _cond_ca(void) { return si_a && db_sync_count(si_a) == ca_target; }
static int _cond_cb(void) { return si_b && db_sync_count(si_b) == cb_target; }
static int _cond_cc(void) { return si_c && db_sync_count(si_c) == cc_target; }
static int _cond_both(void) { return _cond_ca() && _cond_cb(); }
static int _cond_chain_ok(void) { return db_sync_chain_verify(si_a) == 0 && db_sync_chain_verify(si_b) == 0; }
static int _cond_all(void) { return _cond_ca() && _cond_cb() && _cond_cc(); }
static int insert_many(struct DB_SYNC_INSTANCE* si, struct UTUN_INSTANCE* inst, int start, int count) {
char buf[128];
for (int i = start; i < start + count && test_phase == 0; i++) {
snprintf(buf, sizeof(buf), "{\"idx\":%d,\"val\":\"data_%d\",\"pad\":\"%.50s\"}", i, i,
"xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx");
snprintf(buf, sizeof(buf), "{\"idx\":%d,\"val\":\"data_%d\"}", i, i);
uint64_t ts = db_sync_next_timestamp(si);
uint8_t sig_msg[256]; size_t off = 0;
memcpy(sig_msg + off, &ts, 8); off += 8;
size_t jl = strlen(buf); memcpy(sig_msg + off, buf, jl); off += jl;
uint8_t sig[64];
if (sc_ed25519_sign(inst->my_ed25519_privkey, sig_msg, off, sig) != SC_OK) { fprintf(stderr,"sign fail %d\n", i); return -1; }
if (db_sync_insert_signed(si, buf, jl, sig, 64, ts) < 0) { fprintf(stderr,"insert_many fail %d\n", i); return -1; }
if (sc_ed25519_sign(inst->my_ed25519_privkey, sig_msg, off, sig) != SC_OK) { fprintf(stderr,"sign fail\n"); return -1; }
if (db_sync_insert_signed(si, buf, jl, sig, 64, ts) < 0) { fprintf(stderr,"insert fail\n"); return -1; }
}
return 0;
}
@ -146,150 +162,171 @@ static void remove_si(struct DB_SYNC_INSTANCE** psi) {
if (*psi) { db_sync_instance_remove(*psi); *psi = NULL; }
}
// ---- Main test ----
int main(void) {
printf("=== test_db_sync ===\n");
debug_config_init(); debug_set_level(DEBUG_LEVEL_ERROR);
if (create_temp_configs() != 0) { fprintf(stderr, "config creation failed\n"); return 1; }
if (create_temp_configs() != 0) { cleanup_temp_configs(); return 1; }
utun_instance_set_tun_init_enabled(0);
ua = uasync_create();
if (!ua) { cleanup_temp_configs(); return 1; }
ua = uasync_create(); if (!ua) { 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 || utun_instance_init(inst_a) != 0 || utun_instance_init(inst_b) != 0) {
fprintf(stderr, "instance create/init failed\n"); cleanup_temp_configs(); return 1;
}
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; }
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; }
// ===================================================================
// Phase 1: late B — A=50, B=0. B создан ПОСЛЕ conn_up (A не инициирует).
// Только B инициирует: divergence → REFINE → SEND_DATA(1) → SYNC_DONE
// mismatch → A ретраит → dh_match(tp=0,sc=0) → должен отправить хвост [1..49].
// Phase 1: peer_empty. add_quiet (no auto-sync) → insert A=5,
// reinitiate → INIT_SYNC(5) → B INIT_RESP(ch8=0,sc=0) → A sends all 5.
// ===================================================================
printf("Phase 1: late B — dh_match at tp=0 (sc=0) + tail-send 49 records\n");
if (!wait_for("links up", cond_links_init, PHASE_TIMEOUT_TB)) { test_phase=2; goto done; }
si_a = db_sync_instance_add(inst_a, "test", 1);
if (!si_a || insert_many(si_a, inst_a, 0, 50) != 0) { test_phase=2; goto done; }
if (db_sync_count(si_a) != 50) { fprintf(stderr, "FAIL: A!=50\n"); test_phase=2; goto done; }
printf(" A has 50 records, creating B now (B empty, A conn_up already fired)\n");
si_b = db_sync_instance_add(inst_b, "test", 1);
if (!si_b || db_sync_count(si_b) != 0) { test_phase=2; goto done; }
cb_target = 50;
if (!wait_for("B count=50 (dh_match tp=0 tail)", _cond_cb, PHASE_TIMEOUT_TB)) { test_phase=2; goto done; }
printf(" Phase 1 PASS: B=%u A=%u\n", db_sync_count(si_b), db_sync_count(si_a));
printf("Phase 1: peer_empty\n");
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; }
if (insert_many(si_a, inst_a, 0, 5) != 0) { test_phase = 2; goto done; }
db_sync_reinitiate(si_a, inst_b->node_id);
cb_target = 5;
if (!wait_for("B=5", _cond_cb, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; }
if (db_sync_count(si_b) != 5 || db_sync_chain_verify(si_a) || db_sync_chain_verify(si_b))
{ test_phase = 2; goto done; }
printf(" PASS\n");
// ===================================================================
// Phase 2: recreate A after adding 30 more — A=80, B=50.
// A инициирует: dh_match(tp=49, sc>0, sparse hashes) — должен отправить хвост [50..79].
// Phase 2: hash_MATCH tail-send. Insert 3 more on A (PUSH disabled —
// si_a from phase 1 has sync_state=2? No, reinitiate set it to 1, then
// SYNC_DONE set it to 2. Disable PUSH, insert, reinitiate.
// ===================================================================
printf("Phase 2: recreate A — dh_match at tp=49 (sc>0 sparse) + tail-send 30 records\n");
if (insert_many(si_a, inst_a, 50, 30) != 0) { test_phase=2; goto done; }
if (db_sync_count(si_a) != 80) { fprintf(stderr, "FAIL: A!=80\n"); test_phase=2; goto done; }
remove_si(&si_a);
si_a = db_sync_instance_add(inst_a, "test", 1);
if (!si_a || db_sync_count(si_a) != 80) { test_phase=2; goto done; }
printf(" A recreated (80 records), B has 50 — A will initiate with tail\n");
cb_target = 80;
if (!wait_for("B count=80 (dh_match tp=49 sparse)", _cond_cb, PHASE_TIMEOUT_TB)) { test_phase=2; goto done; }
printf(" Phase 2 PASS: B=%u A=%u\n", db_sync_count(si_b), db_sync_count(si_a));
printf("Phase 2: hash_MATCH tail-send\n");
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_a, inst_b->node_id);
cb_target = 8;
if (!wait_for("B=8", _cond_cb, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; }
if (db_sync_count(si_b) != 8 || db_sync_chain_verify(si_a) || db_sync_chain_verify(si_b))
{ test_phase = 2; goto done; }
printf(" PASS\n");
// ===================================================================
// Phase 3: fresh instances, A=30, B=0, both created simultaneously.
// A инициирует: peer empty (peer_dh=0, sc=0) — отправляет все 30.
// B тоже инициирует но расходится через divergence+reinit+dh_match(tp=0) —
// но A уже всё отправил, так что B и так получит через A's sync.
// Phase 3: divergence merge. remove si, add_quiet → insert A=3 B=2 →
// reinitiate both → divergence → REFINE merge to 5.
// ===================================================================
printf("Phase 3: fresh instances A=30 B=0, both initiate — peer empty\n");
remove_si(&si_a);
remove_si(&si_b);
si_a = db_sync_instance_add(inst_a, "test2", 2);
si_b = db_sync_instance_add(inst_b, "test2", 2);
if (!si_a || !si_b) { test_phase=2; goto done; }
if (insert_many(si_a, inst_a, 0, 30) != 0 || db_sync_count(si_a) != 30) { test_phase=2; goto done; }
cb_target = 30;
if (!wait_for("B count=30 (peer empty)", _cond_cb, PHASE_TIMEOUT_TB)) { test_phase=2; goto done; }
printf(" Phase 3 PASS: B=%u A=%u\n", db_sync_count(si_b), db_sync_count(si_a));
printf("Phase 3: divergence merge\n");
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);
if (!si_a || !si_b) { test_phase = 2; goto done; }
if (insert_many(si_a, inst_a, 0, 3) != 0 || insert_many(si_b, inst_b, 10, 2) != 0)
{ test_phase = 2; goto done; }
db_sync_reinitiate(si_a, inst_b->node_id);
db_sync_reinitiate(si_b, inst_a->node_id);
ca_target = 5; cb_target = 5;
if (!wait_for("both=5", _cond_both, 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");
// ===================================================================
// Phase 4: Divergence — A=10 B=5 DIFFERENT records in same table.
// Insert via different hash si (no PUSH between them), then re-add
// with same hash → both have existing records → chains differ →
// REFINE → si_send_data_batch → converge to 15.
// Phase 4: divergence merge. remove si, add_quiet → insert A=2 B=2 →
// reinitiate both → merge to 4.
// ===================================================================
printf("Phase 4: divergence — A=10 B=5 independent records, REFINE/SEND_DATA merge\n");
remove_si(&si_a);
remove_si(&si_b);
{ // Insert with different hashes → no PUSH between A and B
struct DB_SYNC_INSTANCE* tmp_a = db_sync_instance_add(inst_a, "test_div", 10);
struct DB_SYNC_INSTANCE* tmp_b = db_sync_instance_add(inst_b, "test_div", 20);
if (!tmp_a || !tmp_b) { test_phase = 2; goto done; }
if (insert_many(tmp_a, inst_a, 0, 10) != 0 || db_sync_count(tmp_a) != 10) { test_phase = 2; goto done; }
if (insert_many(tmp_b, inst_b, 100, 5) != 0 || db_sync_count(tmp_b) != 5) { test_phase = 2; goto done; }
printf(" A has 10 records (indices 0-9), B has 5 records (indices 100-104)\n");
remove_si(&tmp_a);
remove_si(&tmp_b);
}
si_a = db_sync_instance_add(inst_a, "test_div", 30);
si_b = db_sync_instance_add(inst_b, "test_div", 30);
printf("Phase 4: divergence merge 2\n");
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);
if (!si_a || !si_b) { test_phase = 2; goto done; }
if (insert_many(si_a, inst_a, 0, 2) != 0 || insert_many(si_b, inst_b, 10, 2) != 0)
{ test_phase = 2; goto done; }
db_sync_reinitiate(si_a, inst_b->node_id);
db_sync_reinitiate(si_b, inst_a->node_id);
ca_target = 4; cb_target = 4;
if (!wait_for("both=4", _cond_both, 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");
ca_target = 15; cb_target = 15; // 10 + 5 = 15 after merge
if (!wait_for("both count=15 (divergence resolved)", _cond_both, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; }
if (db_sync_chain_verify(si_a) != 0 || db_sync_chain_verify(si_b) != 0) {
fprintf(stderr, "FAIL: chain_hash mismatch after divergence merge\n"); test_phase = 2; goto done;
// ===================================================================
// Phase 5: PUSH not-at-tail. Use db_sync_instance_add (normal, with auto-sync)
// so PUSH fires. Insert early timestamp record → PUSH → cascade_from.
// ===================================================================
printf("Phase 5: PUSH not-at-tail\n");
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; }
{ 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;
size_t jl = strlen(ebuf); memcpy(msg + moff, ebuf, jl); moff += jl;
uint8_t sig[64];
if (sc_ed25519_sign(inst_a->my_ed25519_privkey, msg, moff, sig) != SC_OK
|| db_sync_insert_signed(si_a, ebuf, jl, sig, 64, early_ts) != 0)
{ test_phase = 2; goto done; }
}
printf(" Phase 4 PASS: B=%u A=%u chains OK\n", db_sync_count(si_b), db_sync_count(si_a));
ca_target = 3; cb_target = 3;
if (!wait_for("both=3", _cond_both, 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");
// ===================================================================
// Phase 5: PUSH not at tail — insert record with early timestamp
// on A → PUSH to B → B must cascade_from to fix chain_hashes
// Phase 6: triple star — cascade notification.
// ===================================================================
printf("Phase 5: PUSH not at tail — early timestamp, cascade verification\n");
printf("Phase 6: triple star\n");
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; }
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);
si_b = db_sync_instance_add(inst_b, "triple", 70, 0);
si_c = db_sync_instance_add(inst_c, "triple", 70, 0);
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_b, inst_c->node_id);
db_sync_reinitiate(si_c, inst_a->node_id); db_sync_reinitiate(si_c, inst_b->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 (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");
// 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 (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");
// 6c: PUSH not-at-tail on A → delivered to B and C → cascade_from, synced_pos reset.
{
char ebuf[128];
snprintf(ebuf, sizeof(ebuf), "{\"idx\":\"early\",\"val\":\"not_at_tail\"}");
uint64_t early_ts = 500;
uint8_t msg[256]; size_t moff = 0;
char ebuf[64]; snprintf(ebuf, sizeof(ebuf), "{\"idx\":\"mid\",\"val\":\"triple_insert\"}");
uint64_t early_ts = 500; uint8_t msg[256]; size_t moff = 0;
memcpy(msg + moff, &early_ts, 8); moff += 8;
size_t jl = strlen(ebuf); memcpy(msg + moff, ebuf, jl); moff += jl;
uint8_t sig[64];
if (sc_ed25519_sign(inst_a->my_ed25519_privkey, msg, moff, sig) != SC_OK)
{ fprintf(stderr, "sign fail\n"); test_phase = 2; goto done; }
if (db_sync_insert_signed(si_a, ebuf, jl, sig, 64, early_ts) != 0)
{ fprintf(stderr, "insert_signed fail\n"); test_phase = 2; goto done; }
}
ca_target = 16; cb_target = 16;
if (!wait_for("both count=16 after PUSH", _cond_both, PHASE_TIMEOUT_TB)) { test_phase = 2; goto done; }
if (db_sync_chain_verify(si_a) != 0 || db_sync_chain_verify(si_b) != 0) {
fprintf(stderr, "FAIL: chain_hash mismatch after PUSH cascade\n"); test_phase = 2; goto done;
if (sc_ed25519_sign(inst_a->my_ed25519_privkey, msg, moff, sig) != SC_OK
|| db_sync_insert_signed(si_a, ebuf, jl, sig, 64, early_ts) != 0)
{ test_phase = 2; goto done; }
}
printf(" Phase 5 PASS: B=%u A=%u chains OK\n", db_sync_count(si_b), db_sync_count(si_a));
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 (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");
done:
remove_si(&si_a);
remove_si(&si_b);
remove_si(&si_a); remove_si(&si_b); remove_si(&si_c);
if (timeout_id && ua) { uasync_cancel_timeout(ua, timeout_id); timeout_id = NULL; }
if (inst_a) { inst_a->running = 0; utun_instance_destroy(inst_a); inst_a = NULL; }
if (inst_c) { inst_c->running = 0; utun_instance_destroy(inst_c); inst_c = NULL; }
if (inst_b) { inst_b->running = 0; utun_instance_destroy(inst_b); inst_b = NULL; }
if (inst_a) { inst_a->running = 0; utun_instance_destroy(inst_a); inst_a = NULL; }
if (ua) { uasync_destroy(ua, 0); ua = NULL; }
cleanup_temp_configs();
if (test_phase == 0) test_phase = 1;
printf("=== %s ===\n", test_phase == 1 ? "PASS" : "FAIL");
return test_phase == 1 ? 0 : 1;

2
tools/chatgui/transport/chat_channel.c

@ -83,7 +83,7 @@ void chat_core_ensure_channel_ready(const char* ch_id) {
uint64_t ch_hash = 0;
{ const uint8_t* chd = (const uint8_t*)ch_id; size_t chl = strlen(ch_id); uint8_t sh[32]; SHA256(chd, chl, sh); memcpy(&ch_hash, sh, 8); }
struct DB_SYNC_INSTANCE* si = db_sync_instance_add(g_cc.inst, tbl_msg, ch_hash);
struct DB_SYNC_INSTANCE* si = db_sync_instance_add(g_cc.inst, tbl_msg, ch_hash, 1);
uint64_t gid = strtoull(ch_id, NULL, 10);
if (gid != 0 && g_cc.inst->topo_groups && !topo_groups_find(g_cc.inst->topo_groups, gid))

Loading…
Cancel
Save