21 changed files with 1562 additions and 1871 deletions
@ -1,149 +1,60 @@
|
||||
# member_sync — синхронизация участников канала |
||||
# member_sync — адаптер участников каналов |
||||
|
||||
## 1. Назначение |
||||
Поверх [merkle_sync](merkle_sync_doc.md) синхронизирует source=0 записи таблиц peers_<channel_id>. |
||||
Транспортный движок не занимается подписями и версиями участников. |
||||
|
||||
Тонкая прослойка над `merkle_sync`, адаптированная под участников (мемберов) чат-каналов. |
||||
Хранит мемберов в SQLite (`peers_<ch_id>`, `nodes`, `node_addresses`), пересчитывает Merkle-дерево |
||||
при каждом изменении и через merkle_sync рассылает дельту другим узлам. |
||||
## Использование |
||||
|
||||
Данные мембера защищены двумя независимо подписанными блоками: |
||||
- **Блок A (мембер)** — идентичность и профиль: `node_id`, `x25519`, `ed25519`, `userinfo` (имя), |
||||
подпись `update_sig` **ключом самого мембера** (`ed25519`); `ver = update_ts` внутри эпохи `join_ts`. |
||||
- **Блок B (владелец)** — атрибуты, управляемые владельцем канала: `adm_tags` (JSON с `ver`, |
||||
`storage`), подпись `adm_tags_sig` **канальным ключом** (`ed25519` канала); `ver = adm_tags.ver`. |
||||
- member_sync_init вызывается при инициализации chat_sync и перестраивает производный индекс. |
||||
- member_sync_start подключает канал к пиру и запрашивает подтверждение совпадения данных. |
||||
- member_sync_put изменяет локальную запись, атомарно обновляет индекс и после commit уведомляет движок. |
||||
- Прямые broadcast/send_to не нужны: локальная версия возвращается второму узлу при обратном проходе. |
||||
- cancel/destroy снимают ожидания; отдельного публичного API рассылки нет. |
||||
|
||||
Оба блока обновляются независимо — только если их `ver` увеличился. Устаревший блок помечается |
||||
stale и автору отправляется свежая запись (send-back). |
||||
## Записи и применение |
||||
|
||||
## 2. Как пользоваться |
||||
У записи два независимо версионируемых блока: участник (join_ts, update_ts) и владелец (adm_tags.ver). |
||||
Удаление используется для повреждённых записей. Конфликт одинаковых версий с разным содержимым |
||||
не разрешается порядком доставки: отсутствие прогресса завершает Merkle-раунд ошибкой конфликта. |
||||
|
||||
```c |
||||
// ── разово (из chat_sync_init) ── |
||||
member_sync_init(inst); // merkle_sync_init(0x31) + verify/purge + rebuild + подписка |
||||
Входящая страница сначала целиком проверяется на корректность раскладки, длин и строк. |
||||
Затем записи применяются под savepoint member_page. Ошибка откатывает всю страницу вместе с индексом. |
||||
Уведомления GUI/топологии и changed выполняются после commit. |
||||
Каждый member_sync_apply_record также имеет savepoint, поэтому ошибка пересчёта не оставляет запись без индекса. |
||||
При внешней транзакции вызывающая сторона отвечает за changed после commit. |
||||
|
||||
// ── подписка на BGP-события канала (BGP → member_sync → node_props_changed) ── |
||||
member_sync_subscribe_group(inst, group); // при создании канала |
||||
Проверяются привязка node_id к X25519, доверенный Ed25519, подписи блока участника и владельца. |
||||
Для дерева приглашений родитель может прийти позже ребёнка; finish повторно проверяет зависимости |
||||
перед подтверждением корня. Это не отдельная карантинная таблица: промежуточные записи существуют в БД |
||||
до finish, поэтому MT_OK нельзя заменять фактом их присутствия. |
||||
|
||||
// ── запуск синхронизации (из cs_on_conn_up) ── |
||||
member_sync_start(inst, peer, ch_id, on_done, ch); |
||||
## Канонический хеш |
||||
|
||||
// ── коллбэк завершения ── |
||||
static void on_done(uint64_t peer, const char* ch_id, int result, void* arg) { |
||||
if (result == MT_OK) ch->synced = CS_SYNC_DONE; |
||||
} |
||||
SHA256 по полям: |
||||
node_id:8, x25519:32, ed25519:32, join_sig:64, join_ts:8, |
||||
update_sig:64, update_ts:8, userinfo_len:1, userinfo, |
||||
adm_tags_len:1, adm_tags, adm_tags_sig:64, signed_by:8, signature:64. |
||||
|
||||
// ── добавление/обновление мембера (локальное, без верификации подписей) ── |
||||
member_sync_put(inst, ch_id, node_id, x25519, ed25519, |
||||
join_sig, join_ts, update_sig, update_ts, userinfo, |
||||
adm_tags, adm_tags_sig, storage); |
||||
Целые big-endian; отсутствующие подписи представлены нулями. |
||||
Обе строки ограничены 255 байтами. Хеш включает userinfo и дерево приглашений. |
||||
Формат сообщений, подписываемых моделью, определяется build_join_msg/build_update_msg/build_sign_msg отдельно. |
||||
|
||||
// ── онлайн-статус: только локальная запись, возвращает 1 если изменилось ── |
||||
if (member_sync_set_online(inst, node_id, 1)) |
||||
merkle_sync_push_update(inst, ch_id, node_id, 0x01, (uint8_t[]){1}, 1); // рассылает caller |
||||
## Страница на проводе |
||||
|
||||
// ── отмена (коллбэк не вызовется) ── |
||||
member_sync_cancel(inst, peer, ch_id); |
||||
Начинается с count:2 big-endian. Далее count записей: |
||||
node_id:8, x25519:32, ed25519:32, flags:1, |
||||
при HAS_JOIN — join_sig:64 и join_ts:8, |
||||
update_sig:64, update_ts:8, userinfo_len:1, userinfo, |
||||
adm_tags_len:1, adm_tags, adm_tags_sig:64, signed_by:8, signature:64. |
||||
|
||||
// ── количество мемберов ── |
||||
int n = member_sync_count(inst, ch_id); |
||||
Целые big-endian. Встроенные NUL и хвостовые байты запрещены. |
||||
get_page перечисляет только один лист, в возрастающем порядке node_id, строго после курсора. |
||||
SQLite использует signed INTEGER; внутри одного листа знак ключей одинаков, поэтому SQL-порядок совпадает с unsigned. |
||||
Следующая запись, не помещающаяся в буфер, остаётся для следующей страницы. |
||||
Ошибки SQL и повреждённые размеры BLOB возвращаются вызывающему коду. |
||||
|
||||
// ── разово (из chat_sync_destroy) ── |
||||
member_sync_destroy(inst); |
||||
``` |
||||
## Диагностика |
||||
|
||||
## 3. Модель данных |
||||
|
||||
- **Таблицы**: `peers_<ch_id>` (per-channel мемберы), `nodes`/`node_addresses` (глобальная |
||||
идентичность и адреса, обновляются отдельно через `topo_node_sqlite_*`). |
||||
- **Хеш мембера** (участвует в Merkle-дереве): |
||||
`SHA256(node_id || x25519 || ed25519 || join_sig || join_ts || update_sig || update_ts || adm_tags_len || adm_tags || adm_tags_sig)`. |
||||
Поля `userinfo` (имя) и `source` в хеш **не** входят. |
||||
- **Колонка `source`**: `source != 0` — плейсхолдер, созданный topo_group (не верифицирован); |
||||
такие записи **исключаются** из хеша и сериализации. `source == 0` — обычный верифицированный |
||||
меркл-рекорд. |
||||
|
||||
### Wire-формат мембера (get_items / apply_items / _serialize_member) |
||||
|
||||
``` |
||||
[count:2][node_id:8][x25519:32][ed25519:32][flags:1] |
||||
(flags&HAS_JOIN → [join_sig:64][join_ts:8]) |
||||
[update_sig:64][update_ts:8] |
||||
[userinfo_len:1][userinfo] |
||||
[adm_tags_len:1][adm_tags][adm_tags_sig:64] |
||||
``` |
||||
|
||||
`flags = PEERS_FLAG_HAS_JOIN` — `join_sig`/`join_ts` присутствуют (join-запись). |
||||
|
||||
## 4. apply_record — по-блочное сравнение и применение |
||||
|
||||
`member_sync_apply_record(inst, ch_id, from_peer, &rec)` — единая точка записи мембера |
||||
(и локальная `put`, и приёмная). Возвращает битовую маску `MS_APPLY_CHANGED | MS_APPLY_STALE`. |
||||
|
||||
1. **Идентичность** (только для приёмных записей): `node_id == sc_derive_node_id_from_pubkey(x25519)` |
||||
и `ed25519` совпадает с доверенным из `nodes`. Несовпадение → forgery, блок A отклоняется. |
||||
2. **Блок A** — проверка `update_sig` (или `join_sig` для join-only). Локальная запись |
||||
(`from_peer == inst->node_id`) доверенная, без верификации. Сравнение `ver`: |
||||
- нет локальной записи → changed; |
||||
- `join_ts > local_join_ts` → changed (**re-join / смена ключа**); |
||||
- `join_ts < local_join_ts` → stale (старая эпоха); |
||||
- иначе по `update_ts`. |
||||
3. **Блок B** — проверка `adm_tags_sig` канальным ключом, `adm_storage = (storage=="yes")`. |
||||
Сравнение `adm_tags.ver` с локальным. |
||||
4. **Запись**: changed-блоки пишутся (`member_block_put` / `member_owner_put`), узел обновляется |
||||
через `node_update_verified`, пересчитывается классификация `nodeinfo_updated` (из |
||||
`node_addresses`, полученных по BGP), затем `merkle_sync_recompute_path`. |
||||
5. **Реальная сверка**: если `recompute_path` вернул 0 (данные идентичны) — `changed` сбрасывается |
||||
в 0 (no-op), чтобы не рассылать. |
||||
|
||||
На приёме (`_member_apply_items`) каждый `MS_APPLY_CHANGED` генерирует `node_props_changed` и |
||||
включается в relay (`merkle_sync_broadcast`); `MS_APPLY_STALE` вызывает `member_sync_send_to` |
||||
(отправить нашу свежую запись автору). |
||||
|
||||
## 5. Подписки и коллбэки |
||||
|
||||
- **`member_sync_subscribe_group(inst, group)`** — подписывает на BGP node-события chat-группы |
||||
(`topo_group_add_node_cbk`). При событии BGP (`_ms_on_bgp_node`) рассылается `node_props_changed` |
||||
с текущим `adm_tags`. REMOVE **не удаляет** мембера. |
||||
- **`member_sync_add_props_cbk/remove_props_cbk`** — многоподписочный `node_props_changed_fn` |
||||
(изменились `adm_tags` любого узла, включая себя). |
||||
- **`member_sync_set_node_updated_cb`** — одиночный `member_sync_node_updated_fn` (online-статус |
||||
узла изменился по инициативе удалённого пира, вызывается из uasync-потока). |
||||
|
||||
## 6. Верификация и удаление |
||||
|
||||
Единственный механизм удаления мембера — битая запись: |
||||
- `member_sync_verify_local_record` — проверяет привязку `node_id = derive(x25519)`, подписи |
||||
`update_sig`/`join_sig`/`adm_tags_sig`. Отсутствие подписи = «не верифицировано» → **не удаляем**. |
||||
- `member_sync_verify_and_purge` / `_all` — прогоняют записи канала (всех `peers_%`), битые удаляют |
||||
через `topo_node_sqlite_member_del` + `recompute_path`. |
||||
- Вызываются при `member_sync_init` и при старте синхронизации канала. Отдельного |
||||
`member_sync_del` нет — conn_down ≠ удаление мембера. |
||||
|
||||
## 7. API |
||||
|
||||
| Функция | Назначение | |
||||
|---|---| |
||||
| `member_sync_init(inst)` | `merkle_sync_init(0x31, &g_member_ops, inst)` + verify/purge + rebuild всех деревьев + подписка на существующие chat-группы | |
||||
| `member_sync_destroy(inst)` | Отписка от групп, `merkle_sync_destroy` | |
||||
| `member_sync_subscribe_group(inst, group)` | Подписка на BGP-события CHAT-группы | |
||||
| `member_sync_start(inst, peer, ch_id, done_cb, arg)` | `verify_and_purge` + `merkle_sync_start` | |
||||
| `member_sync_cancel(inst, peer, ch_id)` | `merkle_sync_cancel` | |
||||
| `member_sync_put(...)` | Локальная запись мембера (trusted) → `apply_record` | |
||||
| `member_sync_apply_record(inst, ch_id, from_peer, &rec)` | Приём/применение рекорда (по-блочно) | |
||||
| `member_sync_send_to(inst, ch_id, member_id, target_peer)` | Send-back нашей свежей записи конкретному пиру | |
||||
| `member_sync_broadcast_one(inst, ch_id, member_id)` | Сериализовать мембера и разослать всем (relay после локального изменения) | |
||||
| `member_sync_set_online(inst, node_id, online)` | Записать `nodes.online`. Возвращает 1 (изменилось) / 0 | |
||||
| `member_sync_count(inst, ch_id)` | `COUNT(*)` из `peers_<ch_id>` | |
||||
| `member_sync_build_join_msg(...)` / `build_update_msg(...)` | Канонические сообщения для подписи/верификации join/update (112 байт и 64+8+userinfo) | |
||||
| `member_sync_get_hash(inst, ch_id, level, prefix64)` | Хеш бакета (для тестов) | |
||||
| `member_sync_add/remove_props_cbk`, `set_node_updated_cb` | Подписки на изменения | |
||||
|
||||
## 8. Нюансы |
||||
|
||||
- **Онлайн-статус не участвует в Merkle-дереве**: не хешируется и не версионируется. `set_online` |
||||
только пишет в БД и возвращает 1/0; рассылку `MSG_ITEM_UPDATE` выполняет **caller** (chat_sync). |
||||
- **`member_sync_put` = локальный `apply_record`** с `from_peer = inst->node_id` — без верификации |
||||
подписей (свои данные доверенные). |
||||
- **`source`-плейсхолдеры** topo_group пропускаются (`src != 0`), в дереве не участвуют. |
||||
- **Re-join / смена ключа** распознаётся по `join_ts` (перекрывает сравнение `update_ts`). |
||||
- Категория логов — `DEBUG_CATEGORY_MEMBER_SYNC`. |
||||
Категория member_sync: ошибки содержат канал, node_id и причину отказа. |
||||
DEBUG показывает применение версий и размеры страниц, TRACE — входы и расчёт листа. |
||||
Хеши родителей вычисляет merkle_tree; адаптер никогда не перечисляет все записи для внутреннего узла. |
||||
|
||||
@ -1,319 +1,61 @@
|
||||
#ifndef MERKLE_SYNC_H |
||||
#define MERKLE_SYNC_H |
||||
|
||||
#include <stdint.h> |
||||
#include "merkle_tree.h" |
||||
#include <stddef.h> |
||||
#include <openssl/evp.h> |
||||
|
||||
struct UTUN_INSTANCE; |
||||
|
||||
#define MT_MAX_LEVEL 5 |
||||
#define MT_BUCKETS 32 |
||||
#define MT_HASH_SIZE 32 |
||||
|
||||
/*
|
||||
* ── Архитектура ── |
||||
* |
||||
* merkle_sync — универсальный протокол синхронизации на Merkle-деревьях. |
||||
* Группирует элементы (items) в префиксное дерево: 5 уровней × 32 бакета, |
||||
* SHA256-хеши. Два пира обмениваются хешами уровней, находят различающиеся |
||||
* бакеты и передают только их содержимое. |
||||
* |
||||
* merkle_sync не знает, что такое "элемент" — потребитель (member_sync) |
||||
* предоставляет четыре коллбэка, описывающих модель данных: |
||||
* |
||||
* update_bucket_hash — перечислить элементы в диапазоне ключей, |
||||
* вычислить хеш каждого и подать в SHA256-контекст |
||||
* get_items — сериализовать элементы бакета в wire-формат |
||||
* apply_items — десериализовать и сохранить полученные элементы |
||||
* apply_update — применить лёгкое обновление (MSG_ITEM_UPDATE) |
||||
* |
||||
* Потребитель (chat_sync/chat_core) работает только с member_sync |
||||
* и не видит merkle_sync напрямую: |
||||
* |
||||
* // ── разово ──
|
||||
* member_sync_init(inst); |
||||
* |
||||
* // ── запустил синхронизацию — забыл ──
|
||||
* member_sync_start(inst, peer, ch_id, on_done, my_ctx); |
||||
* |
||||
* // ── получил результат ──
|
||||
* static void on_done(uint64_t peer, const char* ch_id, int result, void* arg) { |
||||
* if (result == MT_OK) printf("sync ok\n"); |
||||
* if (result == MT_ERR_TIMEOUT) printf("timeout\n"); |
||||
* } |
||||
* |
||||
* // ── изменил данные — дерево само пересчиталось ──
|
||||
* member_sync_put(inst, ch_id, node_id, x25519, ed25519, join_sig, addrs, ac); |
||||
* |
||||
* // ── отменил — коллбэк не вызовется ──
|
||||
* member_sync_cancel(inst, peer, ch_id); |
||||
* |
||||
* // ── разово ──
|
||||
* member_sync_destroy(inst); |
||||
* |
||||
* ── Формат дерева ── |
||||
* |
||||
* 5-уровневое префиксное дерево над 64-битными ключами (node_id). |
||||
* Каждый уровень берёт 5 старших бит ключа: уровень 1 — биты 59-63, |
||||
* уровень 5 — биты 39-63. Каждый узел дерева разбивается на 32 бакета |
||||
* (по 5 бит = 32 комбинации). |
||||
* |
||||
* level 1 [0..31] корень: 32 бакета по 5 бит |
||||
* level 2 [0..31]...[0..31] каждый — ещё 32 бакета |
||||
* ... |
||||
* level 5 [0..31].........[0..31] листья: 32^5 = 33M бакетов макс |
||||
* |
||||
* Хеш бакета = SHA256(хеш_элемента_1 || ... || хеш_элемента_N). |
||||
* Бакет считается терминальным (leaf) на уровне 5 или если в нём < 8 элементов. |
||||
* |
||||
* ── Wire-протокол (сервис ETCP, id задаётся при init) ── |
||||
* |
||||
* Каждое сообщение: [svc_id:1][ns_len:1][ns:var][type:1][payload:var] |
||||
* |
||||
* MSG_HASHES (0x01): level,prefix,is_data, [bitmap+hashes | member_data] |
||||
* MSG_REQUEST (0x02): count, [level,prefix,is_terminal]* |
||||
* MSG_BATCH (0x03): count, [level,prefix,is_terminal,[member_data|bitmap+hashes]]* |
||||
* |
||||
* Алгоритм: |
||||
* A → MSG_HASHES(level=1, bitmap+hashes всех 32 бакетов уровня 2) → B |
||||
* B сравнивает со своим деревом, находит различающиеся бакеты |
||||
* B → MSG_REQUEST(level=2, prefix=X, is_terminal) → A |
||||
* A → MSG_BATCH(данные бакета) → B |
||||
* B сохраняет, пересчитывает хеши |
||||
* Рекурсивно для подбакетов, пока хеши не совпадут. |
||||
* |
||||
* ── Сессии ── |
||||
* |
||||
* Протокол работает поверх надёжного транспорта (ETCP). |
||||
* Таймаутов и ретраев нет — при ошибке сессия завершается. |
||||
*/ |
||||
|
||||
/* ── Data model callbacks ── */ |
||||
|
||||
#define MT_OK 0 |
||||
#define MT_ERR_IO -1 |
||||
#define MT_ERR_PROTOCOL -2 |
||||
#define MT_ERR_DATA -3 |
||||
#define MT_ERR_CONFLICT -4 |
||||
#define MT_ERR_DISCONNECTED -5 |
||||
#define MT_PAGE_SIZE 16384 |
||||
|
||||
/* Namespace — каноническая десятичная строка uint64_t. Движок не знает CHAT-групп.
|
||||
* SHA256 листа обновляется моделью в порядке unsigned key; родителей считает merkle_tree. |
||||
* Все операции выполняются в uasync-потоке экземпляра. */ |
||||
struct merkle_sync_data_ops { |
||||
/*
|
||||
* Подать хеши всех элементов в префиксном диапазоне в sha_ctx. |
||||
* |
||||
* ctx — data_ctx, переданный в merkle_sync_init |
||||
* ns — namespace (например channel_id) |
||||
* level — уровень дерева (1..5) |
||||
* prefix64 — префикс ключа (старшие level*5 бит, остальные нули) |
||||
* sha_ctx — уже инициализирован (EVP_DigestInit_ex), потребитель |
||||
* делает только EVP_DigestUpdate(sha_ctx, item_hash, 32) |
||||
* для каждого элемента |
||||
* |
||||
* Возвращает количество элементов (0 = бакет пуст, -1 = ошибка). |
||||
* |
||||
* Пример реализации для мемберов: |
||||
* SELECT ... FROM peers_<ns> JOIN nodes |
||||
* WHERE (node_id & mask) == prefix64 ORDER BY node_id |
||||
* для каждой строки: _compute_member_hash() → EVP_DigestUpdate() |
||||
*/ |
||||
int (*update_bucket_hash)(void* ctx, const char* ns, uint8_t level, |
||||
uint64_t prefix64, EVP_MD_CTX* sha_ctx); |
||||
|
||||
/*
|
||||
* Сериализовать элементы бакета в wire-формат. |
||||
* |
||||
* buf — буфер для записи (выделяет merkle_sync, мин. 64K) |
||||
* *len — [in] размер буфера, [out] записанный размер |
||||
* |
||||
* Wire-формат для мемберов: |
||||
* [count:2][node_id:8][x25519:32][ed25519:32][flags:1]([join_sig:64][join_ts:8])[update_sig:64][update_ts:8][userinfo_len:1][userinfo:var][adm_tags_len:1][adm_tags:var][adm_tags_sig:64]... |
||||
* |
||||
* Возвращает 0 при успехе, <0 при ошибке, -2 если буфер мал. |
||||
*/ |
||||
int (*get_items)(void* ctx, const char* ns, uint8_t level, |
||||
uint64_t prefix, uint8_t pbytes, uint8_t* buf, size_t* len); |
||||
|
||||
/*
|
||||
* Десериализовать и сохранить элементы из wire-формата. |
||||
* |
||||
* data, len — данные в том же формате, что выдаёт get_items |
||||
* |
||||
* Вызывается при приёме MSG_HASHES(is_data=1) или MSG_BATCH(is_terminal=1). |
||||
* Должна сохранить элементы в БД и вызвать merkle_sync_recompute_path |
||||
* для каждого изменённого ключа (через member_sync_put). |
||||
* |
||||
* Возвращает 0 при успехе, <0 при ошибке. |
||||
*/ |
||||
int (*apply_items)(void* ctx, const char* ns, uint64_t from_peer, |
||||
const uint8_t* data, size_t len); |
||||
|
||||
/*
|
||||
* Применить лёгкое обновление (не влияющее на Merkle-хеш). |
||||
* Вызывается при приёме MSG_ITEM_UPDATE(0x04) — online-статус и т.п. |
||||
* |
||||
* ns — namespace (channel_id) |
||||
* key — node_id изменившегося узла |
||||
* type — тип обновления (0x01 = online-статус) |
||||
* data,len — данные обновления (type=0x01: [online:1]) |
||||
* |
||||
* Возвращает 0 при успехе, <0 при ошибке. |
||||
*/ |
||||
int (*apply_update)(void* ctx, const char* ns, uint64_t key, |
||||
uint8_t type, const uint8_t* data, size_t len); |
||||
|
||||
/*
|
||||
* Проверить, что пир авторизован для namespace ns (напр. его pubkey есть |
||||
* в таблице мемберов канала). Вызывается в _recv_cb для КАЖДОГО входящего |
||||
* пакета до диспатча. Возвращает 1 если разрешено, 0 если нет. |
||||
* Может быть NULL — тогда проверка пропускается (generic-потребитель). |
||||
*/ |
||||
merkle_leaf_hash_fn update_bucket_hash; |
||||
/* Одна страница листа: after_valid=0 начинает обход, иначе только key > after.
|
||||
* len: ёмкость -> размер; next — последний ключ; more — есть продолжение. |
||||
* Пустой результат разрешён только с more=0. Ошибка не кодируется пустой страницей. */ |
||||
int (*get_page)(void* ctx, const char* ns, uint64_t prefix, int after_valid, uint64_t after, |
||||
uint8_t* buf, size_t* len, uint64_t* next, int* more); |
||||
/* Проверить весь формат, применить записи и обновить дерево в одной DB-транзакции.
|
||||
* <0 — ошибка. Никаких broadcast/send-back из этого callback. */ |
||||
int (*apply_items)(void* ctx, const char* ns, uint64_t from_peer, const uint8_t* data, size_t len); |
||||
int (*validate_peer)(void* ctx, const char* ns, uint64_t peer); |
||||
/* Проверить отложенные зависимости модели после обоих проходов, до подтверждения корня. */ |
||||
int (*finish)(void* ctx, const char* ns); |
||||
}; |
||||
|
||||
/*
|
||||
* Коллбэк завершения синхронизации. |
||||
* |
||||
* peer — node_id пира |
||||
* ns — namespace (channel_id) |
||||
* result — MT_OK (0) = успех, MT_ERR_TIMEOUT (-1) = зарезервировано (в протоколе не возникает) |
||||
* arg — пользовательский контекст |
||||
* |
||||
* При отмене через merkle_sync_cancel() коллбэк НЕ вызывается. |
||||
*/ |
||||
|
||||
#define MT_OK 0 |
||||
#define MT_ERR_TIMEOUT -1 |
||||
|
||||
typedef void (*merkle_sync_done_cb)(uint64_t peer, const char* ns, int result, void* arg); |
||||
|
||||
/* ── Lifecycle ── */ |
||||
|
||||
/*
|
||||
* Инициализировать модуль. Создаёт таблицу merkle_tree_hash, биндит |
||||
* ETCP-сервис svc_id (напр. 0x31), запускает фоновую проверку (bg_timer). |
||||
* Вызывается один раз, обычно из member_sync_init(). |
||||
* |
||||
* inst — экземпляр uTun (нужен для etcp_send, uasync, sqlite3) |
||||
* svc_id — ID сервиса на ETCP-роутере |
||||
* ops — коллбэки модели данных (update_bucket_hash, get_items, apply_items) |
||||
* data_ctx — прозрачный контекст, передаваемый в коллбэки первым аргументом |
||||
* |
||||
* Возвращает 0 при успехе, -1 при ошибке. |
||||
*/ |
||||
int merkle_sync_init(struct UTUN_INSTANCE* inst, uint8_t svc_id, |
||||
const struct merkle_sync_data_ops* ops, void* data_ctx); |
||||
|
||||
/*
|
||||
* Завершить модуль. Отвязывает ETCP-сервис, останавливает bg_timer, |
||||
* удаляет все активные сессии (коллбэки НЕ вызываются). |
||||
*/ |
||||
void merkle_sync_destroy(struct UTUN_INSTANCE* inst); |
||||
|
||||
/* ── Async sync ── */ |
||||
|
||||
/*
|
||||
* Запустить синхронизацию namespace ns с пиром peer. |
||||
* Отправляет MSG_HASHES(level=0). Таймаута нет — сессия завершается при |
||||
* совпадении корневого хеша либо удаляется вручную (cancel/destroy). |
||||
* |
||||
* Если сессия для (peer, ns) уже существует — перезапускает её |
||||
* с новым коллбэком (старый коллбэк теряется без вызова). |
||||
* |
||||
* Возвращает 0 при успехе, -1 при ошибке. |
||||
*/ |
||||
int merkle_sync_start(struct UTUN_INSTANCE* inst, uint64_t peer, |
||||
const char* ns, merkle_sync_done_cb done_cb, void* arg); |
||||
|
||||
/*
|
||||
* Отменить синхронизацию. Удаляет сессию, done_cb НЕ вызывается. |
||||
* Безопасно вызывать, если сессия не существует. |
||||
*/ |
||||
/* Подключить namespace к текущему ETCP-соединению и запросить контрольную точку.
|
||||
* Успех — подтверждение одинакового корня конкретного раунда, не вечный статус. |
||||
* Повторные запросы объединяются; каждый callback вызывается один раз, обход не сбрасывается. |
||||
* После первого start изменения дерева автоматически синхронизируются с этим пиром. |
||||
* <0: запрос не принят, callback не вызывается. DOWN/DELETE завершают запросы ошибкой. |
||||
* Явные cancel/destroy освобождают запросы без callback. |
||||
* Callback может отменить сессию; уничтожение всего экземпляра должно быть отложенным. */ |
||||
int merkle_sync_start(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ns, merkle_sync_done_cb cb, void* arg); |
||||
void merkle_sync_cancel(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ns); |
||||
|
||||
/*
|
||||
* Отменить и удалить ВСЕ сессии для namespace ns (независимо от пира). |
||||
* Используется при локальном удалении канала. done_cb НЕ вызывается. |
||||
*/ |
||||
void merkle_sync_cancel_ns(struct UTUN_INSTANCE* inst, const char* ns); |
||||
|
||||
/*
|
||||
* Отменить и удалить ВСЕ сессии с конкретным пиром peer (независимо от ns). |
||||
* Используется при уходе пира в спячку (server-side throttling). done_cb НЕ вызывается. |
||||
*/ |
||||
void merkle_sync_cancel_peer(struct UTUN_INSTANCE* inst, uint64_t peer); |
||||
|
||||
/* ── Recompute tree path after data change ── */ |
||||
|
||||
/*
|
||||
* Пересчитать Merkle-дерево для ключа key — все уровни от 1 до 5. |
||||
* Вызывается после изменения данных: merkle_sync сам не знает, |
||||
* когда данные изменились — потребитель должен вызвать явно. |
||||
* member_sync делает это внутри put/del. |
||||
* |
||||
* Возвращает 1 если хеш хотя бы одного бакета изменился (данные реально |
||||
* изменились), 0 если всё совпало, -1 при ошибке. |
||||
*/ |
||||
/* Обновление индекса вызывается моделью данных.
|
||||
* При внешней транзакции caller вызывает changed ПОСЛЕ commit. |
||||
* Без внешней транзакции пересчёт сам создаёт savepoint и уведомляет после commit. |
||||
* Возвращает 1=изменилось, 0=no-op, <0=ошибка. */ |
||||
int merkle_sync_recompute_path(struct UTUN_INSTANCE* inst, const char* ns, uint64_t key); |
||||
void merkle_sync_changed(struct UTUN_INSTANCE* inst, const char* ns); |
||||
int merkle_sync_read_hash(struct UTUN_INSTANCE* inst, const char* ns, uint8_t level, uint64_t prefix, uint8_t hash[MT_HASH_SIZE]); |
||||
|
||||
/*
|
||||
* Отправить лёгкое обновление всем synced-пирам в namespace ns. |
||||
* Не влияет на Merkle-дерево — используется для online-статуса и т.п. |
||||
* |
||||
* ns — namespace (channel_id) |
||||
* key — node_id изменившегося узла |
||||
* type — тип обновления (0x01 = online-статус) |
||||
* data,len — данные обновления |
||||
*/ |
||||
void merkle_sync_push_update(struct UTUN_INSTANCE* inst, const char* ns, |
||||
uint64_t key, uint8_t type, const uint8_t* data, size_t len); |
||||
|
||||
/*
|
||||
* Переслать данные всем SYNCED-сессиям (кроме from_peer) в namespace ns. |
||||
* Используется для relay: когда apply_items обнаружил реальные изменения, |
||||
* потребитель вызывает эту функцию чтобы разослать дельту остальным пирам. |
||||
* |
||||
* data, len — wire-формат [count:2][key:8][val:4]... как от get_items. |
||||
*/ |
||||
void merkle_sync_broadcast(struct UTUN_INSTANCE* inst, const char* ns, |
||||
uint64_t from_peer, const uint8_t* data, size_t len); |
||||
|
||||
/*
|
||||
* Отправить данные одного мембера конкретному пиру (а не всем). |
||||
* Используется для "send back" при обнаружении stale-версии. |
||||
* data, len — wire-формат [count:2][key:8][val:4]... как от get_items. |
||||
* Возвращает 0 при успехе, -1 если нет сессии/соединения к пиру. |
||||
*/ |
||||
int merkle_sync_send_to(struct UTUN_INSTANCE* inst, const char* ns, |
||||
uint64_t peer, const uint8_t* data, size_t len); |
||||
|
||||
/* ── Background consistency check ── */ |
||||
|
||||
/*
|
||||
* Проверить консистентность одного случайного level-5 бакета в ns. |
||||
* Пересчитывает хеш из данных, сравнивает с сохранённым в merkle_tree_hash. |
||||
* При расхождении — пересчитывает весь путь вверх до корня. |
||||
* Возвращает 1 если были изменения, 0 если совпало, -1 ошибка. |
||||
*/ |
||||
int merkle_sync_bg_check(struct UTUN_INSTANCE* inst, const char* ns); |
||||
|
||||
/* ── Test helper ── */ |
||||
|
||||
/*
|
||||
* Получить сохранённый хеш бакета из merkle_tree_hash. |
||||
* Возвращает указатель на статический буфер (32 байта), перезаписывается |
||||
* при следующем вызове. Если бакет не найден — возвращает нулевой хеш. |
||||
*/ |
||||
const uint8_t* merkle_sync_get_hash(struct UTUN_INSTANCE* inst, |
||||
const char* ns, uint8_t level, uint64_t prefix64); |
||||
|
||||
/* ── Prefix arithmetic (pure, exported for convenience) ── */ |
||||
|
||||
/*
|
||||
* Выделить level*5 старших бит из 64-битного ключа. |
||||
* Уровень 1: биты 59-63 (сдвиг 59) |
||||
* Уровень 5: биты 39-63 (сдвиг 39) |
||||
* Пример: merkle_sync_level_prefix(0x1234567890ABCDEF, 1) = 0x1000000000000000 |
||||
*/ |
||||
uint64_t merkle_sync_level_prefix(uint64_t key, uint8_t level); |
||||
|
||||
/*
|
||||
* Количество байт, необходимое для хранения префикса уровня level. |
||||
* Уровень 1: 1 байт (5 бит), уровень 4: 3 байта (20 бит), уровень 5: 4 байта (25 бит). |
||||
*/ |
||||
uint8_t merkle_sync_prefix_bytes(uint8_t level); |
||||
|
||||
#endif /* MERKLE_SYNC_H */ |
||||
#endif |
||||
|
||||
@ -1,158 +1,100 @@
|
||||
# merkle_sync — синхронизация данных на Merkle-деревьях |
||||
# merkle_sync — протокол 2 |
||||
|
||||
## 1. Назначение |
||||
## Разделение ответственности |
||||
|
||||
Универсальный движок синхронизации между двумя узлами. Группирует элементы в префиксное |
||||
Merkle-дерево (5 уровней × 32 бакета, SHA256). Два пира обмениваются хешами уровней, находят |
||||
различающиеся бакеты и передают **только их содержимое** — вместо полной пересылки всех данных. |
||||
- merkle_tree: производный индекс SQLite, хеши листьев и родителей. |
||||
- merkle_sync: обмен, backpressure и жизненный цикл сессий. |
||||
- member_sync: формат записей, подписи, версии, применение страниц. |
||||
|
||||
merkle_sync **не знает**, что такое «элемент». Потребитель (member_sync) предоставляет набор |
||||
коллбэков `merkle_sync_data_ops`, описывающих модель данных (как перечислить элементы бакета |
||||
и посчитать хеш, как сериализовать элементы в wire-формат, как применить полученные). |
||||
Движок не знает CHAT-групп. Namespace — каноническая десятичная строка uint64_t. |
||||
Все операции выполняются в одном потоке uasync экземпляра. |
||||
|
||||
Работает поверх надёжного транспорта (ETCP), поэтому **таймаутов и ретраев в протоколе нет**: |
||||
сессия завершается при совпадении корневого хеша, либо удаляется вручную (`cancel`/`destroy`). |
||||
## Дерево и транзакции |
||||
|
||||
## 2. Как пользоваться |
||||
Корень имеет уровень 0. Уровни 1–5 разбивают старшие 25 бит ключа на группы по 5 бит. |
||||
SHA256 листа: `0x4d || level:1 || item_hashes`, записи упорядочены по unsigned key. |
||||
SHA256 родителя: `0x4d || level:1 || child_hash[0] || ... || child_hash[31]`. |
||||
Хеш пустого поддерева — 32 нулевых байта; строка индекса для него отсутствует. |
||||
Пересчитывается один лист и его предки, до первого неизменившегося хеша. |
||||
|
||||
Напрямую merkle_sync обычно не используют — его потребляет member_sync. Типовой сценарий |
||||
потребителя выглядит так: |
||||
Индекс создаётся заново при init и заполняется моделью из исходных записей. |
||||
Модель меняет данные и вызывает recompute_path в одной транзакции, затем changed после commit. |
||||
Без внешней транзакции recompute_path создаёт savepoint и уведомляет самостоятельно. |
||||
Ошибки не заменяются пустыми страницами или нулевыми хешами. |
||||
|
||||
```c |
||||
// ── разово: инициализация (из member_sync_init) ── |
||||
merkle_sync_init(inst, 0x31, &ops, data_ctx); // svc_id — ID ETCP-сервиса |
||||
## API |
||||
|
||||
// ── запустил синхронизацию — забыл ── |
||||
merkle_sync_start(inst, peer, ns, on_done, my_arg); |
||||
init принимает service и data_ops: update_bucket_hash, get_page, apply_items, |
||||
необязательные validate_peer и finish. |
||||
update_bucket_hash добавляет хеши одного листа в переданный SHA256-контекст. |
||||
get_page получает ёмкость буфера и курсор; возвращает размер, последний ключ и флаг продолжения. |
||||
apply_items проверяет весь формат и применяет страницу вместе с индексом атомарно. |
||||
finish проверяет отложенные зависимости модели перед подтверждением корня. |
||||
|
||||
// ── коллбэк завершения ── |
||||
static void on_done(uint64_t peer, const char* ns, int result, void* arg) { |
||||
if (result == MT_OK) /* хеши сошлись, данные идентичны */; |
||||
if (result == MT_ERR_TIMEOUT) /* (зарезервировано, в протоколе не возникает) */; |
||||
} |
||||
start подключает namespace к текущему соединению и запрашивает контрольную точку. |
||||
Повторные start объединяются, обход не сбрасывается, каждый принятый callback вызывается один раз. |
||||
После подключения изменения синхронизируются автоматически. Broadcast/send-back не требуются. |
||||
MT_OK означает совпадение корней конкретного раунда, а не отсутствие будущих изменений. |
||||
|
||||
// ── изменил данные — дерево само пересчиталось ── |
||||
int changed = merkle_sync_recompute_path(inst, ns, key); // 1 = изменилось, 0 = no-op |
||||
cancel, cancel_ns, cancel_peer и destroy освобождают запросы без callback. |
||||
DOWN/DELETE/REINIT завершают запросы ошибкой и освобождают сессию. |
||||
После восстановления соединения интеграционный слой повторяет start. |
||||
Callback может отменить сессию; уничтожение всего экземпляра следует отложить через call_soon. |
||||
|
||||
// ── рассылка дельты остальным пирам (relay) ── |
||||
merkle_sync_broadcast(inst, ns, from_peer, data, len); |
||||
merkle_sync_push_update(inst, ns, key, type, data, len); // лёгкое обновление (online и т.п.) |
||||
## Обмен |
||||
|
||||
// ── отменил — коллбэк не вызовется ── |
||||
merkle_sync_cancel(inst, peer, ns); |
||||
Меньший node_id ведёт раунд; второй запрашивает его через WAKE. |
||||
WAKE также перезапускает незавершённый раунд после cancel/start второго узла. |
||||
|
||||
// ── разово: завершение ── |
||||
merkle_sync_destroy(inst); |
||||
``` |
||||
1. BEGIN(new round) → READY. |
||||
2. Ведущий запрашивает хеши детей и обходит различающиеся ветви в глубину. |
||||
3. Лист запрашивается страницами с возрастающим курсором. |
||||
4. TURN меняет направление: второй узел выполняет такой же pull. |
||||
5. Второй проверяет finish и отправляет CHECK(root). |
||||
6. Ведущий сверяет корень, проверяет finish и отправляет DONE(root). |
||||
7. Второй сверяет DONE со своим текущим корнем. |
||||
|
||||
## 3. Модель данных (коллбэки `merkle_sync_data_ops`) |
||||
Ведущий сообщает MT_OK после передачи DONE в надёжный транспорт, если локальная ревизия |
||||
не изменилась после проверки. Второй сообщает MT_OK после проверки полученного DONE. |
||||
Изменения во время обмена вызывают следующий раунд. |
||||
Два завершённых раунда с одной и той же несовпадающей парой корней дают MT_ERR_CONFLICT. |
||||
Детерминированное слияние обеспечивает модель; движок не выбирает победителя конфликта. |
||||
|
||||
| Коллбэк | Что делает | |
||||
|---|---| |
||||
| `update_bucket_hash(ctx, ns, level, prefix64, sha_ctx)` | Перечислить элементы в префиксном диапазоне, вычислить хеш каждого и подать в `sha_ctx` через `EVP_DigestUpdate(sha_ctx, item_hash, 32)`. Возвращает число элементов (0 = бакет пуст, -1 = ошибка) | |
||||
| `get_items(ctx, ns, level, prefix, pbytes, buf, len)` | Сериализовать элементы бакета в wire-формат. Возвращает 0 / -2 (буфер мал) / <0 (ошибка) | |
||||
| `apply_items(ctx, ns, from_peer, data, len)` | Десериализовать и сохранить полученные элементы. Должна пересчитать дерево через `merkle_sync_recompute_path` для каждого изменённого ключа | |
||||
| `apply_update(ctx, ns, key, type, data, len)` | Применить лёгкое обновление (не влияет на Merkle-хеш) — MSG_ITEM_UPDATE | |
||||
|
||||
`ctx` — прозрачный контекст, передаваемый в `merkle_sync_init` и возвращаемый первым аргументом |
||||
в каждый коллбэк. `ns` (namespace) — строка-имя домена данных, обычно `channel_id`. |
||||
|
||||
## 4. Формат дерева |
||||
|
||||
5-уровневое префиксное дерево над 64-битными ключами. Каждый уровень берёт 5 старших бит ключа: |
||||
уровень 1 — биты 59-63, уровень 5 — биты 39-63. Каждый узел разбивается на 32 бакета (5 бит = 32). |
||||
|
||||
``` |
||||
level 1 [0..31] корень: 32 бакета по 5 бит |
||||
level 2 [0..31]...[0..31] каждый — ещё 32 бакета |
||||
... |
||||
level 5 [0..31].........[0..31] листья: 32^5 = 33M бакетов макс |
||||
``` |
||||
|
||||
- Хеш бакета = `SHA256(хеш_элемента_1 || ... || хеш_элемента_N)`. |
||||
- Бакет считается **терминальным** (leaf) на уровне 5 или если в нём `< 8` элементов. |
||||
- Хеши бакетов хранятся в таблице `merkle_tree_hash(namespace, level, prefix64, hash, member_count)`. |
||||
|
||||
Вспомогательные функции префикса: `merkle_sync_level_prefix(key, level)` (выделить `level*5` |
||||
старших бит) и `merkle_sync_prefix_bytes(level)` (байт на хранение префикса: уровень 1 → 1 байт, |
||||
уровень 5 → 4 байта). |
||||
|
||||
## 5. Wire-протокол |
||||
|
||||
Каждое сообщение: `[svc_id:1][ns_len:1][ns:var][type:1][payload:var]`. |
||||
|
||||
| Тип | #define | Назначение | |
||||
|---|---|---| |
||||
| 0x01 | MSG_HASHES | Сравнение бакетов: bitmap + хеши, либо данные (is_data) | |
||||
| 0x02 | MSG_REQUEST | Запрос содержимого различающихся бакетов | |
||||
| 0x03 | MSG_BATCH | Передача данных/подхешей запрошенных бакетов | |
||||
| 0x04 | MSG_ITEM_UPDATE | Лёгкое обновление (online-статус), вне дерева | |
||||
|
||||
Форматы payload: |
||||
- **MSG_HASHES**: `[level:1][prefix_bytes:1][prefix:pb][is_data:1]` + (`[bitmap:4][hash:32]*` | `[dlen:2][data]`) |
||||
- **MSG_REQUEST**: `[count:1][ {level:1, prefix_bytes:1, prefix:pb, is_terminal:1} × count ]` |
||||
- **MSG_BATCH**: `[count:1][ {level:1, prefix_bytes:1, prefix:pb, is_terminal:1, [data | bitmap+hashes]} × count ]` |
||||
- **MSG_ITEM_UPDATE**: `[key:8][type:1][data]` |
||||
|
||||
### Алгоритм обмена |
||||
|
||||
``` |
||||
A → MSG_HASHES(level=0, bitmap+хеши всех 32 бакетов уровня 1) → B |
||||
B сравнивает со своим деревом, находит различающиеся бакеты (differs) |
||||
B → MSG_REQUEST(level=2, prefix=X, is_terminal) → A |
||||
A → MSG_BATCH(данные бакета или подхеши) → B |
||||
B сохраняет, пересчитывает хеши |
||||
Рекурсивно для подбакетов, пока хеши не совпадут. |
||||
differs==0 и level==0 → корневой хеш совпал → сессия SYNCED. |
||||
``` |
||||
|
||||
## 6. Сессии и out-очередь (backpressure) |
||||
|
||||
`struct ms_session` — сессия на пару (peer, ns): `sess_state` (`SESS_SYNCING`/`SESS_SYNCED`), |
||||
`pending[128]` (ожидающие бакеты: `MS_PEND_WAITING`/`MS_PEND_RESOLVED`, с родительским индексом |
||||
и счётчиком подбакетов), `done_cb`/`cb_arg`, `out_q` + `waiter`. |
||||
|
||||
**По-пировая очередь рассылки** `out_q`: рассылка (`broadcast`/`push_update`/`send_to`) не шлёт |
||||
немедленно, а кладёт каждый элемент в `out_q` подходящей сессии. Дренаж через `ms_outq_drain_cb`: |
||||
`queue_data_get(out_q)` → `etcp_send(conn, entry)`. Перед отправкой проверяется занятость |
||||
`conn->send_input_q`; при превышении порога — `queue_waiter_wait(send_input_q, &s->waiter, ...)` |
||||
приостанавливает дренаж до освобождения (backpressure, без дропов). |
||||
|
||||
Отправка: |
||||
- `_send_msg` — синхронно одному пиру (служебные HASHES/REQUEST/BATCH — по прямому соединению). |
||||
- `_broadcast_to_session` — кладёт в `out_q` конкретной сессии (обёртка в MSG_HASHES is_data=1). |
||||
|
||||
## 7. WAL-таймер |
||||
Ожидается только один ответ QUERY. Стек — максимум шесть кадров с сохранёнными соседними ветвями. |
||||
Номера запросов проверяются, ответы другого раунда игнорируются. |
||||
Повторную передачу обеспечивает ETCP, отдельных таймеров повторов нет. |
||||
|
||||
Отдельный редкий таймер `MS_WAL_INTERVAL_MS = 60000` (60 с) — `sqlite3_wal_checkpoint_v2(..., |
||||
PASSIVE)`. Локальная операция без сети. На Android (`UTUN_HAVE_STANDBY`) при спячке откладывается |
||||
через `standby_wait`. Периодической сверки деревьев (bg-таймера) нет — пересчёт событийный. |
||||
## Wire |
||||
|
||||
## 8. API |
||||
Заголовок: `service:1 | namespace:8 | type:1 | round:8 | request:4` (22 байта). |
||||
Все целые big-endian. Максимальный пакет вместе с заголовком — 16384 байта. |
||||
Типы 0x10–0x19: WAKE, BEGIN, READY, QUERY, HASHES, PAGE, TURN, CHECK, DONE, ERROR. |
||||
|
||||
| Функция | Назначение | |
||||
| Тип | Payload | |
||||
|---|---| |
||||
| `merkle_sync_init(inst, svc_id, ops, data_ctx)` | Создать контекст, `_ensure_table`, `etcp_bind(svc_id, _recv_cb)`, WAL-таймер. `inst->msync = ms` | |
||||
| `merkle_sync_destroy(inst)` | Отвязать сервис, отменить WAL-таймер, удалить все сессии (коллбэки не вызываются) | |
||||
| `merkle_sync_start(inst, peer, ns, done_cb, arg)` | Создать/перезапустить сессию, отправить MSG_HASHES(level=0). Перезапуск теряет старый коллбэк без вызова | |
||||
| `merkle_sync_cancel(inst, peer, ns)` | Удалить сессию, `done_cb` не вызывается | |
||||
| `merkle_sync_recompute_path(inst, ns, key)` | Пересчитать все 5 уровней для ключа. Возвращает 1 (изменилось) / 0 (no-op) / -1 | |
||||
| `merkle_sync_push_update(inst, ns, key, type, data, len)` | MSG_ITEM_UPDATE всем **SYNCED**-сессиям ns | |
||||
| `merkle_sync_broadcast(inst, ns, from_peer, data, len)` | Relay: всем сессиям ns, кроме `from_peer` | |
||||
| `merkle_sync_send_to(inst, ns, peer, data, len)` | Данные одному пиру (send-back при stale) | |
||||
| `merkle_sync_bg_check(inst, ns)` | Проверить один level-5 бакет (сравнение с БД), при расхождении — пересчёт пути вверх. 1/0/-1 | |
||||
| `merkle_sync_get_hash(inst, ns, level, prefix64)` | Прочитать хеш бакета (статический буфер, нули если нет) | |
||||
| `merkle_sync_level_prefix(key, level)` / `merkle_sync_prefix_bytes(level)` | Арифметика префикса | |
||||
|
||||
## 9. Нюансы |
||||
|
||||
- **svc_id** задаётся при init (member_sync использует `0x31`); единственный приёмник `_recv_cb`. |
||||
- **Нет таймаута сессии**: `MT_ERR_TIMEOUT` зарезервирован, но не возникает — сессия живёт до |
||||
совпадения корневого хеша или явной отмены. |
||||
- **Перезапуск**: `merkle_sync_start` на существующую пару (peer, ns) перезаписывает `done_cb` |
||||
(старый коллбэк теряется без вызова). |
||||
- **Сессия создаётся и на приёме**: при неожиданном MSG_HASHES от неизвестного пира сессия |
||||
создаётся автоматически (symmetrical start). |
||||
- **Результат recompute** важен для подавления no-op рассылки: `changed==0` означает, что данные |
||||
фактически не изменились. |
||||
- Категория логов — `DEBUG_CATEGORY_MEMBER_SYNC`. |
||||
| WAKE, BEGIN, READY, TURN | пустой; request=0; у WAKE также round=0 | |
||||
| QUERY | level:1, prefix:8, after_valid:1, after:8 | |
||||
| HASHES | 32 хеша детей по 32 байта в порядке индексов | |
||||
| PAGE | more:1, next:8, model_data | |
||||
| CHECK, DONE | корень 32 байта; request=0 | |
||||
| ERROR | положительный код ошибки 4 байта; request=0 | |
||||
|
||||
Курсор продолжения должен возрастать и принадлежать запрошенному листу. |
||||
Длины, состояние, тип, префикс и номер запроса проверяются до применения. |
||||
Формат model_data определяет адаптер. Старый протокол не поддерживается. |
||||
|
||||
## Очереди, диагностика, тесты |
||||
|
||||
Сессия удерживает ETCP_CONN и исходную очередь. У неё один подготовленный пакет, |
||||
один waiter и один отложенный вызов. Передача идёт только после разрешения waiter. |
||||
Изменения объединяются флагом dirty; неограниченной out-очереди нет. |
||||
Отмена снимает ожидающий и отложенный callback waiter независимо от links_up. |
||||
|
||||
Категория member_sync: INFO — attach, начало/завершение раунда, detach; |
||||
DEBUG — TX/RX, номера запросов, размеры, состояние; ERROR — причина отказа. |
||||
Ведущий выводит счётчики переданных и полученных байтов при подтверждении. |
||||
|
||||
test_merkle_sync проверяет индекс, односторонний start, двустороннее слияние, несколько узлов, |
||||
параллельные изменения, большой лист с пагинацией и автоматическое обновление. |
||||
test_merkle_protocol проверяет ошибочные пакеты/курсоры, сохранение соседних ветвей, |
||||
отмену waiter, ошибки apply/SQLite, backpressure, disconnect и повторное подключение. |
||||
|
||||
@ -0,0 +1,122 @@
|
||||
#include "merkle_tree.h" |
||||
#include "../../lib/debug_config.h" |
||||
#include <string.h> |
||||
|
||||
static int tree_error(sqlite3* db, const char* op, const char* ns) { |
||||
DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "merkle_tree: %s ns=%s: %s", op, ns ? ns : "-", db ? sqlite3_errmsg(db) : "no DB"); |
||||
return -1; |
||||
} |
||||
|
||||
uint64_t merkle_sync_level_prefix(uint64_t key, uint8_t level) { |
||||
if (level == 0) return 0; |
||||
if (level > MT_MAX_LEVEL) return 0; |
||||
return key & (~UINT64_C(0) << (64 - level * 5)); |
||||
} |
||||
|
||||
uint8_t merkle_sync_prefix_bytes(uint8_t level) { |
||||
return level <= MT_MAX_LEVEL ? (uint8_t)((level * 5 + 7) / 8) : 0; |
||||
} |
||||
|
||||
int merkle_tree_init(sqlite3* db) { |
||||
const char* sql = "DROP TABLE IF EXISTS merkle_tree_hash;" |
||||
"CREATE TABLE merkle_tree_hash(namespace TEXT NOT NULL, level INTEGER NOT NULL CHECK(level BETWEEN 0 AND 5)," |
||||
"prefix64 INTEGER NOT NULL, hash BLOB NOT NULL CHECK(length(hash)=32)," |
||||
"PRIMARY KEY(namespace,level,prefix64))"; |
||||
if (!db || sqlite3_exec(db, sql, NULL, NULL, NULL) != SQLITE_OK) return tree_error(db, "init", NULL); |
||||
DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "merkle_tree: empty derived index initialized"); |
||||
return 0; |
||||
} |
||||
|
||||
int merkle_tree_get(sqlite3* db, const char* ns, uint8_t level, uint64_t prefix, uint8_t hash[MT_HASH_SIZE]) { |
||||
if (!db || !ns || !hash || level > MT_MAX_LEVEL || prefix != merkle_sync_level_prefix(prefix, level)) |
||||
return tree_error(db, "invalid get", ns); |
||||
memset(hash, 0, MT_HASH_SIZE); |
||||
sqlite3_stmt* st = NULL; |
||||
int rc = sqlite3_prepare_v2(db, "SELECT hash FROM merkle_tree_hash WHERE namespace=? AND level=? AND prefix64=?", -1, &st, NULL); |
||||
if (rc == SQLITE_OK) rc = sqlite3_bind_text(st, 1, ns, -1, SQLITE_STATIC); |
||||
if (rc == SQLITE_OK) rc = sqlite3_bind_int(st, 2, level); |
||||
if (rc == SQLITE_OK) rc = sqlite3_bind_int64(st, 3, (sqlite3_int64)prefix); |
||||
if (rc == SQLITE_OK) rc = sqlite3_step(st); |
||||
if (rc == SQLITE_ROW && sqlite3_column_bytes(st, 0) == MT_HASH_SIZE) memcpy(hash, sqlite3_column_blob(st, 0), MT_HASH_SIZE); |
||||
else if (rc != SQLITE_DONE) { sqlite3_finalize(st); return tree_error(db, "read hash", ns); } |
||||
sqlite3_finalize(st); |
||||
return 0; |
||||
} |
||||
|
||||
int merkle_tree_children(sqlite3* db, const char* ns, uint8_t level, uint64_t prefix, |
||||
uint8_t hashes[MT_BUCKETS][MT_HASH_SIZE]) { |
||||
if (level >= MT_MAX_LEVEL || prefix != merkle_sync_level_prefix(prefix, level)) return tree_error(db, "invalid children", ns); |
||||
if (!db || !ns) return tree_error(db, "invalid children DB", ns); |
||||
memset(hashes, 0, MT_BUCKETS * MT_HASH_SIZE); |
||||
int shift = 64 - (level + 1) * 5; |
||||
const char* sql = level == 0 ? "SELECT prefix64,hash FROM merkle_tree_hash WHERE namespace=? AND level=?" |
||||
: "SELECT prefix64,hash FROM merkle_tree_hash WHERE namespace=? AND level=? AND prefix64>=? AND prefix64<=?"; |
||||
sqlite3_stmt* st = NULL; |
||||
int rc = sqlite3_prepare_v2(db, sql, -1, &st, NULL); |
||||
if (rc == SQLITE_OK) rc = sqlite3_bind_text(st, 1, ns, -1, SQLITE_STATIC); |
||||
if (rc == SQLITE_OK) rc = sqlite3_bind_int(st, 2, level + 1); |
||||
if (level && rc == SQLITE_OK) rc = sqlite3_bind_int64(st, 3, (sqlite3_int64)prefix); |
||||
if (level && rc == SQLITE_OK) rc = sqlite3_bind_int64(st, 4, (sqlite3_int64)(prefix | (UINT64_C(31) << shift))); |
||||
if (rc != SQLITE_OK) { sqlite3_finalize(st); return tree_error(db, "children query", ns); } |
||||
while ((rc = sqlite3_step(st)) == SQLITE_ROW) { |
||||
if (sqlite3_column_bytes(st, 1) != MT_HASH_SIZE) { sqlite3_finalize(st); return tree_error(db, "corrupt child", ns); } |
||||
unsigned i = (unsigned)(((uint64_t)sqlite3_column_int64(st, 0) >> shift) & 31); |
||||
memcpy(hashes[i], sqlite3_column_blob(st, 1), MT_HASH_SIZE); |
||||
} |
||||
sqlite3_finalize(st); |
||||
if (rc != SQLITE_DONE) return tree_error(db, "children read", ns); |
||||
return 0; |
||||
} |
||||
|
||||
static int tree_store(sqlite3* db, const char* ns, uint8_t level, uint64_t prefix, const uint8_t hash[MT_HASH_SIZE]) { |
||||
static const uint8_t empty[MT_HASH_SIZE]; |
||||
uint8_t old[MT_HASH_SIZE]; |
||||
if (merkle_tree_get(db, ns, level, prefix, old) < 0) return -1; |
||||
if (!memcmp(old, hash, MT_HASH_SIZE)) return 0; |
||||
int is_empty = !memcmp(hash, empty, MT_HASH_SIZE); |
||||
const char* sql = is_empty ? "DELETE FROM merkle_tree_hash WHERE namespace=? AND level=? AND prefix64=?" |
||||
: "INSERT OR REPLACE INTO merkle_tree_hash(namespace,level,prefix64,hash) VALUES(?,?,?,?)"; |
||||
sqlite3_stmt* st = NULL; |
||||
int rc = sqlite3_prepare_v2(db, sql, -1, &st, NULL); |
||||
if (rc == SQLITE_OK) rc = sqlite3_bind_text(st, 1, ns, -1, SQLITE_STATIC); |
||||
if (rc == SQLITE_OK) rc = sqlite3_bind_int(st, 2, level); |
||||
if (rc == SQLITE_OK) rc = sqlite3_bind_int64(st, 3, (sqlite3_int64)prefix); |
||||
if (rc == SQLITE_OK && !is_empty) rc = sqlite3_bind_blob(st, 4, hash, MT_HASH_SIZE, SQLITE_STATIC); |
||||
if (rc == SQLITE_OK) rc = sqlite3_step(st); |
||||
sqlite3_finalize(st); |
||||
if (rc != SQLITE_DONE) return tree_error(db, "write hash", ns); |
||||
return 1; |
||||
} |
||||
|
||||
int merkle_tree_update(sqlite3* db, const char* ns, uint64_t key, merkle_leaf_hash_fn leaf_hash, void* ctx) { |
||||
EVP_MD_CTX* sha = EVP_MD_CTX_new(); |
||||
uint8_t hash[MT_HASH_SIZE], children[MT_BUCKETS][MT_HASH_SIZE]; |
||||
static const uint8_t empty[sizeof(children)]; |
||||
int changed = 0; |
||||
if (!db || !ns || !leaf_hash || !sha) goto fail; |
||||
for (int level = MT_MAX_LEVEL; level >= 0; level--) { |
||||
uint64_t prefix = merkle_sync_level_prefix(key, (uint8_t)level); |
||||
uint8_t domain[2] = { 0x4d, (uint8_t)level }; |
||||
if (EVP_DigestInit_ex(sha, EVP_sha256(), NULL) != 1 || EVP_DigestUpdate(sha, domain, sizeof(domain)) != 1) goto fail; |
||||
int count; |
||||
if (level == MT_MAX_LEVEL) { |
||||
count = leaf_hash(ctx, ns, (uint8_t)level, prefix, sha); |
||||
if (count < 0) goto fail; |
||||
} else { |
||||
if (merkle_tree_children(db, ns, (uint8_t)level, prefix, children) < 0) goto fail; |
||||
count = memcmp(children, empty, sizeof(children)) != 0; |
||||
if (EVP_DigestUpdate(sha, children, sizeof(children)) != 1) goto fail; |
||||
} |
||||
if (EVP_DigestFinal_ex(sha, hash, NULL) != 1) goto fail; |
||||
if (!count) memset(hash, 0, sizeof(hash)); |
||||
int rc = tree_store(db, ns, (uint8_t)level, prefix, hash); |
||||
if (rc < 0) goto fail; |
||||
if (rc == 0) break; |
||||
changed = 1; |
||||
} |
||||
EVP_MD_CTX_free(sha); |
||||
return changed; |
||||
fail: |
||||
EVP_MD_CTX_free(sha); |
||||
return tree_error(db, "recompute", ns); |
||||
} |
||||
@ -0,0 +1,24 @@
|
||||
#ifndef MERKLE_TREE_H |
||||
#define MERKLE_TREE_H |
||||
|
||||
#include <stdint.h> |
||||
#include <openssl/evp.h> |
||||
#include <sqlite3.h> |
||||
|
||||
#define MT_MAX_LEVEL 5 |
||||
#define MT_BUCKETS 32 |
||||
#define MT_HASH_SIZE 32 |
||||
|
||||
typedef int (*merkle_leaf_hash_fn)(void* ctx, const char* ns, uint8_t level, uint64_t prefix, EVP_MD_CTX* hash); |
||||
|
||||
/* Дерево — производный индекс. При открытии очищается и восстанавливается из записей модели. */ |
||||
int merkle_tree_init(sqlite3* db); |
||||
int merkle_tree_get(sqlite3* db, const char* ns, uint8_t level, uint64_t prefix, uint8_t hash[MT_HASH_SIZE]); |
||||
int merkle_tree_children(sqlite3* db, const char* ns, uint8_t level, uint64_t prefix, |
||||
uint8_t hashes[MT_BUCKETS][MT_HASH_SIZE]); |
||||
/* Работает внутри транзакции вызывающей стороны. Сначала лист, затем родители до корня. */ |
||||
int merkle_tree_update(sqlite3* db, const char* ns, uint64_t key, merkle_leaf_hash_fn leaf_hash, void* ctx); |
||||
uint64_t merkle_sync_level_prefix(uint64_t key, uint8_t level); |
||||
uint8_t merkle_sync_prefix_bytes(uint8_t level); |
||||
|
||||
#endif |
||||
@ -0,0 +1,68 @@
|
||||
/* Реальный адаптер, SQLite и транзакции; без копии алгоритма и без сетевых таймингов. */ |
||||
#include "../src/chat/member_sync.c" |
||||
#include <assert.h> |
||||
|
||||
static void setup(struct UTUN_INSTANCE* inst, struct UASYNC* ua) { |
||||
memset(inst, 0, sizeof(*inst)); inst->ua = ua; inst->node_id = 1; |
||||
assert(sqlite3_open(":memory:", &inst->topo_sqlite_db) == SQLITE_OK); |
||||
assert(topo_node_sqlite_init(inst->topo_sqlite_db) == 0); |
||||
uint8_t key[32] = {1}, signature[64] = {0}; |
||||
assert(topo_node_sqlite_channel_put(inst->topo_sqlite_db, "1001", "test", 1, key, NULL, key, NULL, signature) == 0); |
||||
assert(member_sync_init(inst) == 0); |
||||
} |
||||
|
||||
static void put(struct UTUN_INSTANCE* inst, uint64_t id) { |
||||
uint8_t key[32] = {1}; |
||||
assert(member_sync_put(inst, "1001", id, key, key, NULL, 0, NULL, 0, "{}", NULL, NULL, 0, 0, NULL) >= 0); |
||||
} |
||||
|
||||
static void root(struct UTUN_INSTANCE* inst, uint8_t hash[32]) { |
||||
assert(merkle_sync_read_hash(inst, "1001", 0, 0, hash) == 0); |
||||
} |
||||
|
||||
int main(void) { |
||||
debug_config_init(); debug_set_level(DEBUG_LEVEL_ERROR); |
||||
size_t baseline = u_get_allocated_count(); |
||||
struct UASYNC* ua = uasync_create(); assert(ua); |
||||
struct UTUN_INSTANCE source, target; |
||||
setup(&source, ua); setup(&target, ua); |
||||
put(&source, 1); put(&source, 2); |
||||
|
||||
uint8_t raw[2049], *page = raw + 1; /* проверяем невыровненный буфер */ |
||||
size_t len = 300; uint64_t next; int more; |
||||
assert(_member_get_page(&source, "1001", 0, 0, 0, page, &len, &next, &more) == 0); |
||||
assert(page[0] == 0 && page[1] == 1 && next == 1 && more == 1 && len <= 300); |
||||
struct member_wire w; size_t used; |
||||
assert(member_decode(page + 2, len - 2, &w, &used) == 0 && w.rec.node_id == 1 && used == len - 2); |
||||
len = 300; |
||||
assert(_member_get_page(&source, "1001", 0, 1, next, page, &len, &next, &more) == 0 && next == 2 && !more); |
||||
puts("PASS: bounded member pages, canonical integers and unaligned buffers"); |
||||
|
||||
len = 2048; |
||||
assert(_member_get_page(&source, "1001", 0, 0, 0, page, &len, &next, &more) == 0 && page[1] == 2 && !more); |
||||
assert(_member_apply_items(&target, "1001", 1, page, len - 1) < 0 && member_sync_count(&target, "1001") == 0); |
||||
sqlite3* db = target.topo_sqlite_db; |
||||
assert(sqlite3_exec(db, "CREATE TRIGGER fail_second BEFORE INSERT ON peers_1001 WHEN NEW.node_id=2 " |
||||
"BEGIN SELECT RAISE(ABORT,'injected second record failure'); END", NULL, NULL, NULL) == SQLITE_OK); |
||||
assert(_member_apply_items(&target, "1001", 1, page, len) < 0); |
||||
uint8_t hash[32], empty[32] = {0}; root(&target, hash); |
||||
assert(member_sync_count(&target, "1001") == 0 && !memcmp(hash, empty, 32) && sqlite3_get_autocommit(db)); |
||||
assert(sqlite3_exec(db, "DROP TRIGGER fail_second", NULL, NULL, NULL) == SQLITE_OK); |
||||
assert(_member_apply_items(&target, "1001", 1, page, len) == 0 && member_sync_count(&target, "1001") == 2); |
||||
uint8_t expected[32]; root(&source, expected); root(&target, hash); assert(!memcmp(hash, expected, 32)); |
||||
puts("PASS: malformed page rejected; second-record error rolls back data and tree; complete page converges"); |
||||
|
||||
assert(sqlite3_exec(db, "UPDATE peers_1001 SET x25519_pubkey=X'01' WHERE node_id=1", NULL, NULL, NULL) == SQLITE_OK); |
||||
assert(merkle_sync_recompute_path(&target, "1001", 1) < 0); |
||||
root(&target, hash); assert(!memcmp(hash, expected, 32)); |
||||
len = 2048; |
||||
assert(_member_get_page(&target, "1001", 0, 0, 0, page, &len, &next, &more) < 0); |
||||
puts("PASS: corrupt BLOB cannot be hashed or serialized"); |
||||
|
||||
member_sync_destroy(&target); member_sync_destroy(&source); |
||||
assert(sqlite3_close(target.topo_sqlite_db) == SQLITE_OK && sqlite3_close(source.topo_sqlite_db) == SQLITE_OK); |
||||
uasync_poll(ua, 0); uasync_destroy(ua, 0); |
||||
assert(u_get_allocated_count() == baseline); |
||||
puts("ALL PASS: member adapter"); |
||||
return 0; |
||||
} |
||||
@ -0,0 +1,215 @@
|
||||
/* Проверяем границы автомата на точных wire-пакетах. Включение реализации позволяет
|
||||
* задать состояние без копирования private-структур и без подмены очередей/таймеров. */ |
||||
#include "../src/chat/merkle_sync.c" |
||||
#include <assert.h> |
||||
|
||||
struct fixture { |
||||
struct UTUN_INSTANCE inst; |
||||
struct ETCP_CONN conn; |
||||
int callbacks, result, apply_error; |
||||
}; |
||||
|
||||
static int t_hash(void* arg, const char* ns, uint8_t level, uint64_t prefix, EVP_MD_CTX* hash) { |
||||
(void)arg; (void)ns; (void)level; (void)prefix; |
||||
return EVP_DigestUpdate(hash, "record", 6) == 1 ? 1 : -1; |
||||
} |
||||
|
||||
static int t_page(void* arg, const char* ns, uint64_t prefix, int has_after, uint64_t after, |
||||
uint8_t* data, size_t* len, uint64_t* next, int* more) { |
||||
(void)arg; (void)ns; (void)prefix; (void)has_after; (void)after; |
||||
memset(data, 0, *len); *next = 1; *more = 1; return 0; |
||||
} |
||||
|
||||
static int t_apply(void* arg, const char* ns, uint64_t peer, const uint8_t* data, size_t len) { |
||||
(void)ns; (void)peer; (void)data; (void)len; |
||||
return ((struct fixture*)arg)->apply_error; |
||||
} |
||||
|
||||
static const struct merkle_sync_data_ops t_ops = { |
||||
.update_bucket_hash = t_hash, .get_page = t_page, .apply_items = t_apply |
||||
}; |
||||
|
||||
static void t_done(uint64_t peer, const char* ns, int result, void* arg) { |
||||
(void)peer; (void)ns; |
||||
struct fixture* f = arg; |
||||
f->callbacks++; f->result = result; |
||||
} |
||||
|
||||
static void t_init(struct fixture* f) { |
||||
memset(f, 0, sizeof(*f)); |
||||
f->inst.node_id = 1; |
||||
f->inst.ua = uasync_create(); assert(f->inst.ua); |
||||
assert(sqlite3_open(":memory:", &f->inst.topo_sqlite_db) == SQLITE_OK); |
||||
f->inst.connections = queue_new(f->inst.ua, 16, 0, 8, "test connections"); assert(f->inst.connections); |
||||
f->conn.instance = &f->inst; f->conn.peer_node_id = 2; f->conn.initialized = 1; f->conn.links_up = 1; |
||||
f->conn.send_input_q = queue_new(f->inst.ua, 0, 0, 0, "test send"); assert(f->conn.send_input_q); |
||||
queue_set_waiter_defer(f->conn.send_input_q, 1); |
||||
struct ll_entry* e = queue_entry_new(sizeof(struct conn_queue_entry)); assert(e); |
||||
struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data; |
||||
ce->peer_node_id = 2; ce->conn = &f->conn; |
||||
assert(queue_data_put_with_index(f->inst.connections, e) == 0); |
||||
assert(merkle_sync_init(&f->inst, 0x72, &t_ops, f) == 0); |
||||
} |
||||
|
||||
static void t_clear_queue(struct ll_queue* q) { |
||||
struct ll_entry* e; |
||||
while ((e = queue_data_get(q))) { queue_dgram_free(e); queue_entry_free(e); } |
||||
queue_free(q); |
||||
} |
||||
|
||||
static void t_destroy(struct fixture* f) { |
||||
merkle_sync_destroy(&f->inst); |
||||
assert(f->conn.ref_count == 0); |
||||
t_clear_queue(f->conn.send_input_q); t_clear_queue(f->inst.connections); |
||||
sqlite3_close(f->inst.topo_sqlite_db); |
||||
uasync_poll(f->inst.ua, 0); uasync_destroy(f->inst.ua, 0); |
||||
} |
||||
|
||||
static struct ms_session* t_session(struct fixture* f) { |
||||
struct ms_session* s = ms_create(f->inst.msync, &f->conn, 1001); assert(s); |
||||
s->round = 7; s->state = MS_PULL; s->request_id = 1; s->waiting = 1; |
||||
s->requests = u_calloc(1, sizeof(*s->requests)); assert(s->requests); |
||||
s->requests->cb = t_done; s->requests->arg = f; |
||||
return s; |
||||
} |
||||
|
||||
static void t_receive(struct fixture* f, uint8_t type, uint64_t round, uint32_t request, const uint8_t* data, size_t len) { |
||||
struct ll_entry* e = ll_alloc_lldgram((uint16_t)(MS_HEADER + len)); assert(e); |
||||
e->len = (uint16_t)(MS_HEADER + len); e->dgram[0] = 0x72; |
||||
ms_write64(e->dgram + 1, 1001); e->dgram[9] = type; |
||||
ms_write64(e->dgram + 10, round); ms_write32(e->dgram + 18, request); |
||||
if (len) memcpy(e->dgram + MS_HEADER, data, len); |
||||
ms_receive(&f->conn, e); |
||||
} |
||||
|
||||
static void test_malformed_and_old_round(void) { |
||||
struct fixture f; t_init(&f); |
||||
struct ms_session* s = t_session(&f); |
||||
uint8_t short_hashes[4] = {1}; |
||||
t_receive(&f, MS_HASHES, 6, 1, short_hashes, sizeof(short_hashes)); |
||||
assert(!f.callbacks && s->waiting && s->round == 7); |
||||
t_receive(&f, MS_HASHES, 7, 1, short_hashes, sizeof(short_hashes)); |
||||
assert(f.callbacks == 1 && f.result == MT_ERR_PROTOCOL && s->state == MS_FAILED); |
||||
t_destroy(&f); |
||||
puts("PASS: stale round ignored; truncated hashes rejected before reading"); |
||||
} |
||||
|
||||
static void test_rejected_data(void) { |
||||
struct fixture f; t_init(&f); |
||||
struct ms_session* s = t_session(&f); s->stack[0].level = MT_MAX_LEVEL; |
||||
f.apply_error = MT_ERR_DATA; |
||||
uint8_t page[11] = {0}; |
||||
t_receive(&f, MS_PAGE, 7, 1, page, sizeof(page)); |
||||
assert(f.callbacks == 1 && f.result == MT_ERR_DATA && s->state == MS_FAILED); |
||||
t_destroy(&f); |
||||
puts("PASS: failed apply cannot produce success"); |
||||
} |
||||
|
||||
static void test_sibling_walk(void) { |
||||
struct fixture f; t_init(&f); |
||||
struct ms_session* s = t_session(&f); |
||||
uint8_t hashes[32 * 32] = {0}; hashes[0] = 1; hashes[32] = 1; |
||||
assert(ms_receive_hashes(s, 1, hashes, sizeof(hashes)) == 0); |
||||
assert(s->depth == 1 && s->stack[1].prefix == 0 && s->stack[0].children == 2); |
||||
ms_discard_tx(s); |
||||
memset(hashes, 0, sizeof(hashes)); |
||||
assert(ms_receive_hashes(s, 2, hashes, sizeof(hashes)) == 0); |
||||
assert(s->depth == 1 && s->stack[1].prefix == (UINT64_C(1) << 59)); |
||||
assert(s->waiting && !f.callbacks); |
||||
t_destroy(&f); |
||||
puts("PASS: completing one branch preserves the sibling request"); |
||||
} |
||||
|
||||
static void test_cancel_waiters(void) { |
||||
for (int blocked = 0; blocked < 2; blocked++) { |
||||
struct fixture f; t_init(&f); |
||||
struct ms_session* s = t_session(&f); |
||||
if (blocked) { struct ll_entry* e = queue_entry_new(0); assert(e); assert(queue_data_put(f.conn.send_input_q, e) == 0); } |
||||
assert(ms_send(s, MS_QUERY, 1, NULL, 0) == 0); |
||||
assert(s->waiter.internal || s->waiter.call_soon_id); |
||||
f.conn.links_up = 0; |
||||
merkle_sync_cancel_peer(&f.inst, 2); |
||||
assert(!f.inst.msync->sessions && !f.conn.send_input_q->waiter_head); |
||||
uasync_poll(f.inst.ua, 0); |
||||
assert(!f.callbacks); |
||||
t_destroy(&f); |
||||
} |
||||
puts("PASS: cancel after links_down removes queued and deferred waiters"); |
||||
} |
||||
|
||||
static void test_page_size_and_cursor(void) { |
||||
struct fixture f; t_init(&f); |
||||
struct ms_session* s = t_session(&f); s->state = MS_SERVE; |
||||
uint8_t query[MS_QUERY_SIZE] = { MT_MAX_LEVEL }; |
||||
assert(ms_serve_query(s, 1, query, sizeof(query)) == 0); |
||||
assert(s->tx && s->tx->len == MT_PAGE_SIZE); |
||||
ms_discard_tx(s); |
||||
query[9] = 1; ms_write64(query + 10, 1); |
||||
assert(ms_serve_query(s, 2, query, sizeof(query)) == MT_ERR_DATA); |
||||
query[0] = 255; |
||||
assert(ms_serve_query(s, 3, query, sizeof(query)) == MT_ERR_PROTOCOL); |
||||
t_destroy(&f); |
||||
puts("PASS: page bounded; non-advancing cursor and invalid level rejected"); |
||||
} |
||||
|
||||
static void test_readonly_tree(void) { |
||||
struct fixture f; t_init(&f); |
||||
assert(sqlite3_exec(f.inst.topo_sqlite_db, "PRAGMA query_only=ON", NULL, NULL, NULL) == SQLITE_OK); |
||||
assert(merkle_sync_recompute_path(&f.inst, "1001", 0) < 0); |
||||
uint8_t hash[32], empty[32] = {0}; |
||||
assert(merkle_sync_read_hash(&f.inst, "1001", 0, 0, hash) == 0 && !memcmp(hash, empty, 32)); |
||||
t_destroy(&f); |
||||
puts("PASS: SQLite failure is reported and cannot publish a root"); |
||||
} |
||||
|
||||
static void test_backpressure_and_disconnect(void) { |
||||
struct fixture f; t_init(&f); |
||||
struct ll_entry* blocker = queue_entry_new(0); assert(blocker); |
||||
assert(queue_data_put(f.conn.send_input_q, blocker) == 0); |
||||
assert(merkle_sync_start(&f.inst, 2, "1001", t_done, &f) == 0); |
||||
uasync_poll(f.inst.ua, 0); |
||||
struct ms_session* s = f.inst.msync->sessions; |
||||
assert(s->tx && s->waiter.internal); |
||||
for (int i = 0; i < 1000; i++) merkle_sync_changed(&f.inst, "1001"); |
||||
assert(s->tx && queue_entry_count(f.conn.send_input_q) == 1 && s->dirty); |
||||
f.conn.links_up = 0; |
||||
ms_conn_status(&f.conn, ETCP_CONN_STATUS_DOWN, f.inst.msync); |
||||
assert(!f.inst.msync->sessions && f.callbacks == 1 && f.result == MT_ERR_DISCONNECTED); |
||||
t_destroy(&f); |
||||
puts("PASS: 1000 changes coalesce behind backpressure; disconnect completes request with error"); |
||||
} |
||||
|
||||
static void test_peer_restart(void) { |
||||
struct fixture f; t_init(&f); |
||||
struct ms_session* s = t_session(&f); |
||||
assert(ms_send(s, MS_QUERY, 1, NULL, 0) == 0); |
||||
t_receive(&f, MS_WAKE, 0, 0, NULL, 0); |
||||
assert(s->round > 7 && s->state == MS_BEGIN_WAIT && s->tx && s->tx->dgram[9] == MS_BEGIN); |
||||
uint8_t hashes[32 * 32] = {0}; |
||||
t_receive(&f, MS_HASHES, 7, 1, hashes, sizeof(hashes)); |
||||
assert(!f.callbacks && s->state == MS_BEGIN_WAIT); |
||||
t_destroy(&f); |
||||
puts("PASS: reattaching follower restarts the round; old responses cannot affect it"); |
||||
} |
||||
|
||||
static void test_send_rejection(void) { |
||||
struct fixture f; t_init(&f); |
||||
struct ms_session* s = t_session(&f); |
||||
queue_set_size_limit(f.conn.send_input_q, 0); |
||||
assert(ms_send(s, MS_QUERY, 1, NULL, 0) == 0); |
||||
uasync_poll(f.inst.ua, 0); |
||||
assert(f.callbacks == 1 && f.result == MT_ERR_IO && !s->tx); |
||||
t_destroy(&f); |
||||
puts("PASS: transport rejection retains ownership; packet is freed exactly once"); |
||||
} |
||||
|
||||
int main(void) { |
||||
debug_config_init(); debug_set_level(DEBUG_LEVEL_ERROR); |
||||
size_t baseline = u_get_allocated_count(); |
||||
test_malformed_and_old_round(); test_rejected_data(); test_sibling_walk(); test_cancel_waiters(); |
||||
test_page_size_and_cursor(); test_readonly_tree(); test_backpressure_and_disconnect(); test_peer_restart(); |
||||
test_send_rejection(); |
||||
assert(u_get_allocated_count() == baseline); |
||||
puts("ALL PASS: protocol regressions, no tracked allocations leaked"); |
||||
return 0; |
||||
} |
||||
Loading…
Reference in new issue