Browse Source
- merkle_sync.h/c: tree (5-level prefix, SHA256 buckets), wire protocol (HASHES/REQUEST/BATCH), sessions (10s timeout, 3 retries), async API with done_cb - member_sync: thin glue layer, implements 3 data model callbacks (update_bucket_hash, get_items, apply_items) for chat members - chat_sync: updated member_sync_start() calls with done_cb callback - member_tree_hash table renamed to merkle_tree_hash - member_sync.h public API preserved, internal tree/wire logic removedtopo_upd
7 changed files with 1155 additions and 786 deletions
File diff suppressed because it is too large
Load Diff
@ -0,0 +1,618 @@ |
|||||||
|
#include "merkle_sync.h" |
||||||
|
|
||||||
|
#include "../../../src/utun_instance.h" |
||||||
|
#include "../../../src/etcp_router.h" |
||||||
|
#include "../../../src/etcp_api.h" |
||||||
|
#include "../../../src/etcp.h" |
||||||
|
#include "../../../src/topo_group.h" |
||||||
|
#include "../../../lib/debug_config.h" |
||||||
|
#include "../../../lib/mem.h" |
||||||
|
#include "../../../lib/u_async.h" |
||||||
|
|
||||||
|
#include <string.h> |
||||||
|
#include <stdlib.h> |
||||||
|
#include <sqlite3.h> |
||||||
|
|
||||||
|
#define MS_ID "merkle_sync" |
||||||
|
#define MS_SYNC_TIMEOUT_MS 10000 |
||||||
|
#define MS_BG_INTERVAL_MS 100 |
||||||
|
|
||||||
|
static struct merkle_sync* g_merkle = NULL; |
||||||
|
|
||||||
|
struct bucket_entry { uint8_t level; uint8_t prefix_bytes; uint64_t prefix; }; |
||||||
|
|
||||||
|
struct ms_session { |
||||||
|
struct ms_session* next; |
||||||
|
char ns[64]; |
||||||
|
uint64_t peer; |
||||||
|
uint8_t active; |
||||||
|
void* timer; |
||||||
|
uint8_t retries; |
||||||
|
merkle_sync_done_cb done_cb; |
||||||
|
void* cb_arg; |
||||||
|
}; |
||||||
|
|
||||||
|
struct merkle_sync { |
||||||
|
struct UTUN_INSTANCE* inst; |
||||||
|
uint8_t svc_id; |
||||||
|
const struct merkle_sync_data_ops* ops; |
||||||
|
void* data_ctx; |
||||||
|
struct ms_session* sessions; |
||||||
|
uint8_t initialized; |
||||||
|
void* bg_timer; |
||||||
|
}; |
||||||
|
|
||||||
|
/* ── DB access ── */ |
||||||
|
|
||||||
|
static sqlite3* _db(struct UTUN_INSTANCE* inst) { |
||||||
|
return inst && inst->topo_groups ? inst->topo_groups->topo_sqlite_db : NULL; |
||||||
|
} |
||||||
|
|
||||||
|
/* ── Prefix arithmetic ── */ |
||||||
|
|
||||||
|
uint64_t merkle_sync_level_prefix(uint64_t key, uint8_t level) { |
||||||
|
int shift = 63 - (int)level * 5; if (shift < 0) shift = 0; |
||||||
|
return (key >> shift) << shift; |
||||||
|
} |
||||||
|
|
||||||
|
uint8_t merkle_sync_prefix_bytes(uint8_t level) { int bits = level * 5; return (uint8_t)((bits + 7) / 8); } |
||||||
|
|
||||||
|
static void _prefix_write(uint8_t* out, uint64_t prefix, uint8_t pb) { |
||||||
|
for (int i = (int)pb - 1; i >= 0; i--) out[pb - 1 - i] = (uint8_t)(prefix >> (i * 8)); |
||||||
|
} |
||||||
|
|
||||||
|
static uint64_t _prefix_read(const uint8_t* in, uint8_t pb) { |
||||||
|
uint64_t v = 0; |
||||||
|
for (uint8_t i = 0; i < pb && i < 8; i++) v = (v << 8) | in[i]; |
||||||
|
return v; |
||||||
|
} |
||||||
|
|
||||||
|
static int _popcount_u32(uint32_t v) { return __builtin_popcount(v); } |
||||||
|
|
||||||
|
/* ── Ensure merkle_tree_hash table ── */ |
||||||
|
|
||||||
|
static void _ensure_table(struct merkle_sync* ms) { |
||||||
|
sqlite3* db = _db(ms->inst); if (!db) return; |
||||||
|
sqlite3_exec(db, |
||||||
|
"CREATE TABLE IF NOT EXISTS merkle_tree_hash (" |
||||||
|
" namespace TEXT NOT NULL," |
||||||
|
" level INTEGER NOT NULL CHECK(level BETWEEN 1 AND 5)," |
||||||
|
" prefix64 INTEGER NOT NULL," |
||||||
|
" hash BLOB NOT NULL," |
||||||
|
" member_count INTEGER NOT NULL," |
||||||
|
" PRIMARY KEY (namespace, level, prefix64))", |
||||||
|
NULL, NULL, NULL); |
||||||
|
} |
||||||
|
|
||||||
|
/* ── Tree operations ── */ |
||||||
|
|
||||||
|
static int _recompute_bucket(struct merkle_sync* ms, const char* ns, |
||||||
|
uint8_t level, uint64_t prefix64) { |
||||||
|
sqlite3* db = _db(ms->inst); if (!db) return -1; |
||||||
|
|
||||||
|
EVP_MD_CTX* ctx = EVP_MD_CTX_new(); |
||||||
|
EVP_DigestInit_ex(ctx, EVP_sha256(), NULL); |
||||||
|
int count = ms->ops->update_bucket_hash(ms->data_ctx, ns, level, prefix64, ctx); |
||||||
|
|
||||||
|
if (count < 0) { EVP_MD_CTX_free(ctx); return -1; } |
||||||
|
if (count == 0) { |
||||||
|
EVP_MD_CTX_free(ctx); |
||||||
|
sqlite3_stmt* ds = NULL; |
||||||
|
sqlite3_prepare_v2(db, |
||||||
|
"DELETE FROM merkle_tree_hash WHERE namespace=? AND level=? AND prefix64=?", |
||||||
|
-1, &ds, NULL); |
||||||
|
if (ds) { sqlite3_bind_text(ds, 1, ns, -1, SQLITE_STATIC); sqlite3_bind_int(ds, 2, level); |
||||||
|
sqlite3_bind_int64(ds, 3, (sqlite3_int64)prefix64); sqlite3_step(ds); sqlite3_finalize(ds); } |
||||||
|
return 0; |
||||||
|
} |
||||||
|
|
||||||
|
uint8_t hash[MT_HASH_SIZE]; |
||||||
|
EVP_DigestFinal_ex(ctx, hash, NULL); |
||||||
|
EVP_MD_CTX_free(ctx); |
||||||
|
|
||||||
|
sqlite3_stmt* is = NULL; |
||||||
|
sqlite3_prepare_v2(db, |
||||||
|
"INSERT OR REPLACE INTO merkle_tree_hash(namespace, level, prefix64, hash, member_count)" |
||||||
|
" VALUES(?,?,?,?,?)", -1, &is, NULL); |
||||||
|
if (is) { |
||||||
|
sqlite3_bind_text(is, 1, ns, -1, SQLITE_STATIC); |
||||||
|
sqlite3_bind_int(is, 2, level); |
||||||
|
sqlite3_bind_int64(is, 3, (sqlite3_int64)prefix64); |
||||||
|
sqlite3_bind_blob(is, 4, hash, MT_HASH_SIZE, SQLITE_STATIC); |
||||||
|
sqlite3_bind_int(is, 5, count); |
||||||
|
sqlite3_step(is); sqlite3_finalize(is); |
||||||
|
} |
||||||
|
return 0; |
||||||
|
} |
||||||
|
|
||||||
|
void merkle_sync_recompute_path(struct UTUN_INSTANCE* inst, const char* ns, uint64_t key) { |
||||||
|
struct merkle_sync* ms = g_merkle; |
||||||
|
if (!ms || !ms->initialized) return; |
||||||
|
for (uint8_t level = 1; level <= MT_MAX_LEVEL; level++) |
||||||
|
_recompute_bucket(ms, ns, level, merkle_sync_level_prefix(key, level)); |
||||||
|
} |
||||||
|
|
||||||
|
const uint8_t* merkle_sync_get_hash(struct UTUN_INSTANCE* inst, const char* ns, |
||||||
|
uint8_t level, uint64_t prefix64) { |
||||||
|
static uint8_t zero[MT_HASH_SIZE]; |
||||||
|
sqlite3* db = _db(inst); if (!db) { memset(zero, 0, MT_HASH_SIZE); return zero; } |
||||||
|
sqlite3_stmt* stmt = NULL; |
||||||
|
sqlite3_prepare_v2(db, |
||||||
|
"SELECT hash FROM merkle_tree_hash WHERE namespace=? AND level=? AND prefix64=?", |
||||||
|
-1, &stmt, NULL); |
||||||
|
if (!stmt) { memset(zero, 0, MT_HASH_SIZE); return zero; } |
||||||
|
sqlite3_bind_text(stmt, 1, ns, -1, SQLITE_STATIC); |
||||||
|
sqlite3_bind_int(stmt, 2, level); |
||||||
|
sqlite3_bind_int64(stmt, 3, (sqlite3_int64)prefix64); |
||||||
|
memset(zero, 0, MT_HASH_SIZE); |
||||||
|
if (sqlite3_step(stmt) == SQLITE_ROW) { |
||||||
|
const void* h = sqlite3_column_blob(stmt, 0); |
||||||
|
if (h) memcpy(zero, h, MT_HASH_SIZE); |
||||||
|
} |
||||||
|
sqlite3_finalize(stmt); |
||||||
|
return zero; |
||||||
|
} |
||||||
|
|
||||||
|
static int _get_level_hashes(struct merkle_sync* ms, const char* ns, |
||||||
|
uint8_t level, uint64_t prefix, uint8_t prefix_bytes, |
||||||
|
uint32_t* bitmap, uint8_t hashes[MT_BUCKETS][MT_HASH_SIZE]) { |
||||||
|
(void)prefix_bytes; |
||||||
|
memset(hashes, 0, sizeof(uint8_t) * MT_BUCKETS * MT_HASH_SIZE); |
||||||
|
*bitmap = 0; |
||||||
|
if (level >= MT_MAX_LEVEL) return 0; |
||||||
|
int next_shift = 63 - ((int)level + 1) * 5; |
||||||
|
|
||||||
|
for (int i = 0; i < MT_BUCKETS; i++) { |
||||||
|
uint64_t child_prefix = prefix | ((uint64_t)i << (next_shift > 0 ? next_shift : 0)); |
||||||
|
const uint8_t* h = merkle_sync_get_hash(ms->inst, ns, (uint8_t)(level + 1), child_prefix); |
||||||
|
int empty = 1; |
||||||
|
for (int j = 0; j < MT_HASH_SIZE; j++) { if (h[j] != 0) { empty = 0; break; } } |
||||||
|
if (!empty) { *bitmap |= (1u << i); memcpy(hashes[i], h, MT_HASH_SIZE); } |
||||||
|
} |
||||||
|
return 0; |
||||||
|
} |
||||||
|
|
||||||
|
/* ── Send helpers ── */ |
||||||
|
|
||||||
|
static int _send_msg(struct merkle_sync* ms, uint64_t peer, |
||||||
|
const uint8_t* payload, size_t len) { |
||||||
|
if (!ms->inst || len < 1) return -1; |
||||||
|
uint8_t* buf = u_malloc(1 + len); |
||||||
|
if (!buf) return -1; |
||||||
|
buf[0] = ms->svc_id; |
||||||
|
memcpy(buf + 1, payload, len); |
||||||
|
struct ll_entry* entry = queue_entry_new(0); |
||||||
|
if (!entry) { u_free(buf); return -1; } |
||||||
|
entry->dgram = buf; entry->len = 1 + len; |
||||||
|
int r = etcp_route_send(ms->inst, peer, entry, 0); |
||||||
|
if (r != 0) { u_free(buf); queue_entry_free(entry); } |
||||||
|
return r; |
||||||
|
} |
||||||
|
|
||||||
|
static int _send_hashes(struct merkle_sync* ms, uint64_t peer, const char* ns, |
||||||
|
uint8_t level, uint64_t prefix, uint8_t prefix_bytes, int is_data) { |
||||||
|
uint8_t ch_len = (uint8_t)strlen(ns); |
||||||
|
size_t max_sz = 1 + 1 + ch_len + 1 + 1 + prefix_bytes + 1 + 4 + MT_BUCKETS * MT_HASH_SIZE + 2 + 65536; |
||||||
|
uint8_t* buf = u_malloc(max_sz); |
||||||
|
if (!buf) return -1; |
||||||
|
uint8_t* p = buf; |
||||||
|
*p++ = 0x01; /* MS_MSG_HASHES */ |
||||||
|
*p++ = ch_len; memcpy(p, ns, ch_len); p += ch_len; |
||||||
|
*p++ = level; |
||||||
|
*p++ = prefix_bytes; _prefix_write(p, prefix, prefix_bytes); p += prefix_bytes; |
||||||
|
*p++ = (uint8_t)(is_data ? 1 : 0); |
||||||
|
|
||||||
|
if (is_data) { |
||||||
|
size_t mlen = 65536; uint8_t* mbuf = u_malloc(mlen); |
||||||
|
if (mbuf) { |
||||||
|
if (ms->ops->get_items(ms->data_ctx, ns, level, prefix, prefix_bytes, mbuf, &mlen) == 0) |
||||||
|
{ memcpy(p, mbuf, mlen); p += mlen; } |
||||||
|
else { uint16_t zero = 0; memcpy(p, &zero, 2); p += 2; } |
||||||
|
u_free(mbuf); |
||||||
|
} |
||||||
|
} else { |
||||||
|
uint32_t bitmap; uint8_t hashes[MT_BUCKETS][MT_HASH_SIZE]; |
||||||
|
_get_level_hashes(ms, ns, level, prefix, prefix_bytes, &bitmap, hashes); |
||||||
|
memcpy(p, &bitmap, 4); p += 4; |
||||||
|
for (int i = 0; i < MT_BUCKETS; i++) |
||||||
|
if (bitmap & (1u << i)) { memcpy(p, hashes[i], MT_HASH_SIZE); p += MT_HASH_SIZE; } |
||||||
|
} |
||||||
|
int r = _send_msg(ms, peer, buf, (size_t)(p - buf)); |
||||||
|
u_free(buf); |
||||||
|
return r; |
||||||
|
} |
||||||
|
|
||||||
|
static int _send_batch(struct merkle_sync* ms, uint64_t peer, const char* ns, |
||||||
|
struct bucket_entry* buckets, int count) { |
||||||
|
uint8_t ch_len = (uint8_t)strlen(ns); |
||||||
|
size_t max_sz = 1 + 1 + ch_len + 1 + (size_t)count * (1 + 1 + 8 + 1 + 4 + MT_BUCKETS * MT_HASH_SIZE + 65536); |
||||||
|
uint8_t* buf = u_malloc(max_sz); |
||||||
|
if (!buf) return -1; |
||||||
|
uint8_t* p = buf; |
||||||
|
*p++ = 0x03; /* MS_MSG_BATCH */ |
||||||
|
*p++ = ch_len; memcpy(p, ns, ch_len); p += ch_len; |
||||||
|
*p++ = (uint8_t)count; |
||||||
|
|
||||||
|
for (int i = 0; i < count; i++) { |
||||||
|
*p++ = buckets[i].level; |
||||||
|
*p++ = buckets[i].prefix_bytes; |
||||||
|
_prefix_write(p, buckets[i].prefix, buckets[i].prefix_bytes); p += buckets[i].prefix_bytes; |
||||||
|
|
||||||
|
EVP_MD_CTX* tctx = EVP_MD_CTX_new(); |
||||||
|
EVP_DigestInit_ex(tctx, EVP_sha256(), NULL); |
||||||
|
int bcount = ms->ops->update_bucket_hash(ms->data_ctx, ns, |
||||||
|
buckets[i].level, buckets[i].prefix, tctx); |
||||||
|
if (bcount < 0) bcount = 0; |
||||||
|
EVP_MD_CTX_free(tctx); |
||||||
|
|
||||||
|
int is_terminal = (buckets[i].level >= MT_MAX_LEVEL || bcount < 8); |
||||||
|
*p++ = (uint8_t)(is_terminal ? 1 : 0); |
||||||
|
|
||||||
|
if (is_terminal) { |
||||||
|
size_t mlen = 65536; uint8_t* mbuf = u_malloc(mlen); |
||||||
|
if (mbuf) { |
||||||
|
if (ms->ops->get_items(ms->data_ctx, ns, buckets[i].level, |
||||||
|
buckets[i].prefix, buckets[i].prefix_bytes, mbuf, &mlen) == 0) |
||||||
|
{ memcpy(p, mbuf, mlen); p += mlen; } |
||||||
|
else { uint16_t z = 0; memcpy(p, &z, 2); p += 2; } |
||||||
|
u_free(mbuf); |
||||||
|
} |
||||||
|
} else { |
||||||
|
uint32_t bm; uint8_t hs[MT_BUCKETS][MT_HASH_SIZE]; |
||||||
|
_get_level_hashes(ms, ns, buckets[i].level, |
||||||
|
buckets[i].prefix, buckets[i].prefix_bytes, &bm, hs); |
||||||
|
memcpy(p, &bm, 4); p += 4; |
||||||
|
for (int j = 0; j < MT_BUCKETS; j++) |
||||||
|
if (bm & (1u << j)) { memcpy(p, hs[j], MT_HASH_SIZE); p += MT_HASH_SIZE; } |
||||||
|
} |
||||||
|
} |
||||||
|
int r = _send_msg(ms, peer, buf, (size_t)(p - buf)); |
||||||
|
u_free(buf); |
||||||
|
return r; |
||||||
|
} |
||||||
|
|
||||||
|
/* ── Sessions ── */ |
||||||
|
|
||||||
|
static void _session_start_timer(struct ms_session* s); |
||||||
|
|
||||||
|
static void _session_timeout_cb(void* arg) { |
||||||
|
struct ms_session* s = (struct ms_session*)arg; |
||||||
|
if (!s || !s->active || !g_merkle || !g_merkle->initialized) return; |
||||||
|
s->timer = NULL; |
||||||
|
s->retries++; |
||||||
|
if (s->retries > 3) { |
||||||
|
DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: sync timeout peer=%016llx ns=%s", MS_ID, |
||||||
|
(unsigned long long)s->peer, s->ns); |
||||||
|
s->active = 0; |
||||||
|
if (s->done_cb) s->done_cb(s->peer, s->ns, MT_ERR_TIMEOUT, s->cb_arg); |
||||||
|
return; |
||||||
|
} |
||||||
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: retry %d peer=%016llx ns=%s", MS_ID, |
||||||
|
s->retries, (unsigned long long)s->peer, s->ns); |
||||||
|
_send_hashes(g_merkle, s->peer, s->ns, 1, 0, 1, 0); |
||||||
|
_session_start_timer(s); |
||||||
|
} |
||||||
|
|
||||||
|
static void _session_start_timer(struct ms_session* s) { |
||||||
|
if (!g_merkle || !g_merkle->inst) return; |
||||||
|
s->timer = uasync_set_timeout(g_merkle->inst->ua, (uint32_t)(MS_SYNC_TIMEOUT_MS * 10), |
||||||
|
s, _session_timeout_cb, "ms_sync"); |
||||||
|
} |
||||||
|
|
||||||
|
static struct ms_session* _session_find(struct merkle_sync* ms, uint64_t peer, const char* ns) { |
||||||
|
for (struct ms_session* s = ms->sessions; s; s = s->next) |
||||||
|
if (s->peer == peer && strcmp(s->ns, ns) == 0) return s; |
||||||
|
return NULL; |
||||||
|
} |
||||||
|
|
||||||
|
static void _session_done(struct ms_session* s, int result) { |
||||||
|
s->active = 0; |
||||||
|
if (s->timer) { uasync_cancel_timeout(g_merkle->inst->ua, s->timer); s->timer = NULL; } |
||||||
|
if (s->done_cb) s->done_cb(s->peer, s->ns, result, s->cb_arg); |
||||||
|
} |
||||||
|
|
||||||
|
/* ── Recv handlers ── */ |
||||||
|
|
||||||
|
static void _handle_hashes(struct merkle_sync* ms, uint64_t peer, const char* ns, |
||||||
|
const uint8_t* pl, size_t plen) { |
||||||
|
if (plen < 3) return; |
||||||
|
uint8_t level = pl[0]; uint8_t pb = pl[1]; |
||||||
|
uint64_t prefix = _prefix_read(pl + 2, pb); |
||||||
|
uint8_t is_data = pl[2 + pb]; |
||||||
|
const uint8_t* payload = pl + 3 + pb; |
||||||
|
size_t paylen = plen - 3 - pb; |
||||||
|
|
||||||
|
struct ms_session* s = _session_find(ms, peer, ns); |
||||||
|
if (s && s->timer) { uasync_cancel_timeout(g_merkle->inst->ua, s->timer); s->timer = NULL; } |
||||||
|
|
||||||
|
if (is_data) { |
||||||
|
if (paylen < 2) return; |
||||||
|
ms->ops->apply_items(ms->data_ctx, ns, payload, paylen); |
||||||
|
if (s && s->active) _session_done(s, MT_OK); |
||||||
|
return; |
||||||
|
} |
||||||
|
|
||||||
|
if (!s) { |
||||||
|
s = u_calloc(1, sizeof(*s)); if (!s) return; |
||||||
|
snprintf(s->ns, sizeof(s->ns), "%s", ns); s->peer = peer; |
||||||
|
s->next = ms->sessions; ms->sessions = s; |
||||||
|
} |
||||||
|
s->active = 1; |
||||||
|
|
||||||
|
if (paylen < 4) return; |
||||||
|
uint32_t remote_bm; memcpy(&remote_bm, payload, 4); |
||||||
|
const uint8_t* rh = payload + 4; |
||||||
|
|
||||||
|
uint32_t local_bm; uint8_t lh[MT_BUCKETS][MT_HASH_SIZE]; |
||||||
|
_get_level_hashes(ms, ns, level, prefix, pb, &local_bm, lh); |
||||||
|
|
||||||
|
uint32_t differs = remote_bm ^ local_bm; |
||||||
|
for (int i = 0; i < MT_BUCKETS; i++) { |
||||||
|
if (!(local_bm & (1u << i)) && !(remote_bm & (1u << i))) continue; |
||||||
|
if ((local_bm & (1u << i)) && (remote_bm & (1u << i))) { |
||||||
|
/* find offset of remote hash for slot i in the received blob */ |
||||||
|
int rh_idx = 0; |
||||||
|
for (int j = 0; j < i; j++) if (remote_bm & (1u << j)) rh_idx++; |
||||||
|
if (memcmp(rh + (size_t)rh_idx * MT_HASH_SIZE, lh[i], MT_HASH_SIZE) == 0) |
||||||
|
differs &= ~(1u << i); |
||||||
|
} |
||||||
|
} |
||||||
|
|
||||||
|
if (differs == 0) { _session_done(s, MT_OK); return; } |
||||||
|
|
||||||
|
int next_shift = 63 - ((int)level + 1) * 5; |
||||||
|
struct bucket_entry requests[MT_BUCKETS]; int rcount = 0; |
||||||
|
for (int i = 0; i < MT_BUCKETS && rcount < 32; i++) { |
||||||
|
if (!(differs & (1u << i))) continue; |
||||||
|
uint8_t nl = (uint8_t)(level + 1); if (nl > MT_MAX_LEVEL) nl = MT_MAX_LEVEL; |
||||||
|
requests[rcount].level = nl; |
||||||
|
requests[rcount].prefix_bytes = merkle_sync_prefix_bytes(nl); |
||||||
|
requests[rcount].prefix = prefix | ((uint64_t)i << (next_shift > 0 ? next_shift : 0)); |
||||||
|
rcount++; |
||||||
|
} |
||||||
|
|
||||||
|
if (rcount > 0) { |
||||||
|
uint8_t ch_len = (uint8_t)strlen(ns); |
||||||
|
size_t rs = 1 + 1 + ch_len + 1; |
||||||
|
for (int i = 0; i < rcount; i++) rs += 1 + 1 + requests[i].prefix_bytes + 1; |
||||||
|
uint8_t* rbuf = u_malloc(rs); |
||||||
|
if (rbuf) { |
||||||
|
uint8_t* wr = rbuf; |
||||||
|
*wr++ = 0x02; /* MS_MSG_REQUEST */ *wr++ = ch_len; |
||||||
|
memcpy(wr, ns, ch_len); wr += ch_len; |
||||||
|
*wr++ = (uint8_t)rcount; |
||||||
|
for (int i = 0; i < rcount; i++) { |
||||||
|
*wr++ = requests[i].level; |
||||||
|
*wr++ = requests[i].prefix_bytes; |
||||||
|
_prefix_write(wr, requests[i].prefix, requests[i].prefix_bytes); |
||||||
|
wr += requests[i].prefix_bytes; |
||||||
|
EVP_MD_CTX* tctx = EVP_MD_CTX_new(); |
||||||
|
EVP_DigestInit_ex(tctx, EVP_sha256(), NULL); |
||||||
|
int cnt = ms->ops->update_bucket_hash(ms->data_ctx, ns, |
||||||
|
requests[i].level, requests[i].prefix, tctx); |
||||||
|
if (cnt < 0) cnt = 0; |
||||||
|
EVP_MD_CTX_free(tctx); |
||||||
|
*wr++ = (uint8_t)(cnt < 8 || requests[i].level >= MT_MAX_LEVEL ? 1 : 0); |
||||||
|
} |
||||||
|
_send_msg(ms, peer, rbuf, (size_t)(wr - rbuf)); |
||||||
|
u_free(rbuf); |
||||||
|
} |
||||||
|
if (s) _session_start_timer(s); |
||||||
|
} else { |
||||||
|
if (s) _session_done(s, MT_OK); |
||||||
|
} |
||||||
|
} |
||||||
|
|
||||||
|
static void _handle_request(struct merkle_sync* ms, uint64_t peer, const char* ns, |
||||||
|
const uint8_t* pl, size_t plen) { |
||||||
|
if (plen < 1) return; |
||||||
|
uint8_t count = pl[0]; const uint8_t* bp = pl + 1; size_t off = 0; |
||||||
|
|
||||||
|
struct bucket_entry buckets[32]; int bc = 0; |
||||||
|
for (uint8_t i = 0; i < count && bc < 32; i++) { |
||||||
|
if (off + 2 > plen - 1) break; |
||||||
|
uint8_t lvl = bp[off++]; uint8_t pb_i = bp[off++]; |
||||||
|
if (off + pb_i > plen - 1) break; |
||||||
|
uint64_t pr = _prefix_read(bp + off, pb_i); off += pb_i; |
||||||
|
buckets[bc].level = lvl; buckets[bc].prefix_bytes = pb_i; buckets[bc].prefix = pr; bc++; |
||||||
|
off++; /* skip is_terminal byte */ |
||||||
|
} |
||||||
|
if (bc > 0) _send_batch(ms, peer, ns, buckets, bc); |
||||||
|
} |
||||||
|
|
||||||
|
static void _handle_batch(struct merkle_sync* ms, uint64_t peer, const char* ns, |
||||||
|
const uint8_t* pl, size_t plen) { |
||||||
|
if (plen < 1) return; |
||||||
|
uint8_t count = pl[0]; const uint8_t* bp = pl + 1; size_t rem = plen - 1; |
||||||
|
|
||||||
|
for (uint8_t i = 0; i < count && rem >= 3; i++) { |
||||||
|
uint8_t lvl = bp[0]; uint8_t pb_i = bp[1]; rem -= 2; bp += 2; |
||||||
|
if (rem < pb_i + 1) break; |
||||||
|
uint64_t pr = _prefix_read(bp, pb_i); bp += pb_i; rem -= pb_i; |
||||||
|
uint8_t is_data = *bp++; rem--; |
||||||
|
|
||||||
|
if (is_data && rem >= 2) { |
||||||
|
uint16_t mc; memcpy(&mc, bp, 2); |
||||||
|
size_t item_len = rem; /* pass entire remainder to apply_items */ |
||||||
|
ms->ops->apply_items(ms->data_ctx, ns, bp, item_len); |
||||||
|
bp += rem; rem = 0; |
||||||
|
} else if (!is_data && rem >= 4) { |
||||||
|
uint8_t sub_pl[4096]; size_t sub_len = 0; |
||||||
|
uint8_t next_lvl = (uint8_t)(lvl < MT_MAX_LEVEL ? lvl + 1 : lvl); |
||||||
|
sub_pl[sub_len++] = next_lvl; |
||||||
|
sub_pl[sub_len++] = pb_i; |
||||||
|
_prefix_write(sub_pl + sub_len, pr, pb_i); sub_len += pb_i; |
||||||
|
sub_pl[sub_len++] = 0; /* is_data=0 */ |
||||||
|
uint32_t bm; memcpy(&bm, bp, 4); bp += 4; rem -= 4; |
||||||
|
memcpy(sub_pl + sub_len, &bm, 4); sub_len += 4; |
||||||
|
int nh = _popcount_u32(bm); |
||||||
|
if (rem >= (size_t)nh * MT_HASH_SIZE) { |
||||||
|
memcpy(sub_pl + sub_len, bp, (size_t)nh * MT_HASH_SIZE); |
||||||
|
sub_len += (size_t)nh * MT_HASH_SIZE; |
||||||
|
bp += (size_t)nh * MT_HASH_SIZE; rem -= (size_t)nh * MT_HASH_SIZE; |
||||||
|
} |
||||||
|
_handle_hashes(ms, peer, ns, sub_pl, sub_len); |
||||||
|
} |
||||||
|
} |
||||||
|
} |
||||||
|
|
||||||
|
static void _recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { |
||||||
|
if (!entry || entry->len < 4) { |
||||||
|
if (entry) { if (entry->dgram) u_free(entry->dgram); queue_entry_free(entry); } |
||||||
|
return; |
||||||
|
} |
||||||
|
if (!g_merkle || !g_merkle->initialized) { u_free(entry->dgram); queue_entry_free(entry); return; } |
||||||
|
uint64_t peer = conn ? conn->peer_node_id : 0; |
||||||
|
const uint8_t* d = entry->dgram; |
||||||
|
size_t dlen = entry->len; |
||||||
|
|
||||||
|
uint8_t ch_len = d[1]; |
||||||
|
if (dlen < (size_t)(2 + ch_len + 1)) { u_free(entry->dgram); queue_entry_free(entry); return; } |
||||||
|
char ns[64]; memcpy(ns, d + 2, ch_len); ns[ch_len] = '\0'; |
||||||
|
uint8_t type = d[2 + ch_len]; |
||||||
|
const uint8_t* pl = d + 3 + ch_len; |
||||||
|
size_t plen = dlen - 3 - ch_len; |
||||||
|
|
||||||
|
switch (type) { |
||||||
|
case 0x01: _handle_hashes(g_merkle, peer, ns, pl, plen); break; |
||||||
|
case 0x02: _handle_request(g_merkle, peer, ns, pl, plen); break; |
||||||
|
case 0x03: _handle_batch(g_merkle, peer, ns, pl, plen); break; |
||||||
|
} |
||||||
|
|
||||||
|
u_free(entry->dgram); queue_entry_free(entry); |
||||||
|
} |
||||||
|
|
||||||
|
/* ── Background consistency check ── */ |
||||||
|
|
||||||
|
int merkle_sync_bg_check(struct UTUN_INSTANCE* inst, const char* ns) { |
||||||
|
struct merkle_sync* ms = g_merkle; |
||||||
|
if (!ms || !ms->initialized) return -1; |
||||||
|
sqlite3* db = _db(inst); if (!db || !ns) return -1; |
||||||
|
|
||||||
|
sqlite3_stmt* stmt = NULL; |
||||||
|
if (sqlite3_prepare_v2(db, |
||||||
|
"SELECT prefix64 FROM merkle_tree_hash WHERE namespace=? AND level=? LIMIT 1", |
||||||
|
-1, &stmt, NULL) != SQLITE_OK) return -1; |
||||||
|
sqlite3_bind_text(stmt, 1, ns, -1, SQLITE_STATIC); |
||||||
|
sqlite3_bind_int(stmt, 2, MT_MAX_LEVEL); |
||||||
|
|
||||||
|
if (sqlite3_step(stmt) != SQLITE_ROW) { sqlite3_finalize(stmt); return -1; } |
||||||
|
uint64_t pref = (uint64_t)sqlite3_column_int64(stmt, 0); |
||||||
|
sqlite3_finalize(stmt); |
||||||
|
|
||||||
|
uint8_t old_hash[MT_HASH_SIZE]; |
||||||
|
memcpy(old_hash, merkle_sync_get_hash(inst, ns, MT_MAX_LEVEL, pref), MT_HASH_SIZE); |
||||||
|
|
||||||
|
_recompute_bucket(ms, ns, MT_MAX_LEVEL, pref); |
||||||
|
|
||||||
|
const uint8_t* new_hash = merkle_sync_get_hash(inst, ns, MT_MAX_LEVEL, pref); |
||||||
|
if (memcmp(old_hash, new_hash, MT_HASH_SIZE) != 0) { |
||||||
|
for (uint8_t lv = (uint8_t)(MT_MAX_LEVEL - 1); lv >= 1; lv--) |
||||||
|
_recompute_bucket(ms, ns, lv, merkle_sync_level_prefix(pref, lv)); |
||||||
|
return 1; |
||||||
|
} |
||||||
|
return 0; |
||||||
|
} |
||||||
|
|
||||||
|
static void _bg_timer_cb(void* arg) { |
||||||
|
struct merkle_sync* ms = (struct merkle_sync*)arg; |
||||||
|
if (!ms || !ms->initialized || !ms->inst) return; |
||||||
|
sqlite3* db = _db(ms->inst); if (!db) return; |
||||||
|
|
||||||
|
sqlite3_stmt* cs = NULL; |
||||||
|
sqlite3_prepare_v2(db, "SELECT channel_id FROM channels", -1, &cs, NULL); |
||||||
|
if (!cs) return; |
||||||
|
while (sqlite3_step(cs) == SQLITE_ROW) { |
||||||
|
const char* ch_id = (const char*)sqlite3_column_text(cs, 0); |
||||||
|
if (!ch_id) continue; |
||||||
|
int rc = merkle_sync_bg_check(ms->inst, ch_id); |
||||||
|
if (rc < 0) continue; |
||||||
|
break; |
||||||
|
} |
||||||
|
sqlite3_finalize(cs); |
||||||
|
|
||||||
|
ms->bg_timer = uasync_set_timeout(ms->inst->ua, (uint32_t)(MS_BG_INTERVAL_MS * 10), |
||||||
|
ms, _bg_timer_cb, "ms_bg"); |
||||||
|
} |
||||||
|
|
||||||
|
/* ── Public lifecycle ── */ |
||||||
|
|
||||||
|
int merkle_sync_init(struct UTUN_INSTANCE* inst, uint8_t svc_id, |
||||||
|
const struct merkle_sync_data_ops* ops, void* data_ctx) { |
||||||
|
if (!inst || !ops) return -1; |
||||||
|
struct merkle_sync* ms = u_calloc(1, sizeof(*ms)); |
||||||
|
if (!ms) return -1; |
||||||
|
ms->inst = inst; ms->svc_id = svc_id; ms->ops = ops; ms->data_ctx = data_ctx; |
||||||
|
ms->initialized = 1; |
||||||
|
g_merkle = ms; |
||||||
|
|
||||||
|
_ensure_table(ms); |
||||||
|
|
||||||
|
etcp_router_bind(inst, svc_id, _recv_cb); |
||||||
|
|
||||||
|
ms->bg_timer = uasync_set_timeout(inst->ua, (uint32_t)(MS_BG_INTERVAL_MS * 10), |
||||||
|
ms, _bg_timer_cb, "ms_bg"); |
||||||
|
|
||||||
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: initialized svc=%02x", MS_ID, svc_id); |
||||||
|
return 0; |
||||||
|
} |
||||||
|
|
||||||
|
void merkle_sync_destroy(struct UTUN_INSTANCE* inst) { |
||||||
|
struct merkle_sync* ms = g_merkle; |
||||||
|
if (!ms || !inst) return; |
||||||
|
ms->initialized = 0; g_merkle = NULL; |
||||||
|
|
||||||
|
etcp_router_bind(inst, ms->svc_id, NULL); |
||||||
|
|
||||||
|
if (ms->bg_timer) { uasync_cancel_timeout(inst->ua, ms->bg_timer); ms->bg_timer = NULL; } |
||||||
|
|
||||||
|
struct ms_session* s = ms->sessions; |
||||||
|
while (s) { struct ms_session* next = s->next; |
||||||
|
if (s->timer) uasync_cancel_timeout(inst->ua, s->timer); |
||||||
|
u_free(s); s = next; } |
||||||
|
|
||||||
|
u_free(ms); |
||||||
|
} |
||||||
|
|
||||||
|
/* ── Public async sync ── */ |
||||||
|
|
||||||
|
int merkle_sync_start(struct UTUN_INSTANCE* inst, uint64_t peer, |
||||||
|
const char* ns, merkle_sync_done_cb done_cb, void* arg) { |
||||||
|
struct merkle_sync* ms = g_merkle; |
||||||
|
if (!inst || !ns || !ms || !ms->initialized) return -1; |
||||||
|
_ensure_table(ms); |
||||||
|
|
||||||
|
struct ms_session* s = _session_find(ms, peer, ns); |
||||||
|
if (!s) { |
||||||
|
s = u_calloc(1, sizeof(*s)); |
||||||
|
if (!s) return -1; |
||||||
|
snprintf(s->ns, sizeof(s->ns), "%s", ns); |
||||||
|
s->peer = peer; |
||||||
|
s->next = ms->sessions; |
||||||
|
ms->sessions = s; |
||||||
|
} else { |
||||||
|
/* replace callback if session already exists */ |
||||||
|
if (s->timer) { uasync_cancel_timeout(inst->ua, s->timer); s->timer = NULL; } |
||||||
|
} |
||||||
|
s->active = 1; s->retries = 0; |
||||||
|
s->done_cb = done_cb; s->cb_arg = arg; |
||||||
|
|
||||||
|
_send_hashes(ms, peer, ns, 1, 0, 1, 0); |
||||||
|
_session_start_timer(s); |
||||||
|
return 0; |
||||||
|
} |
||||||
|
|
||||||
|
void merkle_sync_cancel(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ns) { |
||||||
|
struct merkle_sync* ms = g_merkle; |
||||||
|
if (!ms || !ns) return; |
||||||
|
struct ms_session** p = &ms->sessions; |
||||||
|
while (*p) { |
||||||
|
struct ms_session* s = *p; |
||||||
|
if (s->peer == peer && strcmp(s->ns, ns) == 0) { |
||||||
|
if (s->timer) { uasync_cancel_timeout(inst->ua, s->timer); s->timer = NULL; } |
||||||
|
*p = s->next; u_free(s); |
||||||
|
return; |
||||||
|
} |
||||||
|
p = &(*p)->next; |
||||||
|
} |
||||||
|
} |
||||||
@ -0,0 +1,250 @@ |
|||||||
|
#ifndef MERKLE_SYNC_H |
||||||
|
#define MERKLE_SYNC_H |
||||||
|
|
||||||
|
#include <stdint.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 — десериализовать и сохранить полученные элементы |
||||||
|
* |
||||||
|
* Потребитель (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_set_online(inst, node_id, 1); |
||||||
|
* |
||||||
|
* // ── отменил — коллбэк не вызовется ──
|
||||||
|
* 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 сохраняет, пересчитывает хеши |
||||||
|
* Рекурсивно для подбакетов, пока хеши не совпадут. |
||||||
|
* |
||||||
|
* ── Сессии ── |
||||||
|
* |
||||||
|
* Таймаут 10 секунд, 3 ретрая. При таймауте сессия перезапускает |
||||||
|
* MSG_HASHES с level=1. После 3 ретраев — done_cb с MT_ERR_TIMEOUT. |
||||||
|
*/ |
||||||
|
|
||||||
|
/* ── Data model callbacks ── */ |
||||||
|
|
||||||
|
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][join_sig:64][online:1][addr_cnt:1][addrs:var]... |
||||||
|
* |
||||||
|
* Возвращает 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, |
||||||
|
const uint8_t* data, size_t len); |
||||||
|
}; |
||||||
|
|
||||||
|
/*
|
||||||
|
* Коллбэк завершения синхронизации. |
||||||
|
* |
||||||
|
* 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_route_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=1) и запускает таймаут 10s. |
||||||
|
* |
||||||
|
* Если сессия для (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 НЕ вызывается. |
||||||
|
* Безопасно вызывать, если сессия не существует. |
||||||
|
*/ |
||||||
|
void merkle_sync_cancel(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ns); |
||||||
|
|
||||||
|
/* ── Recompute tree path after data change ── */ |
||||||
|
|
||||||
|
/*
|
||||||
|
* Пересчитать Merkle-дерево для ключа key — все уровни от 1 до 5. |
||||||
|
* Вызывается после изменения данных: merkle_sync сам не знает, |
||||||
|
* когда данные изменились — потребитель должен вызвать явно. |
||||||
|
* member_sync делает это внутри put/del/set_online. |
||||||
|
*/ |
||||||
|
void merkle_sync_recompute_path(struct UTUN_INSTANCE* inst, const char* ns, uint64_t key); |
||||||
|
|
||||||
|
/* ── 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 */ |
||||||
Loading…
Reference in new issue