Browse Source

chat: передавать group_id в бинарном виде в протоколе синхронизации

db_sync/merkle_sync/chat_sync шлют group_id (uint64 BE) вместо текстового
channel_id/namespace. Приёмная сторона находит группу/таблицу через
topo_groups_find(group_id) и требует существующую CHAT-группу. db_sync:
instance key заменён с SHA256(ch_id) на strtoull(ch_id)=group_id.
v2
evgeny 4 weeks ago
parent
commit
7b73d99b5b
  1. 10
      src/chat/chat_channel.c
  2. 40
      src/chat/chat_sync.c
  3. 80
      src/chat/db_sync.c
  4. 5
      src/chat/db_sync.h
  5. 69
      src/chat/merkle_sync.c
  6. 137
      tests/test_merkle_sync.c

10
src/chat/chat_channel.c

@ -58,11 +58,9 @@ void chat_core_ensure_channel_ready(const char* ch_id) {
char tbl_msg[80]; msg_table_name(ch_id, tbl_msg, sizeof(tbl_msg)); char tbl_msg[80]; msg_table_name(ch_id, tbl_msg, sizeof(tbl_msg));
uint64_t ch_hash = 0;
{ const uint8_t* chd = (const uint8_t*)ch_id; size_t chl = strlen(ch_id); uint8_t sh[32]; SHA256(chd, chl, sh); memcpy(&ch_hash, sh, 8); }
struct DB_SYNC_INSTANCE* si = db_sync_instance_add(g_cc.inst, tbl_msg, ch_hash, 1);
uint64_t gid = strtoull(ch_id, NULL, 10); uint64_t gid = strtoull(ch_id, NULL, 10);
struct DB_SYNC_INSTANCE* si = db_sync_instance_add(g_cc.inst, tbl_msg, gid, 1);
if (gid != 0 && g_cc.inst->topo_groups && !topo_groups_find(g_cc.inst->topo_groups, gid)) { if (gid != 0 && g_cc.inst->topo_groups && !topo_groups_find(g_cc.inst->topo_groups, gid)) {
struct TOPO_GROUP* g = topo_groups_create_group(g_cc.inst->topo_groups, gid, TOPO_GROUP_TYPE_CHAT, ch_id); struct TOPO_GROUP* g = topo_groups_create_group(g_cc.inst->topo_groups, gid, TOPO_GROUP_TYPE_CHAT, ch_id);
if (g) member_sync_subscribe_group(g_cc.inst, g); if (g) member_sync_subscribe_group(g_cc.inst, g);
@ -70,8 +68,8 @@ void chat_core_ensure_channel_ready(const char* ch_id) {
if (si) { if (si) {
si_register(si, ch_id); si_register(si, ch_id);
db_sync_set_insert_cb(si, on_msg_inserted, u_strdup(ch_id)); db_sync_set_insert_cb(si, on_msg_inserted, u_strdup(ch_id));
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "%s: channel ready ch=%s tbl=%s hash=0x%016llx", DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "%s: channel ready ch=%s tbl=%s group=0x%016llx",
CC_ID, ch_id, tbl_msg, (unsigned long long)ch_hash); CC_ID, ch_id, tbl_msg, (unsigned long long)gid);
} else { } else {
DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: db_sync_instance_add failed for ch=%s", DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: db_sync_instance_add failed for ch=%s",
CC_ID, ch_id); CC_ID, ch_id);

40
src/chat/chat_sync.c

@ -102,18 +102,33 @@ static struct ETCP_CONN* cs_find_conn_for_node(struct UTUN_INSTANCE* inst, uint6
/* ── Send ── */ /* ── Send ── */
/* group_id в wire: [group_id:8] вместо [ch_len:1][ch_id] */
static void cs_write_group_id(uint8_t* p, const char* ch_id) {
uint64_t gb = htobe64(strtoull(ch_id, NULL, 10));
memcpy(p, &gb, 8);
}
/* group_id → ch_id через группу канала (требует существующую CHAT-группу). 0=ок, -1=нет группы */
static int cs_group_id_to_ch_id(struct UTUN_INSTANCE* inst, uint64_t group_id, char* ch_id, size_t sz) {
if (!inst || !inst->topo_groups) return -1;
struct TOPO_GROUP* g = topo_groups_find(inst->topo_groups, group_id);
if (!g || !g->channel_id[0]) {
DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: recv — no CHAT group for group=%016llx, drop", CS_ID, (unsigned long long)group_id);
return -1;
}
snprintf(ch_id, sz, "%s", g->channel_id);
return 0;
}
static int cs_send(struct chat_sync* cs, const char* ch_id, uint64_t dst, static int cs_send(struct chat_sync* cs, const char* ch_id, uint64_t dst,
const uint8_t* payload, size_t len) { const uint8_t* payload, size_t len) {
if (len < 1) { DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: cs_send len<1", CS_ID); return -1; } if (len < 1) { DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: cs_send len<1", CS_ID); return -1; }
uint8_t ch_len = (uint8_t)strlen(ch_id); size_t total = 1 + 8 + len;
if (ch_len > 63) { DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: cs_send ch_len=%u > 63", CS_ID, ch_len); return -1; }
size_t total = 1 + 1 + ch_len + len;
uint8_t* buf = u_malloc(total); uint8_t* buf = u_malloc(total);
if (!buf) { DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: cs_send alloc(%zu) failed", CS_ID, total); return -1; } if (!buf) { DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: cs_send alloc(%zu) failed", CS_ID, total); return -1; }
uint8_t* p = buf; uint8_t* p = buf;
*p++ = ETCP_RT_ID_CHAT_SYNC; *p++ = ETCP_RT_ID_CHAT_SYNC;
*p++ = ch_len; cs_write_group_id(p, ch_id); p += 8;
memcpy(p, ch_id, ch_len); p += ch_len;
memcpy(p, payload, len); memcpy(p, payload, len);
struct ll_entry* entry = queue_entry_new(0); struct ll_entry* entry = queue_entry_new(0);
if (!entry) { DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: cs_send queue_entry_new failed", CS_ID); u_free(buf); return -1; } if (!entry) { DEBUG_ERROR(DEBUG_CATEGORY_MEMBER_SYNC, "%s: cs_send queue_entry_new failed", CS_ID); u_free(buf); return -1; }
@ -160,7 +175,7 @@ static int cs_send_channel_invite(struct chat_sync* cs, struct UTUN_INSTANCE* in
/* ── Recv dispatcher ── */ /* ── Recv dispatcher ── */
static void chat_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { static void chat_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
if (!entry || entry->len < 4) { if (!entry || entry->len < 10) {
if (entry) { if (entry) {
DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: recv short entry len=%zu conn=%s", CS_ID, entry->len, conn ? conn->log_name : "?"); DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: recv short entry len=%zu conn=%s", CS_ID, entry->len, conn ? conn->log_name : "?");
if (entry->dgram) u_free(entry->dgram); queue_entry_free(entry); if (entry->dgram) u_free(entry->dgram); queue_entry_free(entry);
@ -173,15 +188,14 @@ static void chat_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
const uint8_t* d = entry->dgram; const uint8_t* d = entry->dgram;
size_t dlen = entry->len; size_t dlen = entry->len;
uint8_t ch_len = d[1]; uint64_t group_id = be64toh(*(const uint64_t*)(d + 1));
if (dlen < (size_t)(2 + ch_len + 1)) { u_free(entry->dgram); queue_entry_free(entry); return; } char ch_id[64];
if (cs_group_id_to_ch_id(g_cs->inst, group_id, ch_id, sizeof(ch_id)) != 0) { u_free(entry->dgram); queue_entry_free(entry); return; }
char ch_id[64]; memcpy(ch_id, d + 2, ch_len); ch_id[ch_len] = '\0'; uint8_t type = d[9];
uint8_t type = d[2 + ch_len];
DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: recv %s(%02x) from=%016llx ch=%s len=%zu", DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: recv %s(%02x) from=%016llx ch=%s len=%zu",
CS_ID, cs_msg_name(type), type, (unsigned long long)peer, ch_id, dlen); CS_ID, cs_msg_name(type), type, (unsigned long long)peer, ch_id, dlen);
const uint8_t* pl = d + 3 + ch_len; const uint8_t* pl = d + 10;
size_t plen = dlen - 3 - ch_len; size_t plen = dlen - 10;
DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: RECV %s from=%016llx ch=%s len=%zu", DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: RECV %s from=%016llx ch=%s len=%zu",
CS_ID, cs_msg_name(type), (unsigned long long)peer, ch_id, plen); CS_ID, cs_msg_name(type), (unsigned long long)peer, ch_id, plen);

80
src/chat/db_sync.c

@ -33,7 +33,7 @@ static void db_sync_peer_check_cb(void* arg);
static void db_sync_resume_peer_check(struct DB_SYNC* db); static void db_sync_resume_peer_check(struct DB_SYNC* db);
static void db_sync_instance_ttl_cb(void* arg); static void db_sync_instance_ttl_cb(void* arg);
static int db_sync_send(struct DB_SYNC_INSTANCE* si, uint64_t dst_node_id, const uint8_t* payload, size_t len); static int db_sync_send(struct DB_SYNC_INSTANCE* si, uint64_t dst_node_id, const uint8_t* payload, size_t len);
static int db_sync_send_hash(struct DB_SYNC* db, uint64_t node_id, uint64_t hash, const uint8_t* payload, size_t len); static int db_sync_send_gid(struct DB_SYNC* db, uint64_t node_id, uint64_t group_id, const uint8_t* payload, size_t len);
static void db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si, uint64_t peer_node_id); static void db_sync_initiate_sync(struct DB_SYNC_INSTANCE* si, uint64_t peer_node_id);
static int si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32_t from, uint32_t count, uint32_t vp, uint32_t want_from); static int si_send_data_batch(struct DB_SYNC_INSTANCE* si, uint64_t dst, uint32_t from, uint32_t count, uint32_t vp, uint32_t want_from);
static void si_delivery_update(struct DB_SYNC_INSTANCE* si, uint64_t ts, uint64_t author, uint64_t peer_id); static void si_delivery_update(struct DB_SYNC_INSTANCE* si, uint64_t ts, uint64_t author, uint64_t peer_id);
@ -54,7 +54,7 @@ struct SI_PEER {
struct DB_SYNC_INSTANCE { struct DB_SYNC_INSTANCE {
struct DB_SYNC* db_sync; struct DB_SYNC* db_sync;
uint64_t hash; uint64_t group_id;
char table_name[64]; char table_name[64];
uint64_t next_id; uint64_t next_id;
uint64_t last_timestamp_ms; uint64_t last_timestamp_ms;
@ -128,7 +128,7 @@ void db_sync_remove_done_cbk(struct UTUN_INSTANCE* inst, db_sync_done_fn fn, voi
#define SI_TBL(si) ((si)->table_name) #define SI_TBL(si) ((si)->table_name)
#define SI_SHRT(si) ({ \ #define SI_SHRT(si) ({ \
static char _tbuf[28]; \ static char _tbuf[28]; \
snprintf(_tbuf, sizeof(_tbuf), "%s.%04llX", SI_TBL(si), (unsigned long long)((si)->hash & 0xFFFF)); \ snprintf(_tbuf, sizeof(_tbuf), "%s.%04llX", SI_TBL(si), (unsigned long long)((si)->group_id & 0xFFFF)); \
_tbuf; \ _tbuf; \
}) })
@ -220,11 +220,11 @@ static void db_chain_hash_compute(const uint8_t prev_chain_hash[32],
// Instance management // Instance management
// ============================================================ // ============================================================
// Ищет активный инстанс синхронизации по 64-битному хешу // Ищет активный инстанс синхронизации по group_id канала
static struct DB_SYNC_INSTANCE* db_instance_find(struct DB_SYNC* db, uint64_t hash) static struct DB_SYNC_INSTANCE* db_instance_find(struct DB_SYNC* db, uint64_t group_id)
{ {
for (int i = 0; i < db->instance_count; i++) { for (int i = 0; i < db->instance_count; i++) {
if (db->instances[i].hash == hash && db->instances[i].enabled) return &db->instances[i]; if (db->instances[i].group_id == group_id && db->instances[i].enabled) return &db->instances[i];
} }
return NULL; return NULL;
} }
@ -672,24 +672,24 @@ static int db_record_insert_cascade(struct DB_SYNC_INSTANCE* si,
// Send // Send
// ============================================================ // ============================================================
// Отправляет сообщение узлу через ETCP с заголовком: [ETCP_RT_ID_DB_SYNC:1][hash_be:8][payload]. Возвращает код etcp_send // Отправляет сообщение узлу через ETCP с заголовком: [ETCP_RT_ID_DB_SYNC:1][group_id_be:8][payload]. Возвращает код etcp_send
static int db_sync_send_hash(struct DB_SYNC* db, uint64_t node_id, uint64_t hash, const uint8_t* payload, size_t plen) static int db_sync_send_gid(struct DB_SYNC* db, uint64_t node_id, uint64_t group_id, const uint8_t* payload, size_t plen)
{ {
struct ll_entry* entry = queue_entry_new(0); struct ll_entry* entry = queue_entry_new(0);
if (!entry) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "queue_entry_new"); return -1; } if (!entry) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "queue_entry_new"); return -1; }
uint8_t* buf = u_malloc(plen + 9); uint8_t* buf = u_malloc(plen + 9);
if (!buf) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "db_sync_send_hash: u_malloc(%zu) failed for dst=%016llx", plen + 9, (unsigned long long)node_id); queue_entry_free(entry); return -1; } if (!buf) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "db_sync_send_gid: u_malloc(%zu) failed for dst=%016llx", plen + 9, (unsigned long long)node_id); queue_entry_free(entry); return -1; }
buf[0] = ETCP_RT_ID_DB_SYNC; buf[0] = ETCP_RT_ID_DB_SYNC;
uint64_t hb = htobe64(hash); uint64_t gb = htobe64(group_id);
memcpy(buf + 1, &hb, 8); memcpy(buf + 1, &gb, 8);
memcpy(buf + 9, payload, plen); memcpy(buf + 9, payload, plen);
entry->dgram = buf; entry->dgram = buf;
entry->len = plen + 9; entry->len = plen + 9;
struct ETCP_CONN* conn = instance_find_conn(db->inst, node_id); struct ETCP_CONN* conn = instance_find_conn(db->inst, node_id);
if (!conn || !conn->links_up) { if (!conn || !conn->links_up) {
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync: send DROP — no conn to %04llX (hash=%016llx, type=%02x)", DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync: send DROP — no conn to %04llX (group=%016llx, type=%02x)",
(unsigned long long)(node_id >> 16), hash, plen > 0 ? payload[0] : 0, (unsigned long long)(node_id >> 16), group_id, plen > 0 ? payload[0] : 0,
(void*)conn, conn ? conn->links_up : -1); (void*)conn, conn ? conn->links_up : -1);
queue_entry_free(entry); queue_entry_free(entry);
return -1; return -1;
@ -700,10 +700,10 @@ static int db_sync_send_hash(struct DB_SYNC* db, uint64_t node_id, uint64_t hash
return ret; return ret;
} }
// Отправляет сообщение через ETCP от имени конкретного инстанса (подставляет хеш инстанса) // Отправляет сообщение через ETCP от имени конкретного инстанса (подставляет group_id инстанса)
static int db_sync_send(struct DB_SYNC_INSTANCE* si, uint64_t dst_node_id, const uint8_t* payload, size_t len) static int db_sync_send(struct DB_SYNC_INSTANCE* si, uint64_t dst_node_id, const uint8_t* payload, size_t len)
{ {
return db_sync_send_hash(si->db_sync, dst_node_id, si->hash, payload, len); return db_sync_send_gid(si->db_sync, dst_node_id, si->group_id, payload, len);
} }
// ============================================================ // ============================================================
@ -1202,15 +1202,15 @@ static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry)
if (!db || !db->enabled) { queue_dgram_free(entry); queue_entry_free(entry); return; } if (!db || !db->enabled) { queue_dgram_free(entry); queue_entry_free(entry); return; }
uint64_t src = conn ? conn->peer_node_id : 0; uint64_t src = conn ? conn->peer_node_id : 0;
uint64_t hash = be64toh(*(uint64_t*)(entry->dgram + 1)); uint64_t group_id = be64toh(*(uint64_t*)(entry->dgram + 1));
uint8_t type = entry->dgram[9]; uint8_t type = entry->dgram[9];
const uint8_t* payload = entry->dgram + 10; const uint8_t* payload = entry->dgram + 10;
size_t plen = entry->len - 10; size_t plen = entry->len - 10;
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "recv type=%02x from=%016llx hash=%016llx len=%zu", type, (unsigned long long)src, (unsigned long long)hash, plen); DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "recv type=%02x from=%016llx group=%016llx len=%zu", type, (unsigned long long)src, (unsigned long long)group_id, plen);
struct DB_SYNC_INSTANCE* si = db_instance_find(db, hash); struct DB_SYNC_INSTANCE* si = db_instance_find(db, group_id);
if (!si) { if (!si) {
/* lazy-register: find table msg_<ch_id> with matching SHA256 */ /* lazy-register: найти таблицу msg_<ch_id>, у которой strtoull(ch_id) == group_id */
sqlite3_stmt* st = NULL; sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db->db, if (sqlite3_prepare_v2(db->db,
"SELECT name, substr(name,5) FROM sqlite_master WHERE type='table' AND name LIKE 'msg_%'", "SELECT name, substr(name,5) FROM sqlite_master WHERE type='table' AND name LIKE 'msg_%'",
@ -1219,11 +1219,11 @@ static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry)
const char* tbl = (const char*)sqlite3_column_text(st, 0); const char* tbl = (const char*)sqlite3_column_text(st, 0);
const char* ch_id = (const char*)sqlite3_column_text(st, 1); const char* ch_id = (const char*)sqlite3_column_text(st, 1);
if (!ch_id || !ch_id[0]) continue; if (!ch_id || !ch_id[0]) continue;
uint64_t h; { uint8_t sh[32]; SHA256((const uint8_t*)ch_id, strlen(ch_id), sh); memcpy(&h, sh, 8); } uint64_t gid = strtoull(ch_id, NULL, 10);
if (h == hash) { if (gid == group_id) {
si = db_sync_instance_add(db->inst, tbl, hash, 1); si = db_sync_instance_add(db->inst, tbl, group_id, 1);
if (si) { si_peer_add(si, src); } if (si) { si_peer_add(si, src); }
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "lazy-register: tbl=%s ch=%s hash=%016llx", tbl, ch_id, (unsigned long long)hash); DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "lazy-register: tbl=%s ch=%s group=%016llx", tbl, ch_id, (unsigned long long)group_id);
break; break;
} }
} }
@ -1231,17 +1231,17 @@ static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry)
} }
if (!si) { if (!si) {
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC,
"recv msg type=0x%02x from %016llx hash=%016llx — instance NOT FOUND; sending DB_ERR_NOT_FOUND", "recv msg type=0x%02x from %016llx group=%016llx — instance NOT FOUND; sending DB_ERR_NOT_FOUND",
type, (unsigned long long)src, (unsigned long long)hash); type, (unsigned long long)src, (unsigned long long)group_id);
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "→ instance not found for hash=%016llx", (unsigned long long)hash); DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "→ instance not found for group=%016llx", (unsigned long long)group_id);
uint8_t err[2]; err[0] = DB_MSG_ERROR; err[1] = DB_ERR_NOT_FOUND; uint8_t err[2]; err[0] = DB_MSG_ERROR; err[1] = DB_ERR_NOT_FOUND;
db_sync_send_hash(db, src, hash, err, 2); db_sync_send_gid(db, src, group_id, err, 2);
queue_dgram_free(entry); queue_entry_free(entry); return; queue_dgram_free(entry); queue_entry_free(entry); return;
} }
} }
if (!si->enabled) { if (!si->enabled) {
uint8_t err[2]; err[0] = DB_MSG_ERROR; err[1] = DB_ERR_DISABLED; uint8_t err[2]; err[0] = DB_MSG_ERROR; err[1] = DB_ERR_DISABLED;
db_sync_send_hash(db, src, hash, err, 2); db_sync_send_gid(db, src, group_id, err, 2);
queue_dgram_free(entry); queue_entry_free(entry); return; queue_dgram_free(entry); queue_entry_free(entry); return;
} }
@ -1255,8 +1255,8 @@ static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry)
case DB_MSG_SYNC_DONE: DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "→ handle SYNC_DONE"); db_handle_sync_done(si, src, payload, plen); break; case DB_MSG_SYNC_DONE: DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "→ handle SYNC_DONE"); db_handle_sync_done(si, src, payload, plen); break;
case DB_MSG_ERROR: DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "→ handle ERROR"); db_handle_error(si, src, payload, plen); break; case DB_MSG_ERROR: DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "→ handle ERROR"); db_handle_error(si, src, payload, plen); break;
default: default:
DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC, "sync: unknown msg type 0x%02x from %04llX (hash=%016llx) — dropped", DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC, "sync: unknown msg type 0x%02x from %04llX (group=%016llx) — dropped",
type, (unsigned long long)(src >> 16), hash); type, (unsigned long long)(src >> 16), group_id);
break; break;
} }
queue_dgram_free(entry); queue_dgram_free(entry);
@ -1707,19 +1707,19 @@ void db_sync_destroy(struct UTUN_INSTANCE* inst)
} }
// Создаёт SQLite-таблицу для синхронизации, регистрирует пиров, запускает первичную синхронизацию и TTL-таймер // Создаёт SQLite-таблицу для синхронизации, регистрирует пиров, запускает первичную синхронизацию и TTL-таймер
struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const char* table_name, uint64_t hash, int auto_sync) struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const char* table_name, uint64_t group_id, int auto_sync)
{ {
if (!inst || !inst->db_sync || !table_name || !table_name[0]) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "instance_add invalid args"); return NULL; } if (!inst || !inst->db_sync || !table_name || !table_name[0]) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "instance_add invalid args"); return NULL; }
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "table=%s hash=%016llx", table_name, (unsigned long long)hash); DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "table=%s group=%016llx", table_name, (unsigned long long)group_id);
struct DB_SYNC* db = inst->db_sync; struct DB_SYNC* db = inst->db_sync;
if (!db->enabled || !db->db) return NULL; if (!db->enabled || !db->db) return NULL;
struct DB_SYNC_INSTANCE* si = db_instance_find(db, hash); struct DB_SYNC_INSTANCE* si = db_instance_find(db, group_id);
if (si) { DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "instance already exists hash=%016llx", (unsigned long long)hash); return si; } if (si) { DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "instance already exists group=%016llx", (unsigned long long)group_id); return si; }
si = db_instance_alloc(db); si = db_instance_alloc(db);
if (!si) return NULL; if (!si) return NULL;
si->hash = hash; si->group_id = group_id;
snprintf(si->table_name, sizeof(si->table_name), "%s", table_name); snprintf(si->table_name, sizeof(si->table_name), "%s", table_name);
// Create table — PK is (timestamp, node_id), author_signature NOT NULL // Create table — PK is (timestamp, node_id), author_signature NOT NULL
@ -1762,7 +1762,7 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const
db_verify_chain(si); db_verify_chain(si);
DEBUG_DEBUG(DEBUG_CATEGORY_CHAT_SYNC, "instance_add: tbl=%s mc=%u next_id=%llu hash=%016llx", SI_TBL(si), db_count(si), (unsigned long long)si->next_id, (unsigned long long)si->hash); DEBUG_DEBUG(DEBUG_CATEGORY_CHAT_SYNC, "instance_add: tbl=%s mc=%u next_id=%llu group=%016llx", SI_TBL(si), db_count(si), (unsigned long long)si->next_id, (unsigned long long)si->group_id);
si->ttl_timer = uasync_set_timeout(inst->ua, DB_SYNC_TTL_INTERVAL * 10000u, si, db_sync_instance_ttl_cb, "db_sync_ttl"); si->ttl_timer = uasync_set_timeout(inst->ua, DB_SYNC_TTL_INTERVAL * 10000u, si, db_sync_instance_ttl_cb, "db_sync_ttl");
@ -1783,8 +1783,8 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const
} }
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC,
"instance added table=%s tbl=%s hash=%016llx next_id=%llu peers_found=%d peers_synced=%d mc=%u", "instance added table=%s tbl=%s group=%016llx next_id=%llu peers_found=%d peers_synced=%d mc=%u",
table_name, si->table_name, (unsigned long long)hash, table_name, si->table_name, (unsigned long long)group_id,
(unsigned long long)si->next_id, peers_found, peers_synced, db_count(si)); (unsigned long long)si->next_id, peers_found, peers_synced, db_count(si));
return si; return si;
} }
@ -1793,7 +1793,7 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const
void db_sync_instance_remove(struct DB_SYNC_INSTANCE* si) void db_sync_instance_remove(struct DB_SYNC_INSTANCE* si)
{ {
if (!si) return; if (!si) return;
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "tbl=%s hash=%016llx", si->table_name, (unsigned long long)si->hash); DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "tbl=%s group=%016llx", si->table_name, (unsigned long long)si->group_id);
struct DB_SYNC* db = si->db_sync; struct DB_SYNC* db = si->db_sync;
struct UTUN_INSTANCE* inst = db->inst; struct UTUN_INSTANCE* inst = db->inst;
@ -1808,7 +1808,7 @@ void db_sync_instance_remove(struct DB_SYNC_INSTANCE* si)
(db->instance_count - idx - 1) * sizeof(*db->instances)); (db->instance_count - idx - 1) * sizeof(*db->instances));
db->instance_count--; db->instance_count--;
} }
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "instance removed tbl=%s hash=%016llx", si->table_name, (unsigned long long)si->hash); DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "instance removed tbl=%s group=%016llx", si->table_name, (unsigned long long)si->group_id);
} }
// Вставляет подписанную запись: проверяет подпись, вставляет в БД, рассылает PUSH всем синхронизированным пирам // Вставляет подписанную запись: проверяет подпись, вставляет в БД, рассылает PUSH всем синхронизированным пирам

