Browse Source

merkle_sync: fix BATCH parse, synthetic HASHES, initiator done_cb

- Add 2-byte dlen prefix to terminal BATCH/HASHES data so receiver knows
  exact byte count — fixes rem=0 bug that truncated multi-item BATCHes.
- Fix synthetic HASHES prefix_bytes: use merkle_sync_prefix_bytes(next_lvl)
  for both declaration and _prefix_write, not current level's pb_i.
- Fix initiator never getting _session_done: responder sends confirmation
  HASHES(is_data=1,level=0) before _session_done so initiator's done_cb fires.
- Add level==0 guard in _member_get_items/_member_update_bucket_hash
  to prevent SQL queries with bogus mask on virtual root level.
topo_upd
Evgeny 2 months ago
parent
commit
6bf8482558
  1. 4
      tools/chatgui/transport/member_sync.c
  2. 43
      tools/chatgui/transport/merkle_sync.c

4
tools/chatgui/transport/member_sync.c

@ -88,9 +88,10 @@ static void _compute_member_hash(uint64_t node_id, const uint8_t* x25519,
/* ── merkle_sync_data_ops implementation ── */
static int _member_update_bucket_hash(void* ctx, const char* ns, uint8_t level,
uint64_t prefix64, EVP_MD_CTX* sha_ctx) {
uint64_t prefix64, EVP_MD_CTX* sha_ctx) {
struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)ctx;
sqlite3* db = _db(inst); if (!db) return -1;
if (level == 0) return 0;
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: bucket_hash ns=%s L%d/P%016llx", MS_ID, ns, level, (unsigned long long)prefix64);
char peers_tbl[128]; _peers_table(ns, peers_tbl, sizeof(peers_tbl));
@ -131,6 +132,7 @@ static int _member_get_items(void* ctx, const char* ns, uint8_t level,
uint8_t* buf, size_t* len) {
struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)ctx;
sqlite3* db = _db(inst); if (!db || !buf || !len) return -1;
if (level == 0) { *len = 0; return 0; }
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: get_items ns=%s L%d/P%016llx", MS_ID, ns, level, (unsigned long long)prefix);
char peers_tbl[128]; _peers_table(ns, peers_tbl, sizeof(peers_tbl));

43
tools/chatgui/transport/merkle_sync.c

@ -250,9 +250,10 @@ static int _send_hashes(struct merkle_sync* ms, uint64_t peer, const char* ns,
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; }
if (ms->ops->get_items(ms->data_ctx, ns, level, prefix, prefix_bytes, mbuf, &mlen) == 0) {
uint16_t dlen = (uint16_t)mlen; memcpy(p, &dlen, 2); p += 2;
memcpy(p, mbuf, mlen); p += mlen;
} else { uint16_t zero = 0; memcpy(p, &zero, 2); p += 2; }
u_free(mbuf);
}
} else {
@ -299,9 +300,10 @@ static int _send_batch(struct merkle_sync* ms, uint64_t peer, const char* ns,
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; }
buckets[i].prefix, buckets[i].prefix_bytes, mbuf, &mlen) == 0) {
uint16_t dlen = (uint16_t)mlen; memcpy(p, &dlen, 2); p += 2;
memcpy(p, mbuf, mlen); p += mlen;
} else { uint16_t z = 0; memcpy(p, &z, 2); p += 2; }
u_free(mbuf);
}
} else {
@ -349,7 +351,10 @@ static void _handle_hashes(struct merkle_sync* ms, uint64_t peer, const char* ns
if (is_data) {
if (paylen < 2) return;
ms->ops->apply_items(ms->data_ctx, ns, payload, paylen);
uint16_t dlen; memcpy(&dlen, payload, 2);
const uint8_t* pd = payload + 2;
if (paylen - 2 < dlen) return;
ms->ops->apply_items(ms->data_ctx, ns, pd, dlen);
if (s && s->active) _session_done(s, MT_OK);
return;
}
@ -382,7 +387,11 @@ static void _handle_hashes(struct merkle_sync* ms, uint64_t peer, const char* ns
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: handle_hashes peer=%016llx ns=%s level=%d remote_bm=%08x local_bm=%08x differs=%08x is_data=%d",
MS_ID, (unsigned long long)peer, ns, level, remote_bm, local_bm, differs, is_data);
if (differs == 0) { _session_done(s, MT_OK); return; }
if (differs == 0) {
_send_hashes(ms, peer, ns, 0, 0, 0, 1);
_session_done(s, MT_OK);
return;
}
int next_shift = 63 - ((int)level + 1) * 5;
struct bucket_entry requests[MT_BUCKETS]; int rcount = 0;
@ -465,18 +474,19 @@ static void _handle_batch(struct merkle_sync* ms, uint64_t peer, const char* ns,
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 item_count; memcpy(&item_count, bp, 2);
if (rem < 2 + (size_t)item_count * 2) break;
ms->ops->apply_items(ms->data_ctx, ns, bp, rem);
bp += rem; rem = 0;
if (is_data && rem >= 4) {
uint16_t dlen; memcpy(&dlen, bp, 2); bp += 2; rem -= 2;
if (rem < dlen) break;
ms->ops->apply_items(ms->data_ctx, ns, bp, dlen);
bp += dlen; rem -= dlen;
} else if (!is_data && rem >= 4) {
all_terminal = 0;
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++] = merkle_sync_prefix_bytes(next_lvl);
_prefix_write(sub_pl + sub_len, pr, pb_i); sub_len += pb_i;
_prefix_write(sub_pl + sub_len, pr, merkle_sync_prefix_bytes(next_lvl));
sub_len += merkle_sync_prefix_bytes(next_lvl);
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;
@ -491,7 +501,10 @@ static void _handle_batch(struct merkle_sync* ms, uint64_t peer, const char* ns,
}
if (all_terminal) {
struct ms_session* s = _session_find(ms, peer, ns);
if (s && s->active) _session_done(s, MT_OK);
if (s && s->active) {
_send_hashes(ms, peer, ns, 0, 0, 0, 1);
_session_done(s, MT_OK);
}
}
}

Loading…
Cancel
Save