From 7b73d99b5bed9ef1e85f947abfbddab1528ec333 Mon Sep 17 00:00:00 2001 From: evgeny Date: Fri, 4 Sep 2026 09:02:46 +0300 Subject: [PATCH] =?UTF-8?q?chat:=20=D0=BF=D0=B5=D1=80=D0=B5=D0=B4=D0=B0?= =?UTF-8?q?=D0=B2=D0=B0=D1=82=D1=8C=20group=5Fid=20=D0=B2=20=D0=B1=D0=B8?= =?UTF-8?q?=D0=BD=D0=B0=D1=80=D0=BD=D0=BE=D0=BC=20=D0=B2=D0=B8=D0=B4=D0=B5?= =?UTF-8?q?=20=D0=B2=20=D0=BF=D1=80=D0=BE=D1=82=D0=BE=D0=BA=D0=BE=D0=BB?= =?UTF-8?q?=D0=B5=20=D1=81=D0=B8=D0=BD=D1=85=D1=80=D0=BE=D0=BD=D0=B8=D0=B7?= =?UTF-8?q?=D0=B0=D1=86=D0=B8=D0=B8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 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. --- src/chat/chat_channel.c | 10 ++- src/chat/chat_sync.c | 40 ++++++++---- src/chat/db_sync.c | 80 +++++++++++------------ src/chat/db_sync.h | 5 +- src/chat/merkle_sync.c | 69 ++++++++++++-------- tests/test_merkle_sync.c | 137 ++++++++++++++++++++++----------------- 6 files changed, 194 insertions(+), 147 deletions(-) diff --git a/src/chat/chat_channel.c b/src/chat/chat_channel.c index 2bff382b..ec6a2cb7 100644 --- a/src/chat/chat_channel.c +++ b/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)); - 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); + 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)) { 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); @@ -70,8 +68,8 @@ void chat_core_ensure_channel_ready(const char* ch_id) { if (si) { si_register(si, 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", - CC_ID, ch_id, tbl_msg, (unsigned long long)ch_hash); + 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)gid); } else { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: db_sync_instance_add failed for ch=%s", CC_ID, ch_id); diff --git a/src/chat/chat_sync.c b/src/chat/chat_sync.c index 4cfa279a..9e756516 100644 --- a/src/chat/chat_sync.c +++ b/src/chat/chat_sync.c @@ -102,18 +102,33 @@ static struct ETCP_CONN* cs_find_conn_for_node(struct UTUN_INSTANCE* inst, uint6 /* ── 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, 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; } - uint8_t ch_len = (uint8_t)strlen(ch_id); - 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; + size_t total = 1 + 8 + len; 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; } uint8_t* p = buf; *p++ = ETCP_RT_ID_CHAT_SYNC; - *p++ = ch_len; - memcpy(p, ch_id, ch_len); p += ch_len; + cs_write_group_id(p, ch_id); p += 8; memcpy(p, payload, len); 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; } @@ -160,7 +175,7 @@ static int cs_send_channel_invite(struct chat_sync* cs, struct UTUN_INSTANCE* in /* ── Recv dispatcher ── */ 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) { 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); @@ -173,15 +188,14 @@ static void chat_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { 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 ch_id[64]; memcpy(ch_id, d + 2, ch_len); ch_id[ch_len] = '\0'; - uint8_t type = d[2 + ch_len]; + uint64_t group_id = be64toh(*(const uint64_t*)(d + 1)); + 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; } + uint8_t type = d[9]; 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); - const uint8_t* pl = d + 3 + ch_len; - size_t plen = dlen - 3 - ch_len; + const uint8_t* pl = d + 10; + size_t plen = dlen - 10; 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); diff --git a/src/chat/db_sync.c b/src/chat/db_sync.c index f942335f..3ae8ff89 100644 --- a/src/chat/db_sync.c +++ b/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_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_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 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); @@ -54,7 +54,7 @@ struct SI_PEER { struct DB_SYNC_INSTANCE { struct DB_SYNC* db_sync; - uint64_t hash; + uint64_t group_id; char table_name[64]; uint64_t next_id; 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_SHRT(si) ({ \ 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; \ }) @@ -220,11 +220,11 @@ static void db_chain_hash_compute(const uint8_t prev_chain_hash[32], // Instance management // ============================================================ -// Ищет активный инстанс синхронизации по 64-битному хешу -static struct DB_SYNC_INSTANCE* db_instance_find(struct DB_SYNC* db, uint64_t hash) +// Ищет активный инстанс синхронизации по group_id канала +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++) { - 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; } @@ -672,24 +672,24 @@ static int db_record_insert_cascade(struct DB_SYNC_INSTANCE* si, // Send // ============================================================ -// Отправляет сообщение узлу через ETCP с заголовком: [ETCP_RT_ID_DB_SYNC:1][hash_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) +// Отправляет сообщение узлу через ETCP с заголовком: [ETCP_RT_ID_DB_SYNC:1][group_id_be:8][payload]. Возвращает код etcp_send +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); if (!entry) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "queue_entry_new"); return -1; } 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; - uint64_t hb = htobe64(hash); - memcpy(buf + 1, &hb, 8); + uint64_t gb = htobe64(group_id); + memcpy(buf + 1, &gb, 8); memcpy(buf + 9, payload, plen); entry->dgram = buf; entry->len = plen + 9; struct ETCP_CONN* conn = instance_find_conn(db->inst, node_id); if (!conn || !conn->links_up) { - DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync: send DROP — no conn to %04llX (hash=%016llx, type=%02x)", - (unsigned long long)(node_id >> 16), hash, plen > 0 ? payload[0] : 0, + DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync: send DROP — no conn to %04llX (group=%016llx, type=%02x)", + (unsigned long long)(node_id >> 16), group_id, plen > 0 ? payload[0] : 0, (void*)conn, conn ? conn->links_up : -1); queue_entry_free(entry); 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; } -// Отправляет сообщение через 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) { - 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; } 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]; const uint8_t* payload = entry->dgram + 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) { - /* lazy-register: find table msg_ with matching SHA256 */ + /* lazy-register: найти таблицу msg_, у которой strtoull(ch_id) == group_id */ sqlite3_stmt* st = NULL; if (sqlite3_prepare_v2(db->db, "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* ch_id = (const char*)sqlite3_column_text(st, 1); 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); } - if (h == hash) { - si = db_sync_instance_add(db->inst, tbl, hash, 1); + uint64_t gid = strtoull(ch_id, NULL, 10); + if (gid == group_id) { + si = db_sync_instance_add(db->inst, tbl, group_id, 1); 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; } } @@ -1231,17 +1231,17 @@ static void db_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) } if (!si) { DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, - "recv msg type=0x%02x from %016llx hash=%016llx — instance NOT FOUND; sending DB_ERR_NOT_FOUND", - type, (unsigned long long)src, (unsigned long long)hash); - DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "→ instance not found for hash=%016llx", (unsigned long long)hash); + "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)group_id); + 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; - 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; } } if (!si->enabled) { 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; } @@ -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_ERROR: DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "→ handle ERROR"); db_handle_error(si, src, payload, plen); break; default: - DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC, "sync: unknown msg type 0x%02x from %04llX (hash=%016llx) — dropped", - type, (unsigned long long)(src >> 16), hash); + DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC, "sync: unknown msg type 0x%02x from %04llX (group=%016llx) — dropped", + type, (unsigned long long)(src >> 16), group_id); break; } queue_dgram_free(entry); @@ -1707,19 +1707,19 @@ void db_sync_destroy(struct UTUN_INSTANCE* inst) } // Создаёт 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; } - 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; if (!db->enabled || !db->db) return NULL; - struct DB_SYNC_INSTANCE* si = db_instance_find(db, hash); - if (si) { DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "instance already exists hash=%016llx", (unsigned long long)hash); return si; } + struct DB_SYNC_INSTANCE* si = db_instance_find(db, group_id); + 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); if (!si) return NULL; - si->hash = hash; + si->group_id = group_id; snprintf(si->table_name, sizeof(si->table_name), "%s", table_name); // 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); - 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"); @@ -1783,8 +1783,8 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const } 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", - table_name, si->table_name, (unsigned long long)hash, + "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)group_id, (unsigned long long)si->next_id, peers_found, peers_synced, db_count(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) { 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 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--; } - 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 всем синхронизированным пирам diff --git a/src/chat/db_sync.h b/src/chat/db_sync.h index 6b28ccca..ca16654b 100644 --- a/src/chat/db_sync.h +++ b/src/chat/db_sync.h @@ -24,8 +24,7 @@ // // Синхронизация: // - При поднятии ETCP-соединения с пиром для каждого инстанса запускается полная синхронизация -// - Каждое сообщение содержит instance_hash (первые 64 бита SHA256(name||id_be)), -// что позволяет маршрутизировать сообщения к нужному инстансу на приёмной стороне +// - Каждое сообщение содержит group_id (uint64, binary), что позволяет маршрутизировать сообщения к нужному инстансу на приёмной стороне // - Если пир не имеет инстанса с таким хешем — возвращает DB_MSG_ERROR(0x01) // - Если пир добавляет инстанс позже — сам инициирует sync в нашу сторону // - Протокол: сравнение chain_hash для поиска расхождений (7 типов сообщений) @@ -97,7 +96,7 @@ int db_sync_init(struct UTUN_INSTANCE* inst); void db_sync_destroy(struct UTUN_INSTANCE* inst); // Instance management -struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const char* table_name, uint64_t hash, 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); // Data operations (per-instance) diff --git a/src/chat/merkle_sync.c b/src/chat/merkle_sync.c index 3297fe54..1520faf6 100644 --- a/src/chat/merkle_sync.c +++ b/src/chat/merkle_sync.c @@ -3,10 +3,12 @@ #include "../utun_instance.h" #include "../transport_layer/etcp_api.h" #include "../transport_layer/etcp.h" +#include "../routing_layer/topo_group.h" #include "../../lib/debug_config.h" #include "../../lib/mem.h" #include "../../lib/u_async.h" #include "../../lib/ll_queue.h" +#include "../../lib/platform_compat.h" #ifdef UTUN_HAVE_STANDBY #include "standby.h" #endif @@ -240,6 +242,29 @@ static int _get_level_hashes(struct merkle_sync* ms, const char* ns, 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 ── */ 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, const char* ns, const uint8_t* item_data, size_t item_len) { - uint8_t ch_len = (uint8_t)strlen(ns); - size_t sz = 1 + ch_len + 1 + 1 + 1 + 1 + 2 + item_len; + size_t sz = 1 + 8 + 1 + 1 + 1 + 1 + 2 + item_len; uint8_t* buf = u_malloc(sz); if (!buf) return -1; 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++ = 0; /* level=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, 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); - 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; + size_t max_sz = 1 + 8 + 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++ = ch_len; memcpy(p, ns, ch_len); p += ch_len; + ms_write_group_id(p, ns); p += 8; *p++ = 0x01; /* MS_MSG_HASHES */ *p++ = level; *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, 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); - 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); + size_t max_sz = 1 + 8 + 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++ = ch_len; memcpy(p, ns, ch_len); p += ch_len; + ms_write_group_id(p, ns); p += 8; *p++ = 0x03; /* MS_MSG_BATCH */ *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) s->pending[parent_idx].sub_count++; } - uint8_t ch_len = (uint8_t)strlen(ns); - size_t rs = 1 + 1 + ch_len + 1; + size_t rs = 1 + 8 + 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++ = ch_len; - memcpy(wr, ns, ch_len); wr += ch_len; + ms_write_group_id(wr, ns); wr += 8; *wr++ = 0x02; /* MS_MSG_REQUEST */ *wr++ = (uint8_t)rcount; 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) { - if (!entry || entry->len < 4) { + if (!entry || entry->len < 10) { if (entry) { if (entry->dgram) u_free(entry->dgram); queue_entry_free(entry); } return; } @@ -673,14 +693,13 @@ static void _recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { const uint8_t* d = entry->dgram; size_t dlen = entry->len; - uint8_t ch_len = 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; } - 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]; - 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 + 3 + ch_len; - size_t plen = dlen - 3 - ch_len; + uint64_t group_id = be64toh(*(const uint64_t*)(d + 1)); + char ns[64]; + if (ms_group_id_to_ns(inst, group_id, ns, sizeof(ns)) != 0) { u_free(entry->dgram); queue_entry_free(entry); return; } + uint8_t type = d[9]; + 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); + const uint8_t* pl = d + 10; + size_t plen = dlen - 10; switch (type) { 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", MS_ID, ns, (unsigned long long)key, type, len); - uint8_t ch_len = (uint8_t)strlen(ns); - if (ch_len > 63) return; - size_t pkt = 1 + ch_len + 1 + 8 + 1 + len; + size_t pkt = 1 + 8 + 1 + 8 + 1 + len; uint8_t* buf = u_malloc(pkt); if (!buf) return; 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 */ memcpy(p, &key, 8); p += 8; *p++ = type; diff --git a/tests/test_merkle_sync.c b/tests/test_merkle_sync.c index 08f68fbe..4b79370f 100644 --- a/tests/test_merkle_sync.c +++ b/tests/test_merkle_sync.c @@ -11,6 +11,7 @@ #include "../lib/socket_compat.h" #include "merkle_sync.h" #include "../src/utun_instance.h" +#include "../src/routing_layer/topo_group.h" #include "../src/config_updater.h" #include "../src/tun_if.h" #include "etcp.h" @@ -19,6 +20,21 @@ #include #include +/* 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. * Must match struct ms_session prefix in merkle_sync.c exactly. */ 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; inst.msync = &ms; - merkle_sync_recompute_path(&inst, "test", 0x1234ULL); - if (_db_count_rows(db, "test") == 0) PASS(); else FAIL("got %d rows", _db_count_rows(db, "test")); + merkle_sync_recompute_path(&inst, MS_NS_TEST, 0x1234ULL); + 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); } @@ -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; inst.msync = &ms; - merkle_sync_recompute_path(&inst, "test", key); - int rows = _db_count_rows(db, "test"); + merkle_sync_recompute_path(&inst, MS_NS_TEST, key); + int rows = _db_count_rows(db, MS_NS_TEST); if (rows != 5) { FAIL("expected 5 rows got %d", rows); sqlite3_close(db); return; } /* check hashes are non-zero */ int ok = 1; for (uint8_t lv = 1; lv <= 5 && ok; 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; 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; } @@ -378,8 +394,8 @@ static void test_tree_multi_item_same_bucket(void) { inst.msync = &ms; /* recompute for both */ - merkle_sync_recompute_path(&inst, "test", k1); - merkle_sync_recompute_path(&inst, "test", k2); + merkle_sync_recompute_path(&inst, MS_NS_TEST, k1); + merkle_sync_recompute_path(&inst, MS_NS_TEST, k2); uint64_t pf1 = merkle_sync_level_prefix(k1, 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_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"); 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; inst.msync = &ms; - merkle_sync_recompute_path(&inst, "test", k); - if (_db_count_rows(db, "test") != 5) { FAIL("expected 5 rows"); sqlite3_close(db); return; } + merkle_sync_recompute_path(&inst, MS_NS_TEST, k); + if (_db_count_rows(db, MS_NS_TEST) != 5) { FAIL("expected 5 rows"); sqlite3_close(db); return; } /* remove item and recompute */ data.count = 0; - merkle_sync_recompute_path(&inst, "test", k); - if (_db_count_rows(db, "test") != 0) { FAIL("expected 0 rows got %d", _db_count_rows(db, "test")); sqlite3_close(db); return; } + merkle_sync_recompute_path(&inst, MS_NS_TEST, k); + 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); - 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; 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"); @@ -694,6 +710,7 @@ static int _intg_init_two(void) { _intg_setup_contexts(); 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 || merkle_sync_init(i_b, MS_SVC_ID, &g_test_ops, &ctx_b) != 0) { _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); 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) { @@ -736,14 +753,14 @@ static void test_peer_empty(void) { if (_intg_init_two() != 0) { FAIL("setup failed"); return; } _intg_ins_many(&data_b, i_b, 0, 5); 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)) { FAIL("sync timeout"); _intg_cleanup(); return; } printf(" A data=%d B data=%d, A tree=%d B tree=%d\n", data_a.count, data_b.count, - _db_count_rows(i_a->topo_sqlite_db, "test"), - _db_count_rows(i_b->topo_sqlite_db, "test")); - if (_intg_verify_all(i_a, i_b, &data_a, &data_b, "test") != 0) + _db_count_rows(i_a->topo_sqlite_db, MS_NS_TEST), + _db_count_rows(i_b->topo_sqlite_db, MS_NS_TEST)); + if (_intg_verify_all(i_a, i_b, &data_a, &data_b, MS_NS_TEST) != 0) { FAIL("verify failed"); _intg_cleanup(); return; } PASS(); _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_b, i_b, 0, 5); 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); 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; } PASS(); _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_b, i_b, 10, 2); 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)) { 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; } PASS(); _intg_cleanup(); @@ -806,17 +823,17 @@ static void test_randomized_two(void) { _data_insert(&data_b, k, (uint32_t)rand()); } _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_b.count; i++) merkle_sync_recompute_path(i_b, "rnd", data_b.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, MS_NS_RND, data_b.items[i].key); 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 (_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", data_a.count, data_b.count, - _db_count_rows(i_a->topo_sqlite_db, "rnd"), - _db_count_rows(i_b->topo_sqlite_db, "rnd")); + _db_count_rows(i_a->topo_sqlite_db, MS_NS_RND), + _db_count_rows(i_b->topo_sqlite_db, MS_NS_RND)); /* dump items only in A */ for (int ai = 0; ai < data_a.count; ai++) { int found = 0; @@ -839,16 +856,16 @@ static void test_randomized_two(void) { sqlite3_prepare_v2(i_a->topo_sqlite_db, "SELECT level,prefix64 FROM merkle_tree_hash WHERE namespace=? ORDER BY level,prefix64", -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) { uint8_t lv = (uint8_t)sqlite3_column_int(st, 0); uint64_t pf = (uint64_t)sqlite3_column_int64(st, 1); EVP_MD_CTX* c = EVP_MD_CTX_new(); 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]; 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); 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); _intg_setup_contexts(); 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 || 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) @@ -890,8 +908,8 @@ static int _intg_init_three(void) { static int _cond_all_synced(void) { 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_b->topo_sqlite_db, i_c->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, MS_NS_RND) != 0) return 0; 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 < 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()); } - _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_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_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_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, 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, MS_NS_RND, data_c.items[i].key); done_sync = 0; i_phase = 0; /* 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; } done_sync = 0; i_phase = 0; /* 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; } /* wait for broadcast relay */ 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) { 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 || - _db_compare_trees(i_b->topo_sqlite_db, i_c->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, MS_NS_RND) != 0) { FAIL("tree mismatch iter %d", iter); _intg_cleanup(); break; } _intg_cleanup(); } @@ -936,9 +954,9 @@ static void test_cancel(void) { if (_intg_init_two() != 0) { FAIL("setup failed"); return; } _intg_ins_many(&data_b, i_b, 0, 5); 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); - 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); if (!done_sync) PASS(); else FAIL("done_cb was called"); _intg_cleanup(); @@ -953,8 +971,8 @@ static void test_start_overwrite(void) { if (_intg_init_two() != 0) { FAIL("setup failed"); return; } _intg_ins_many(&data_b, i_b, 0, 3); 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, "test", _sw_done2, 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, MS_NS_TEST, _sw_done2, NULL); 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); _intg_cleanup(); @@ -1043,8 +1061,8 @@ static void _str_spam_cb(void* arg) { /* spammer: bump own member, sync to hub */ uint64_t key = ((uint64_t)idx) << 60; _data_insert(&gs->data[idx], key, (uint32_t)(gs->data[idx].count + 1)); - merkle_sync_recompute_path(gs->inst[idx], "stress", key); - merkle_sync_start(gs->inst[idx], gs->inst[STRESS_HUB]->node_id, "stress", NULL, NULL); + merkle_sync_recompute_path(gs->inst[idx], MS_NS_STRESS, key); + merkle_sync_start(gs->inst[idx], gs->inst[STRESS_HUB]->node_id, MS_NS_STRESS, NULL, NULL); _str_spam_schedule(gs, idx); } @@ -1053,7 +1071,7 @@ static void _str_obs_spam_cb(void* arg) { if (!gs || !gs->spam_active) return; 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 */ 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].inst = s->inst[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) { 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++) { uint64_t key = ((uint64_t)i) << 60; _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 */ s->sync_pending = STRESS_N - 1; 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); /* observers: sync with hub too */ @@ -1236,7 +1255,7 @@ static void _str_start_final_round(void) { s->final_pending = STRESS_N - 1; printf(" final round %d: syncing %d spokes -> hub\n", s->final_round, s->final_pending); 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) { @@ -1251,17 +1270,17 @@ static void _str_verify(void) { int ok = 1; /* 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++) { - 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; } } /* 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); 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; } } @@ -1271,16 +1290,16 @@ static void _str_verify(void) { sqlite3_prepare_v2(s->inst[i]->topo_sqlite_db, "SELECT level,prefix64 FROM merkle_tree_hash WHERE namespace=? ORDER BY level,prefix64", -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) { uint8_t lv = (uint8_t)sqlite3_column_int(st, 0); uint64_t pf = (uint64_t)sqlite3_column_int64(st, 1); EVP_MD_CTX* c = EVP_MD_CTX_new(); 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]; 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); if (memcmp(expected, scopy, MT_HASH_SIZE) != 0) { 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) { struct str_state* s = gs; 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]; memcpy(root_copy, root0, MT_HASH_SIZE); int converged = 1; 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 (converged) {