5
src/chat/db_sync.h

@ -24,8 +24,7 @@
// //
// Синхронизация: // Синхронизация:
// - При поднятии ETCP-соединения с пиром для каждого инстанса запускается полная синхронизация // - При поднятии ETCP-соединения с пиром для каждого инстанса запускается полная синхронизация
// - Каждое сообщение содержит instance_hash (первые 64 бита SHA256(name||id_be)), // - Каждое сообщение содержит group_id (uint64, binary), что позволяет маршрутизировать сообщения к нужному инстансу на приёмной стороне
// что позволяет маршрутизировать сообщения к нужному инстансу на приёмной стороне
// - Если пир не имеет инстанса с таким хешем — возвращает DB_MSG_ERROR(0x01) // - Если пир не имеет инстанса с таким хешем — возвращает DB_MSG_ERROR(0x01)
// - Если пир добавляет инстанс позже — сам инициирует sync в нашу сторону // - Если пир добавляет инстанс позже — сам инициирует sync в нашу сторону
// - Протокол: сравнение chain_hash для поиска расхождений (7 типов сообщений) // - Протокол: сравнение chain_hash для поиска расхождений (7 типов сообщений)
@ -97,7 +96,7 @@ int db_sync_init(struct UTUN_INSTANCE* inst);
void db_sync_destroy(struct UTUN_INSTANCE* inst); void db_sync_destroy(struct UTUN_INSTANCE* inst);
// Instance management // Instance management
struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const char* table_name, uint64_t hash, int auto_sync); struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const char* table_name, uint64_t group_id, int auto_sync);
void db_sync_instance_remove(struct DB_SYNC_INSTANCE* si); void db_sync_instance_remove(struct DB_SYNC_INSTANCE* si);
// Data operations (per-instance) // Data operations (per-instance)

69
src/chat/merkle_sync.c

@ -3,10 +3,12 @@
#include "../utun_instance.h" #include "../utun_instance.h"
#include "../transport_layer/etcp_api.h" #include "../transport_layer/etcp_api.h"
#include "../transport_layer/etcp.h" #include "../transport_layer/etcp.h"
#include "../routing_layer/topo_group.h"
#include "../../lib/debug_config.h" #include "../../lib/debug_config.h"
#include "../../lib/mem.h" #include "../../lib/mem.h"
#include "../../lib/u_async.h" #include "../../lib/u_async.h"
#include "../../lib/ll_queue.h" #include "../../lib/ll_queue.h"
#include "../../lib/platform_compat.h"
#ifdef UTUN_HAVE_STANDBY #ifdef UTUN_HAVE_STANDBY
#include "standby.h" #include "standby.h"
#endif #endif
@ -240,6 +242,29 @@ static int _get_level_hashes(struct merkle_sync* ms, const char* ns,
return 0; return 0;
} }
/* ── group_id в wire: [group_id:8] вместо [ns_len:1][ns] ── */
static uint64_t ms_ns_to_group_id(const char* ns) {
return strtoull(ns, NULL, 10);
}
static void ms_write_group_id(uint8_t* p, const char* ns) {
uint64_t gb = htobe64(ms_ns_to_group_id(ns));
memcpy(p, &gb, 8);
}
/* group_id → ns через группу канала (требует существующую CHAT-группу). 0=ок, -1=нет группы */
static int ms_group_id_to_ns(struct UTUN_INSTANCE* inst, uint64_t group_id, char* ns, size_t sz) {
if (!inst || !inst->topo_groups) return -1;
struct TOPO_GROUP* g = topo_groups_find(inst->topo_groups, group_id);
if (!g || !g->channel_id[0]) {
DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: recv — no CHAT group for group=%016llx, drop", MS_ID, (unsigned long long)group_id);
return -1;
}
snprintf(ns, sz, "%s", g->channel_id);
return 0;
}
/* ── Send helpers ── */ /* ── Send helpers ── */
static struct ETCP_CONN* ms_find_conn_for_node(struct UTUN_INSTANCE* inst, uint64_t node_id) { static struct ETCP_CONN* ms_find_conn_for_node(struct UTUN_INSTANCE* inst, uint64_t node_id) {
@ -342,11 +367,10 @@ static int _send_msg(struct merkle_sync* ms, uint64_t peer,
static int _broadcast_to_session(struct merkle_sync* ms, struct ms_session* s, static int _broadcast_to_session(struct merkle_sync* ms, struct ms_session* s,
const char* ns, const uint8_t* item_data, size_t item_len) { const char* ns, const uint8_t* item_data, size_t item_len) {
uint8_t ch_len = (uint8_t)strlen(ns); size_t sz = 1 + 8 + 1 + 1 + 1 + 1 + 2 + item_len;
size_t sz = 1 + ch_len + 1 + 1 + 1 + 1 + 2 + item_len;
uint8_t* buf = u_malloc(sz); if (!buf) return -1; uint8_t* buf = u_malloc(sz); if (!buf) return -1;
uint8_t* p = buf; uint8_t* p = buf;
*p++ = ch_len; memcpy(p, ns, ch_len); p += ch_len; ms_write_group_id(p, ns); p += 8;
*p++ = 0x01; /* MSG_HASHES */ *p++ = 0x01; /* MSG_HASHES */
*p++ = 0; /* level=0 */ *p++ = 0; /* level=0 */
*p++ = 0; /* prefix_bytes=0 */ *p++ = 0; /* prefix_bytes=0 */
@ -361,12 +385,11 @@ static int _broadcast_to_session(struct merkle_sync* ms, struct ms_session* s,
static int _send_hashes(struct merkle_sync* ms, uint64_t peer, const char* ns, 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 level, uint64_t prefix, uint8_t prefix_bytes, int is_data) {
DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: send_hashes peer=%016llx ns=%s L%d P%016llx is_data=%d", MS_ID, (unsigned long long)peer, ns, level, (unsigned long long)prefix, is_data); DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: send_hashes peer=%016llx ns=%s L%d P%016llx is_data=%d", MS_ID, (unsigned long long)peer, ns, level, (unsigned long long)prefix, is_data);
uint8_t ch_len = (uint8_t)strlen(ns); size_t max_sz = 1 + 8 + 1 + 1 + prefix_bytes + 1 + 4 + MT_BUCKETS * MT_HASH_SIZE + 2 + 65536;
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); uint8_t* buf = u_malloc(max_sz);
if (!buf) return -1; if (!buf) return -1;
uint8_t* p = buf; uint8_t* p = buf;
*p++ = ch_len; memcpy(p, ns, ch_len); p += ch_len; ms_write_group_id(p, ns); p += 8;
*p++ = 0x01; /* MS_MSG_HASHES */ *p++ = 0x01; /* MS_MSG_HASHES */
*p++ = level; *p++ = level;
*p++ = prefix_bytes; _prefix_write(p, prefix, prefix_bytes); p += prefix_bytes; *p++ = prefix_bytes; _prefix_write(p, prefix, prefix_bytes); p += prefix_bytes;
@ -397,12 +420,11 @@ static int _send_hashes(struct merkle_sync* ms, uint64_t peer, const char* ns,
static int _send_batch(struct merkle_sync* ms, uint64_t peer, const char* ns, static int _send_batch(struct merkle_sync* ms, uint64_t peer, const char* ns,
struct bucket_entry* buckets, int count) { struct bucket_entry* buckets, int count) {
DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: send_batch peer=%016llx ns=%s count=%d", MS_ID, (unsigned long long)peer, ns, count); DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: send_batch peer=%016llx ns=%s count=%d", MS_ID, (unsigned long long)peer, ns, count);
uint8_t ch_len = (uint8_t)strlen(ns); size_t max_sz = 1 + 8 + 1 + (size_t)count * (1 + 1 + 8 + 1 + 4 + MT_BUCKETS * MT_HASH_SIZE + 65536);
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); uint8_t* buf = u_malloc(max_sz);
if (!buf) return -1; if (!buf) return -1;
uint8_t* p = buf; uint8_t* p = buf;
*p++ = ch_len; memcpy(p, ns, ch_len); p += ch_len; ms_write_group_id(p, ns); p += 8;
*p++ = 0x03; /* MS_MSG_BATCH */ *p++ = 0x03; /* MS_MSG_BATCH */
*p++ = (uint8_t)count; *p++ = (uint8_t)count;
@ -556,14 +578,12 @@ static void _handle_hashes(struct merkle_sync* ms, uint64_t peer, const char* ns
if (parent_idx >= 0 && parent_idx < MS_PENDING_MAX) if (parent_idx >= 0 && parent_idx < MS_PENDING_MAX)
s->pending[parent_idx].sub_count++; s->pending[parent_idx].sub_count++;
} }
uint8_t ch_len = (uint8_t)strlen(ns); size_t rs = 1 + 8 + 1;
size_t rs = 1 + 1 + ch_len + 1;
for (int i = 0; i < rcount; i++) rs += 1 + 1 + requests[i].prefix_bytes + 1; for (int i = 0; i < rcount; i++) rs += 1 + 1 + requests[i].prefix_bytes + 1;
uint8_t* rbuf = u_malloc(rs); uint8_t* rbuf = u_malloc(rs);
if (rbuf) { if (rbuf) {
uint8_t* wr = rbuf; uint8_t* wr = rbuf;
*wr++ = ch_len; ms_write_group_id(wr, ns); wr += 8;
memcpy(wr, ns, ch_len); wr += ch_len;
*wr++ = 0x02; /* MS_MSG_REQUEST */ *wr++ = 0x02; /* MS_MSG_REQUEST */
*wr++ = (uint8_t)rcount; *wr++ = (uint8_t)rcount;
for (int i = 0; i < rcount; i++) { for (int i = 0; i < rcount; i++) {
@ -662,7 +682,7 @@ static void _handle_batch(struct merkle_sync* ms, uint64_t peer, const char* ns,
} }
static void _recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { static void _recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
if (!entry || entry->len < 4) { if (!entry || entry->len < 10) {
if (entry) { if (entry->dgram) u_free(entry->dgram); queue_entry_free(entry); } if (entry) { if (entry->dgram) u_free(entry->dgram); queue_entry_free(entry); }
return; return;
} }
@ -673,14 +693,13 @@ static void _recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
const uint8_t* d = entry->dgram; const uint8_t* d = entry->dgram;
size_t dlen = entry->len; size_t dlen = entry->len;
uint8_t ch_len = d[1]; uint64_t group_id = be64toh(*(const uint64_t*)(d + 1));
if (ch_len >= 64) { DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: recv ch_len too large (%u) from=%016llx", MS_ID, ch_len, (unsigned long long)peer); u_free(entry->dgram); queue_entry_free(entry); return; } char ns[64];
if (dlen < (size_t)(2 + ch_len + 1)) { u_free(entry->dgram); queue_entry_free(entry); return; } if (ms_group_id_to_ns(inst, group_id, ns, sizeof(ns)) != 0) { 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[9];
uint8_t type = d[2 + ch_len]; DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: recv type=%02x from=%016llx group=%016llx ns=%s len=%zu", MS_ID, type, (unsigned long long)peer, (unsigned long long)group_id, ns, dlen);
DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: recv type=%02x from=%016llx ns=%s len=%zu", MS_ID, type, (unsigned long long)peer, ns, dlen); const uint8_t* pl = d + 10;
const uint8_t* pl = d + 3 + ch_len; size_t plen = dlen - 10;
size_t plen = dlen - 3 - ch_len;
switch (type) { switch (type) {
case 0x01: _handle_hashes(ms, peer, ns, pl, plen, -1); break; case 0x01: _handle_hashes(ms, peer, ns, pl, plen, -1); break;
@ -709,13 +728,11 @@ void merkle_sync_push_update(struct UTUN_INSTANCE* inst, const char* ns,
DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "%s: push_update ns=%s key=%016llx type=%02x len=%zu", DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "%s: push_update ns=%s key=%016llx type=%02x len=%zu",
MS_ID, ns, (unsigned long long)key, type, len); MS_ID, ns, (unsigned long long)key, type, len);
uint8_t ch_len = (uint8_t)strlen(ns); size_t pkt = 1 + 8 + 1 + 8 + 1 + len;
if (ch_len > 63) return;
size_t pkt = 1 + ch_len + 1 + 8 + 1 + len;
uint8_t* buf = u_malloc(pkt); uint8_t* buf = u_malloc(pkt);
if (!buf) return; if (!buf) return;
uint8_t* p = buf; uint8_t* p = buf;
*p++ = ch_len; memcpy(p, ns, ch_len); p += ch_len; ms_write_group_id(p, ns); p += 8;
*p++ = 0x04; /* MSG_ITEM_UPDATE */ *p++ = 0x04; /* MSG_ITEM_UPDATE */
memcpy(p, &key, 8); p += 8; memcpy(p, &key, 8); p += 8;
*p++ = type; *p++ = type;

137
tests/test_merkle_sync.c

@ -11,6 +11,7 @@
#include "../lib/socket_compat.h" #include "../lib/socket_compat.h"
#include "merkle_sync.h" #include "merkle_sync.h"
#include "../src/utun_instance.h" #include "../src/utun_instance.h"
#include "../src/routing_layer/topo_group.h"
#include "../src/config_updater.h" #include "../src/config_updater.h"
#include "../src/tun_if.h" #include "../src/tun_if.h"
#include "etcp.h" #include "etcp.h"
@ -19,6 +20,21 @@
#include <sqlite3.h> #include <sqlite3.h>
#include <openssl/evp.h> #include <openssl/evp.h>
/* numeric channel_id == group_id для CHAT-групп (recv требует существующую группу) */
#define MS_NS_TEST "1001"
#define MS_NS_RND "1002"
#define MS_NS_STRESS "1003"
static void _ms_ensure_groups(struct UTUN_INSTANCE* inst) {
if (!inst || !inst->topo_groups) return;
const char* nss[] = { MS_NS_TEST, MS_NS_RND, MS_NS_STRESS };
for (int i = 0; i < 3; i++) {
uint64_t gid = strtoull(nss[i], NULL, 10);
if (!topo_groups_find(inst->topo_groups, gid))
topo_groups_create_group(inst->topo_groups, gid, TOPO_GROUP_TYPE_CHAT, nss[i]);
}
}
/* merkle_sync private struct — needed for Stage 2 unit tests. /* merkle_sync private struct — needed for Stage 2 unit tests.
* Must match struct ms_session prefix in merkle_sync.c exactly. */ * Must match struct ms_session prefix in merkle_sync.c exactly. */
struct ms_session { struct ms_session {
@ -320,8 +336,8 @@ static void test_tree_empty_ns(void) {
ms.inst = &inst; ms.ops = &g_test_ops; ms.data_ctx = &ctx; ms.initialized = 1; ms.inst = &inst; ms.ops = &g_test_ops; ms.data_ctx = &ctx; ms.initialized = 1;
inst.msync = &ms; inst.msync = &ms;
merkle_sync_recompute_path(&inst, "test", 0x1234ULL); merkle_sync_recompute_path(&inst, MS_NS_TEST, 0x1234ULL);
if (_db_count_rows(db, "test") == 0) PASS(); else FAIL("got %d rows", _db_count_rows(db, "test")); if (_db_count_rows(db, MS_NS_TEST) == 0) PASS(); else FAIL("got %d rows", _db_count_rows(db, MS_NS_TEST));
sqlite3_close(db); sqlite3_close(db);
} }
@ -342,15 +358,15 @@ static void test_tree_single_item(void) {
ms.inst = &inst; ms.ops = &g_test_ops; ms.data_ctx = &ctx; ms.initialized = 1; ms.inst = &inst; ms.ops = &g_test_ops; ms.data_ctx = &ctx; ms.initialized = 1;
inst.msync = &ms; inst.msync = &ms;
merkle_sync_recompute_path(&inst, "test", key); merkle_sync_recompute_path(&inst, MS_NS_TEST, key);
int rows = _db_count_rows(db, "test"); int rows = _db_count_rows(db, MS_NS_TEST);
if (rows != 5) { FAIL("expected 5 rows got %d", rows); sqlite3_close(db); return; } if (rows != 5) { FAIL("expected 5 rows got %d", rows); sqlite3_close(db); return; }
/* check hashes are non-zero */ /* check hashes are non-zero */
int ok = 1; int ok = 1;
for (uint8_t lv = 1; lv <= 5 && ok; lv++) { for (uint8_t lv = 1; lv <= 5 && ok; lv++) {
uint64_t pf = merkle_sync_level_prefix(key, lv); uint64_t pf = merkle_sync_level_prefix(key, lv);
const uint8_t* h = merkle_sync_get_hash(&inst, "test", lv, pf); const uint8_t* h = merkle_sync_get_hash(&inst, MS_NS_TEST, lv, pf);
int zero = 1; int zero = 1;
for (int i = 0; i < MT_HASH_SIZE; i++) if (h[i] != 0) zero = 0; for (int i = 0; i < MT_HASH_SIZE; i++) if (h[i] != 0) zero = 0;
if (zero) { FAIL("L%d zero hash", lv); ok = 0; } if (zero) { FAIL("L%d zero hash", lv); ok = 0; }
@ -378,8 +394,8 @@ static void test_tree_multi_item_same_bucket(void) {
inst.msync = &ms; inst.msync = &ms;
/* recompute for both */ /* recompute for both */
merkle_sync_recompute_path(&inst, "test", k1); merkle_sync_recompute_path(&inst, MS_NS_TEST, k1);
merkle_sync_recompute_path(&inst, "test", k2); merkle_sync_recompute_path(&inst, MS_NS_TEST, k2);
uint64_t pf1 = merkle_sync_level_prefix(k1, 1); uint64_t pf1 = merkle_sync_level_prefix(k1, 1);
uint64_t pf2 = merkle_sync_level_prefix(k2, 1); uint64_t pf2 = merkle_sync_level_prefix(k2, 1);
@ -395,7 +411,7 @@ static void test_tree_multi_item_same_bucket(void) {
EVP_DigestFinal_ex(ctx, expected_hash, NULL); EVP_DigestFinal_ex(ctx, expected_hash, NULL);
EVP_MD_CTX_free(ctx); EVP_MD_CTX_free(ctx);
} }
const uint8_t* stored = merkle_sync_get_hash(&inst, "test", 1, pf1); const uint8_t* stored = merkle_sync_get_hash(&inst, MS_NS_TEST, 1, pf1);
if (memcmp(expected_hash, stored, MT_HASH_SIZE) == 0) PASS(); else FAIL("hash mismatch"); if (memcmp(expected_hash, stored, MT_HASH_SIZE) == 0) PASS(); else FAIL("hash mismatch");
sqlite3_close(db); sqlite3_close(db);
@ -449,16 +465,16 @@ static void test_tree_delete_empty_bucket(void) {
ms.inst = &inst; ms.ops = &g_test_ops; ms.data_ctx = &ctx; ms.initialized = 1; ms.inst = &inst; ms.ops = &g_test_ops; ms.data_ctx = &ctx; ms.initialized = 1;
inst.msync = &ms; inst.msync = &ms;
merkle_sync_recompute_path(&inst, "test", k); merkle_sync_recompute_path(&inst, MS_NS_TEST, k);
if (_db_count_rows(db, "test") != 5) { FAIL("expected 5 rows"); sqlite3_close(db); return; } if (_db_count_rows(db, MS_NS_TEST) != 5) { FAIL("expected 5 rows"); sqlite3_close(db); return; }
/* remove item and recompute */ /* remove item and recompute */
data.count = 0; data.count = 0;
merkle_sync_recompute_path(&inst, "test", k); merkle_sync_recompute_path(&inst, MS_NS_TEST, k);
if (_db_count_rows(db, "test") != 0) { FAIL("expected 0 rows got %d", _db_count_rows(db, "test")); sqlite3_close(db); return; } if (_db_count_rows(db, MS_NS_TEST) != 0) { FAIL("expected 0 rows got %d", _db_count_rows(db, MS_NS_TEST)); sqlite3_close(db); return; }
uint64_t pf = merkle_sync_level_prefix(k, 3); uint64_t pf = merkle_sync_level_prefix(k, 3);
const uint8_t* h = merkle_sync_get_hash(&inst, "test", 3, pf); const uint8_t* h = merkle_sync_get_hash(&inst, MS_NS_TEST, 3, pf);
int zero = 1; int zero = 1;
for (int i = 0; i < MT_HASH_SIZE; i++) if (h[i] != 0) zero = 0; for (int i = 0; i < MT_HASH_SIZE; i++) if (h[i] != 0) zero = 0;
if (zero) PASS(); else FAIL("hash not zero after delete"); if (zero) PASS(); else FAIL("hash not zero after delete");
@ -694,6 +710,7 @@ static int _intg_init_two(void) {
_intg_setup_contexts(); _intg_setup_contexts();
ctx_a.inst = i_a; ctx_b.inst = i_b; ctx_a.inst = i_a; ctx_b.inst = i_b;
_ms_ensure_groups(i_a); _ms_ensure_groups(i_b);
if (merkle_sync_init(i_a, MS_SVC_ID, &g_test_ops, &ctx_a) != 0 || if (merkle_sync_init(i_a, MS_SVC_ID, &g_test_ops, &ctx_a) != 0 ||
merkle_sync_init(i_b, MS_SVC_ID, &g_test_ops, &ctx_b) != 0) merkle_sync_init(i_b, MS_SVC_ID, &g_test_ops, &ctx_b) != 0)
{ _intg_cleanup(); return -1; } { _intg_cleanup(); return -1; }
@ -709,7 +726,7 @@ static void _intg_ins_many(struct ms_data* d, struct UTUN_INSTANCE* inst, int ba
} }
_data_sort(d); _data_sort(d);
for (int i = 0; i < d->count; i++) for (int i = 0; i < d->count; i++)
merkle_sync_recompute_path(inst, "test", d->items[i].key); merkle_sync_recompute_path(inst, MS_NS_TEST, d->items[i].key);
} }
static int _intg_compare_data(struct ms_data* a, struct ms_data* b) { static int _intg_compare_data(struct ms_data* a, struct ms_data* b) {
@ -736,14 +753,14 @@ static void test_peer_empty(void) {
if (_intg_init_two() != 0) { FAIL("setup failed"); return; } if (_intg_init_two() != 0) { FAIL("setup failed"); return; }
_intg_ins_many(&data_b, i_b, 0, 5); _intg_ins_many(&data_b, i_b, 0, 5);
done_sync = 0; done_sync = 0;
merkle_sync_start(i_a, i_b->node_id, "test", _on_sync_done, NULL); merkle_sync_start(i_a, i_b->node_id, MS_NS_TEST, _on_sync_done, NULL);
if (!_wait_for("sync done", _cond_done, PHASE_TIMEOUT_TB)) if (!_wait_for("sync done", _cond_done, PHASE_TIMEOUT_TB))
{ FAIL("sync timeout"); _intg_cleanup(); return; } { FAIL("sync timeout"); _intg_cleanup(); return; }
printf(" A data=%d B data=%d, A tree=%d B tree=%d\n", printf(" A data=%d B data=%d, A tree=%d B tree=%d\n",
data_a.count, data_b.count, data_a.count, data_b.count,
_db_count_rows(i_a->topo_sqlite_db, "test"), _db_count_rows(i_a->topo_sqlite_db, MS_NS_TEST),
_db_count_rows(i_b->topo_sqlite_db, "test")); _db_count_rows(i_b->topo_sqlite_db, MS_NS_TEST));
if (_intg_verify_all(i_a, i_b, &data_a, &data_b, "test") != 0) if (_intg_verify_all(i_a, i_b, &data_a, &data_b, MS_NS_TEST) != 0)
{ FAIL("verify failed"); _intg_cleanup(); return; } { FAIL("verify failed"); _intg_cleanup(); return; }
PASS(); PASS();
_intg_cleanup(); _intg_cleanup();
@ -755,10 +772,10 @@ static void test_hash_match(void) {
_intg_ins_many(&data_a, i_a, 0, 5); _intg_ins_many(&data_a, i_a, 0, 5);
_intg_ins_many(&data_b, i_b, 0, 5); _intg_ins_many(&data_b, i_b, 0, 5);
done_sync = 0; done_sync = 0;
merkle_sync_start(i_a, i_b->node_id, "test", _on_sync_done, NULL); merkle_sync_start(i_a, i_b->node_id, MS_NS_TEST, _on_sync_done, NULL);
for (int i = 0; i < 200 && !done_sync && i_phase == 0; i++) uasync_poll(i_ua, 1); for (int i = 0; i < 200 && !done_sync && i_phase == 0; i++) uasync_poll(i_ua, 1);
if (!done_sync) { FAIL("sync did not complete quickly"); _intg_cleanup(); return; } if (!done_sync) { FAIL("sync did not complete quickly"); _intg_cleanup(); return; }
if (_intg_verify_all(i_a, i_b, &data_a, &data_b, "test") != 0) if (_intg_verify_all(i_a, i_b, &data_a, &data_b, MS_NS_TEST) != 0)
{ FAIL("verify failed"); _intg_cleanup(); return; } { FAIL("verify failed"); _intg_cleanup(); return; }
PASS(); PASS();
_intg_cleanup(); _intg_cleanup();
@ -770,10 +787,10 @@ static void test_divergence_merge(void) {
_intg_ins_many(&data_a, i_a, 0, 3); _intg_ins_many(&data_a, i_a, 0, 3);
_intg_ins_many(&data_b, i_b, 10, 2); _intg_ins_many(&data_b, i_b, 10, 2);
done_sync = 0; done_sync = 0;
merkle_sync_start(i_a, i_b->node_id, "test", _on_sync_done, NULL); merkle_sync_start(i_a, i_b->node_id, MS_NS_TEST, _on_sync_done, NULL);
if (!_wait_for("sync done", _cond_done, PHASE_TIMEOUT_TB)) if (!_wait_for("sync done", _cond_done, PHASE_TIMEOUT_TB))
{ FAIL("sync timeout"); _intg_cleanup(); return; } { FAIL("sync timeout"); _intg_cleanup(); return; }
if (_intg_verify_all(i_a, i_b, &data_a, &data_b, "test") != 0) if (_intg_verify_all(i_a, i_b, &data_a, &data_b, MS_NS_TEST) != 0)
{ FAIL("verify failed"); _intg_cleanup(); return; } { FAIL("verify failed"); _intg_cleanup(); return; }
PASS(); PASS();
_intg_cleanup(); _intg_cleanup();
@ -806,17 +823,17 @@ static void test_randomized_two(void) {
_data_insert(&data_b, k, (uint32_t)rand()); _data_insert(&data_b, k, (uint32_t)rand());
} }
_data_sort(&data_a); _data_sort(&data_b); _data_sort(&data_a); _data_sort(&data_b);
for (int i = 0; i < data_a.count; i++) merkle_sync_recompute_path(i_a, "rnd", data_a.items[i].key); for (int i = 0; i < data_a.count; i++) merkle_sync_recompute_path(i_a, MS_NS_RND, data_a.items[i].key);
for (int i = 0; i < data_b.count; i++) merkle_sync_recompute_path(i_b, "rnd", data_b.items[i].key); for (int i = 0; i < data_b.count; i++) merkle_sync_recompute_path(i_b, MS_NS_RND, data_b.items[i].key);
done_sync = 0; done_sync = 0;
merkle_sync_start(i_a, i_b->node_id, "rnd", _on_sync_done, NULL); merkle_sync_start(i_a, i_b->node_id, MS_NS_RND, _on_sync_done, NULL);
if (!_wait_for("sync", _cond_done, PHASE_TIMEOUT_TB)) { FAIL("timeout iter %d", iter); _intg_cleanup(); break; } if (!_wait_for("sync", _cond_done, PHASE_TIMEOUT_TB)) { FAIL("timeout iter %d", iter); _intg_cleanup(); break; }
if (_intg_verify_all(i_a, i_b, &data_a, &data_b, "rnd") != 0) { if (_intg_verify_all(i_a, i_b, &data_a, &data_b, MS_NS_RND) != 0) {
printf(" A data=%d B data=%d A tree=%d B tree=%d\n", printf(" A data=%d B data=%d A tree=%d B tree=%d\n",
data_a.count, data_b.count, data_a.count, data_b.count,
_db_count_rows(i_a->topo_sqlite_db, "rnd"), _db_count_rows(i_a->topo_sqlite_db, MS_NS_RND),
_db_count_rows(i_b->topo_sqlite_db, "rnd")); _db_count_rows(i_b->topo_sqlite_db, MS_NS_RND));
/* dump items only in A */ /* dump items only in A */
for (int ai = 0; ai < data_a.count; ai++) { for (int ai = 0; ai < data_a.count; ai++) {
int found = 0; int found = 0;
@ -839,16 +856,16 @@ static void test_randomized_two(void) {
sqlite3_prepare_v2(i_a->topo_sqlite_db, sqlite3_prepare_v2(i_a->topo_sqlite_db,
"SELECT level,prefix64 FROM merkle_tree_hash WHERE namespace=? ORDER BY level,prefix64", "SELECT level,prefix64 FROM merkle_tree_hash WHERE namespace=? ORDER BY level,prefix64",
-1, &st, NULL); -1, &st, NULL);
sqlite3_bind_text(st, 1, "rnd", -1, SQLITE_STATIC); sqlite3_bind_text(st, 1, MS_NS_RND, -1, SQLITE_STATIC);
while (sqlite3_step(st) == SQLITE_ROW && ok) { while (sqlite3_step(st) == SQLITE_ROW && ok) {
uint8_t lv = (uint8_t)sqlite3_column_int(st, 0); uint8_t lv = (uint8_t)sqlite3_column_int(st, 0);
uint64_t pf = (uint64_t)sqlite3_column_int64(st, 1); uint64_t pf = (uint64_t)sqlite3_column_int64(st, 1);
EVP_MD_CTX* c = EVP_MD_CTX_new(); EVP_MD_CTX* c = EVP_MD_CTX_new();
EVP_DigestInit_ex(c, EVP_sha256(), NULL); EVP_DigestInit_ex(c, EVP_sha256(), NULL);
_bucket_hash(&ctx_a, "rnd", lv, pf, c); _bucket_hash(&ctx_a, MS_NS_RND, lv, pf, c);
uint8_t expected[MT_HASH_SIZE]; uint8_t expected[MT_HASH_SIZE];
EVP_DigestFinal_ex(c, expected, NULL); EVP_MD_CTX_free(c); EVP_DigestFinal_ex(c, expected, NULL); EVP_MD_CTX_free(c);
const uint8_t* stored = merkle_sync_get_hash(i_a, "rnd", lv, pf); const uint8_t* stored = merkle_sync_get_hash(i_a, MS_NS_RND, lv, pf);
uint8_t scopy[MT_HASH_SIZE]; memcpy(scopy, stored, MT_HASH_SIZE); uint8_t scopy[MT_HASH_SIZE]; memcpy(scopy, stored, MT_HASH_SIZE);
if (memcmp(expected, scopy, MT_HASH_SIZE) != 0) { ok = 0; } if (memcmp(expected, scopy, MT_HASH_SIZE) != 0) { ok = 0; }
} }
@ -879,6 +896,7 @@ static int _intg_init_three(void) {
for (int i = 0; i < 20; i++) uasync_poll(i_ua, 5); for (int i = 0; i < 20; i++) uasync_poll(i_ua, 5);
_intg_setup_contexts(); _intg_setup_contexts();
ctx_a.inst = i_a; ctx_b.inst = i_b; ctx_c.inst = i_c; ctx_a.inst = i_a; ctx_b.inst = i_b; ctx_c.inst = i_c;
_ms_ensure_groups(i_a); _ms_ensure_groups(i_b); _ms_ensure_groups(i_c);
if (merkle_sync_init(i_a, MS_SVC_ID, &g_test_ops, &ctx_a) != 0 || if (merkle_sync_init(i_a, MS_SVC_ID, &g_test_ops, &ctx_a) != 0 ||
merkle_sync_init(i_b, MS_SVC_ID, &g_test_ops, &ctx_b) != 0 || merkle_sync_init(i_b, MS_SVC_ID, &g_test_ops, &ctx_b) != 0 ||
merkle_sync_init(i_c, MS_SVC_ID, &g_test_ops, &ctx_c) != 0) merkle_sync_init(i_c, MS_SVC_ID, &g_test_ops, &ctx_c) != 0)
@ -890,8 +908,8 @@ static int _intg_init_three(void) {
static int _cond_all_synced(void) { static int _cond_all_synced(void) {
if (data_a.count != data_b.count || data_b.count != data_c.count) return 0; if (data_a.count != data_b.count || data_b.count != data_c.count) return 0;
if (_db_compare_trees(i_a->topo_sqlite_db, i_b->topo_sqlite_db, "rnd") != 0) return 0; if (_db_compare_trees(i_a->topo_sqlite_db, i_b->topo_sqlite_db, MS_NS_RND) != 0) return 0;
if (_db_compare_trees(i_b->topo_sqlite_db, i_c->topo_sqlite_db, "rnd") != 0) return 0; if (_db_compare_trees(i_b->topo_sqlite_db, i_c->topo_sqlite_db, MS_NS_RND) != 0) return 0;
return 1; return 1;
} }
@ -905,24 +923,24 @@ static void test_randomized_three(void) {
for (int i = 0; i < na; i++) { uint64_t k = ((uint64_t)rand()<<32)|(uint64_t)rand(); _data_insert(&data_a, k, (uint32_t)rand()); } for (int i = 0; i < na; i++) { uint64_t k = ((uint64_t)rand()<<32)|(uint64_t)rand(); _data_insert(&data_a, k, (uint32_t)rand()); }
for (int i = 0; i < nb; i++) { uint64_t k = ((uint64_t)rand()<<32)|(uint64_t)rand(); _data_insert(&data_b, k, (uint32_t)rand()); } for (int i = 0; i < nb; i++) { uint64_t k = ((uint64_t)rand()<<32)|(uint64_t)rand(); _data_insert(&data_b, k, (uint32_t)rand()); }
for (int i = 0; i < nc; i++) { uint64_t k = ((uint64_t)rand()<<32)|(uint64_t)rand(); _data_insert(&data_c, k, (uint32_t)rand()); } for (int i = 0; i < nc; i++) { uint64_t k = ((uint64_t)rand()<<32)|(uint64_t)rand(); _data_insert(&data_c, k, (uint32_t)rand()); }
_data_sort(&data_a); for (int i = 0; i < data_a.count; i++) merkle_sync_recompute_path(i_a, "rnd", data_a.items[i].key); _data_sort(&data_a); for (int i = 0; i < data_a.count; i++) merkle_sync_recompute_path(i_a, MS_NS_RND, data_a.items[i].key);
_data_sort(&data_b); for (int i = 0; i < data_b.count; i++) merkle_sync_recompute_path(i_b, "rnd", data_b.items[i].key); _data_sort(&data_b); for (int i = 0; i < data_b.count; i++) merkle_sync_recompute_path(i_b, MS_NS_RND, data_b.items[i].key);
_data_sort(&data_c); for (int i = 0; i < data_c.count; i++) merkle_sync_recompute_path(i_c, "rnd", data_c.items[i].key); _data_sort(&data_c); for (int i = 0; i < data_c.count; i++) merkle_sync_recompute_path(i_c, MS_NS_RND, data_c.items[i].key);
done_sync = 0; i_phase = 0; done_sync = 0; i_phase = 0;
/* B→A: B syncs with A, gets A's data, B: SYNCED */ /* B→A: B syncs with A, gets A's data, B: SYNCED */
merkle_sync_start(i_b, i_a->node_id, "rnd", _on_sync_done, NULL); merkle_sync_start(i_b, i_a->node_id, MS_NS_RND, _on_sync_done, NULL);
if (!_wait_for("sync B", _cond_done, PHASE_TIMEOUT_TB)) { FAIL("timeout B iter %d", iter); _intg_cleanup(); break; } if (!_wait_for("sync B", _cond_done, PHASE_TIMEOUT_TB)) { FAIL("timeout B iter %d", iter); _intg_cleanup(); break; }
done_sync = 0; i_phase = 0; done_sync = 0; i_phase = 0;
/* C→A: C syncs with A, gets A∪B, A broadcasts C's new items to B (relay) */ /* C→A: C syncs with A, gets A∪B, A broadcasts C's new items to B (relay) */
merkle_sync_start(i_c, i_a->node_id, "rnd", _on_sync_done, NULL); merkle_sync_start(i_c, i_a->node_id, MS_NS_RND, _on_sync_done, NULL);
if (!_wait_for("sync C", _cond_done, PHASE_TIMEOUT_TB)) { FAIL("timeout C iter %d", iter); _intg_cleanup(); break; } if (!_wait_for("sync C", _cond_done, PHASE_TIMEOUT_TB)) { FAIL("timeout C iter %d", iter); _intg_cleanup(); break; }
/* wait for broadcast relay */ /* wait for broadcast relay */
for (int k = 0; k < 100; k++) uasync_poll(i_ua, 10); for (int k = 0; k < 100; k++) uasync_poll(i_ua, 10);
if (_intg_compare_data(&data_a, &data_b) != 0 || _intg_compare_data(&data_b, &data_c) != 0) if (_intg_compare_data(&data_a, &data_b) != 0 || _intg_compare_data(&data_b, &data_c) != 0)
{ FAIL("data mismatch iter %d", iter); _intg_cleanup(); break; } { FAIL("data mismatch iter %d", iter); _intg_cleanup(); break; }
if (_db_compare_trees(i_a->topo_sqlite_db, i_b->topo_sqlite_db, "rnd") != 0 || if (_db_compare_trees(i_a->topo_sqlite_db, i_b->topo_sqlite_db, MS_NS_RND) != 0 ||
_db_compare_trees(i_b->topo_sqlite_db, i_c->topo_sqlite_db, "rnd") != 0) _db_compare_trees(i_b->topo_sqlite_db, i_c->topo_sqlite_db, MS_NS_RND) != 0)
{ FAIL("tree mismatch iter %d", iter); _intg_cleanup(); break; } { FAIL("tree mismatch iter %d", iter); _intg_cleanup(); break; }
_intg_cleanup(); _intg_cleanup();
} }
@ -936,9 +954,9 @@ static void test_cancel(void) {
if (_intg_init_two() != 0) { FAIL("setup failed"); return; } if (_intg_init_two() != 0) { FAIL("setup failed"); return; }
_intg_ins_many(&data_b, i_b, 0, 5); _intg_ins_many(&data_b, i_b, 0, 5);
done_sync = 0; done_sync = 0;
merkle_sync_start(i_a, i_b->node_id, "test", _on_sync_done, NULL); merkle_sync_start(i_a, i_b->node_id, MS_NS_TEST, _on_sync_done, NULL);
uasync_poll(i_ua, 1); uasync_poll(i_ua, 1);
merkle_sync_cancel(i_a, i_b->node_id, "test"); merkle_sync_cancel(i_a, i_b->node_id, MS_NS_TEST);
for (int i = 0; i < 100 && i_phase == 0; i++) uasync_poll(i_ua, 10); for (int i = 0; i < 100 && i_phase == 0; i++) uasync_poll(i_ua, 10);
if (!done_sync) PASS(); else FAIL("done_cb was called"); if (!done_sync) PASS(); else FAIL("done_cb was called");
_intg_cleanup(); _intg_cleanup();
@ -953,8 +971,8 @@ static void test_start_overwrite(void) {
if (_intg_init_two() != 0) { FAIL("setup failed"); return; } if (_intg_init_two() != 0) { FAIL("setup failed"); return; }
_intg_ins_many(&data_b, i_b, 0, 3); _intg_ins_many(&data_b, i_b, 0, 3);
sw_cb1_called = sw_cb2_called = 0; sw_cb1_called = sw_cb2_called = 0;
merkle_sync_start(i_a, i_b->node_id, "test", _sw_done1, NULL); merkle_sync_start(i_a, i_b->node_id, MS_NS_TEST, _sw_done1, NULL);
merkle_sync_start(i_a, i_b->node_id, "test", _sw_done2, NULL); merkle_sync_start(i_a, i_b->node_id, MS_NS_TEST, _sw_done2, NULL);
for (int i = 0; i < 500 && !sw_cb2_called && i_phase == 0; i++) uasync_poll(i_ua, 10); for (int i = 0; i < 500 && !sw_cb2_called && i_phase == 0; i++) uasync_poll(i_ua, 10);
if (sw_cb2_called) PASS(); else FAIL("cb1=%d cb2=%d", sw_cb1_called, sw_cb2_called); if (sw_cb2_called) PASS(); else FAIL("cb1=%d cb2=%d", sw_cb1_called, sw_cb2_called);
_intg_cleanup(); _intg_cleanup();
@ -1043,8 +1061,8 @@ static void _str_spam_cb(void* arg) {
/* spammer: bump own member, sync to hub */ /* spammer: bump own member, sync to hub */
uint64_t key = ((uint64_t)idx) << 60; uint64_t key = ((uint64_t)idx) << 60;
_data_insert(&gs->data[idx], key, (uint32_t)(gs->data[idx].count + 1)); _data_insert(&gs->data[idx], key, (uint32_t)(gs->data[idx].count + 1));
merkle_sync_recompute_path(gs->inst[idx], "stress", key); merkle_sync_recompute_path(gs->inst[idx], MS_NS_STRESS, key);
merkle_sync_start(gs->inst[idx], gs->inst[STRESS_HUB]->node_id, "stress", NULL, NULL); merkle_sync_start(gs->inst[idx], gs->inst[STRESS_HUB]->node_id, MS_NS_STRESS, NULL, NULL);
_str_spam_schedule(gs, idx); _str_spam_schedule(gs, idx);
} }
@ -1053,7 +1071,7 @@ static void _str_obs_spam_cb(void* arg) {
if (!gs || !gs->spam_active) return; if (!gs || !gs->spam_active) return;
int idx = (int)(intptr_t)arg; int idx = (int)(intptr_t)arg;
merkle_sync_start(gs->inst[idx], gs->inst[STRESS_HUB]->node_id, "stress", NULL, NULL); merkle_sync_start(gs->inst[idx], gs->inst[STRESS_HUB]->node_id, MS_NS_STRESS, NULL, NULL);
int delay_tb = ((rand() % 10) + 1) * 10; /* 1-10ms */ int delay_tb = ((rand() % 10) + 1) * 10; /* 1-10ms */
uasync_set_timeout(gs->ua, (uint32_t)delay_tb, arg, _str_obs_spam_cb, "obs_spam"); uasync_set_timeout(gs->ua, (uint32_t)delay_tb, arg, _str_obs_spam_cb, "obs_spam");
@ -1160,6 +1178,7 @@ static int _str_init(void) {
s->ctx[i].data = &s->data[i]; s->ctx[i].data = &s->data[i];
s->ctx[i].inst = s->inst[i]; s->ctx[i].inst = s->inst[i];
s->ctx[i].tag = (char)('A' + i); s->ctx[i].tag = (char)('A' + i);
_ms_ensure_groups(s->inst[i]);
if (merkle_sync_init(s->inst[i], MS_SVC_ID, &g_test_ops, &s->ctx[i]) != 0) { if (merkle_sync_init(s->inst[i], MS_SVC_ID, &g_test_ops, &s->ctx[i]) != 0) {
printf(" FAIL: merkle_sync_init %d\n", i); goto fail; printf(" FAIL: merkle_sync_init %d\n", i); goto fail;
} }
@ -1170,12 +1189,12 @@ static int _str_init(void) {
for (int i = 1; i < STRESS_N; i++) { for (int i = 1; i < STRESS_N; i++) {
uint64_t key = ((uint64_t)i) << 60; uint64_t key = ((uint64_t)i) << 60;
_data_insert(&s->data[i], key, (uint32_t)i); _data_insert(&s->data[i], key, (uint32_t)i);
merkle_sync_recompute_path(s->inst[i], "stress", key); merkle_sync_recompute_path(s->inst[i], MS_NS_STRESS, key);
} }
/* initial sync: each spoke → hub */ /* initial sync: each spoke → hub */
s->sync_pending = STRESS_N - 1; s->sync_pending = STRESS_N - 1;
for (int i = 1; i < STRESS_N; i++) for (int i = 1; i < STRESS_N; i++)
merkle_sync_start(s->inst[i], s->inst[STRESS_HUB]->node_id, "stress", _str_sync_done, s); merkle_sync_start(s->inst[i], s->inst[STRESS_HUB]->node_id, MS_NS_STRESS, _str_sync_done, s);
_str_wait_for("initial sync", _str_cond_sync_done, STRESS_SYNC_TB); _str_wait_for("initial sync", _str_cond_sync_done, STRESS_SYNC_TB);
/* observers: sync with hub too */ /* observers: sync with hub too */
@ -1236,7 +1255,7 @@ static void _str_start_final_round(void) {
s->final_pending = STRESS_N - 1; s->final_pending = STRESS_N - 1;
printf(" final round %d: syncing %d spokes -> hub\n", s->final_round, s->final_pending); printf(" final round %d: syncing %d spokes -> hub\n", s->final_round, s->final_pending);
for (int i = 1; i < STRESS_N; i++) for (int i = 1; i < STRESS_N; i++)
merkle_sync_start(s->inst[i], s->inst[STRESS_HUB]->node_id, "stress", _str_final_done, NULL); merkle_sync_start(s->inst[i], s->inst[STRESS_HUB]->node_id, MS_NS_STRESS, _str_final_done, NULL);
} }
static void _str_final_done(uint64_t peer, const char* ns, int result, void* arg) { static void _str_final_done(uint64_t peer, const char* ns, int result, void* arg) {
@ -1251,17 +1270,17 @@ static void _str_verify(void) {
int ok = 1; int ok = 1;
/* tree sizes match pairwise */ /* tree sizes match pairwise */
int na = _db_count_rows(s->inst[STRESS_HUB]->topo_sqlite_db, "stress"); int na = _db_count_rows(s->inst[STRESS_HUB]->topo_sqlite_db, MS_NS_STRESS);
for (int i = 1; i < STRESS_N && ok; i++) { for (int i = 1; i < STRESS_N && ok; i++) {
int nb = _db_count_rows(s->inst[i]->topo_sqlite_db, "stress"); int nb = _db_count_rows(s->inst[i]->topo_sqlite_db, MS_NS_STRESS);
if (na != nb) { printf(" tree size mismatch: hub=%d inst[%d]=%d\n", na, i, nb); ok = 0; } if (na != nb) { printf(" tree size mismatch: hub=%d inst[%d]=%d\n", na, i, nb); ok = 0; }
} }
/* root (level=1, prefix=0) hashes identical */ /* root (level=1, prefix=0) hashes identical */
const uint8_t* root0 = merkle_sync_get_hash(s->inst[STRESS_HUB], "stress", 1, 0); const uint8_t* root0 = merkle_sync_get_hash(s->inst[STRESS_HUB], MS_NS_STRESS, 1, 0);
uint8_t root_copy[MT_HASH_SIZE]; memcpy(root_copy, root0, MT_HASH_SIZE); uint8_t root_copy[MT_HASH_SIZE]; memcpy(root_copy, root0, MT_HASH_SIZE);
for (int i = 1; i < STRESS_N && ok; i++) { for (int i = 1; i < STRESS_N && ok; i++) {
const uint8_t* ri = merkle_sync_get_hash(s->inst[i], "stress", 1, 0); const uint8_t* ri = merkle_sync_get_hash(s->inst[i], MS_NS_STRESS, 1, 0);
if (memcmp(root_copy, ri, MT_HASH_SIZE) != 0) { printf(" ROOT HASH MISMATCH inst[%d]\n", i); ok = 0; } if (memcmp(root_copy, ri, MT_HASH_SIZE) != 0) { printf(" ROOT HASH MISMATCH inst[%d]\n", i); ok = 0; }
} }
@ -1271,16 +1290,16 @@ static void _str_verify(void) {
sqlite3_prepare_v2(s->inst[i]->topo_sqlite_db, sqlite3_prepare_v2(s->inst[i]->topo_sqlite_db,
"SELECT level,prefix64 FROM merkle_tree_hash WHERE namespace=? ORDER BY level,prefix64", "SELECT level,prefix64 FROM merkle_tree_hash WHERE namespace=? ORDER BY level,prefix64",
-1, &st, NULL); -1, &st, NULL);
sqlite3_bind_text(st, 1, "stress", -1, SQLITE_STATIC); sqlite3_bind_text(st, 1, MS_NS_STRESS, -1, SQLITE_STATIC);
while (sqlite3_step(st) == SQLITE_ROW && ok) { while (sqlite3_step(st) == SQLITE_ROW && ok) {
uint8_t lv = (uint8_t)sqlite3_column_int(st, 0); uint8_t lv = (uint8_t)sqlite3_column_int(st, 0);
uint64_t pf = (uint64_t)sqlite3_column_int64(st, 1); uint64_t pf = (uint64_t)sqlite3_column_int64(st, 1);
EVP_MD_CTX* c = EVP_MD_CTX_new(); EVP_MD_CTX* c = EVP_MD_CTX_new();
EVP_DigestInit_ex(c, EVP_sha256(), NULL); EVP_DigestInit_ex(c, EVP_sha256(), NULL);
_bucket_hash(&s->ctx[i], "stress", lv, pf, c); _bucket_hash(&s->ctx[i], MS_NS_STRESS, lv, pf, c);
uint8_t expected[MT_HASH_SIZE]; uint8_t expected[MT_HASH_SIZE];
EVP_DigestFinal_ex(c, expected, NULL); EVP_MD_CTX_free(c); EVP_DigestFinal_ex(c, expected, NULL); EVP_MD_CTX_free(c);
const uint8_t* stored = merkle_sync_get_hash(s->inst[i], "stress", lv, pf); const uint8_t* stored = merkle_sync_get_hash(s->inst[i], MS_NS_STRESS, lv, pf);
uint8_t scopy[MT_HASH_SIZE]; memcpy(scopy, stored, MT_HASH_SIZE); uint8_t scopy[MT_HASH_SIZE]; memcpy(scopy, stored, MT_HASH_SIZE);
if (memcmp(expected, scopy, MT_HASH_SIZE) != 0) { if (memcmp(expected, scopy, MT_HASH_SIZE) != 0) {
printf(" consistency fail inst[%d] L%d P%016llx\n", i, lv, (unsigned long long)pf); printf(" consistency fail inst[%d] L%d P%016llx\n", i, lv, (unsigned long long)pf);
@ -1301,12 +1320,12 @@ static void _str_verify(void) {
static void _str_check_roots(void) { static void _str_check_roots(void) {
struct str_state* s = gs; struct str_state* s = gs;
if (!s) return; if (!s) return;
const uint8_t* root0 = merkle_sync_get_hash(s->inst[STRESS_HUB], "stress", 1, 0); const uint8_t* root0 = merkle_sync_get_hash(s->inst[STRESS_HUB], MS_NS_STRESS, 1, 0);
uint8_t root_copy[MT_HASH_SIZE]; uint8_t root_copy[MT_HASH_SIZE];
memcpy(root_copy, root0, MT_HASH_SIZE); memcpy(root_copy, root0, MT_HASH_SIZE);
int converged = 1; int converged = 1;
for (int i = 1; i < STRESS_N; i++) { for (int i = 1; i < STRESS_N; i++) {
const uint8_t* ri = merkle_sync_get_hash(s->inst[i], "stress", 1, 0); const uint8_t* ri = merkle_sync_get_hash(s->inst[i], MS_NS_STRESS, 1, 0);
if (memcmp(root_copy, ri, MT_HASH_SIZE) != 0) { converged = 0; break; } if (memcmp(root_copy, ri, MT_HASH_SIZE) != 0) { converged = 0; break; }
} }
if (converged) { if (converged) {

Loading…
Cancel
Save