diff --git a/src/utun_instance.c b/src/utun_instance.c index 4d2741fb..f41417de 100644 --- a/src/utun_instance.c +++ b/src/utun_instance.c @@ -194,9 +194,6 @@ static int instance_init_common(struct UTUN_INSTANCE* instance, struct UASYNC* u return -1; } - // db_sync — распределённая таблица с репликацией (после etcp_router) - db_sync_init(instance); - // Bind DATA handler via etcp_router (after etcp_router_init) if (routing_bind(instance) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_ROUTING, "Failed to bind DATA via etcp_router"); @@ -549,7 +546,10 @@ void utun_instance_stop(struct UTUN_INSTANCE *instance) { int utun_instance_init(struct UTUN_INSTANCE *instance) { if (!instance) return -1; - + + // db_sync — распределённая таблица с репликацией + db_sync_init(instance); + // Set TUN interface in routing module if (instance->tun) { routing_set_tun(instance); diff --git a/tools/chatgui/CMakeLists.txt b/tools/chatgui/CMakeLists.txt index f0aa5336..cc3168b7 100644 --- a/tools/chatgui/CMakeLists.txt +++ b/tools/chatgui/CMakeLists.txt @@ -76,7 +76,6 @@ add_executable(chatgui transport/chat_sync.c transport/member_sync.c transport/merkle_sync.c - transport/db_sync_stub.c db/db_manager.cpp ../../lib/sqlite3.c resources/chatgui.qrc @@ -84,7 +83,7 @@ add_executable(chatgui target_include_directories(chatgui PRIVATE ${CMAKE_SOURCE_DIR}/../../lib ${CMAKE_SOURCE_DIR}/../../src ${CMAKE_SOURCE_DIR}/db) target_compile_definitions(chatgui PRIVATE SQLITE_THREADSAFE=1 USE_SQLITE) -set_source_files_properties(../../lib/sqlite3.c transport/chat_core.c transport/chat_sync.c transport/member_sync.c transport/merkle_sync.c transport/db_sync_stub.c PROPERTIES LANGUAGE C) +set_source_files_properties(../../lib/sqlite3.c transport/chat_core.c transport/chat_sync.c transport/member_sync.c transport/merkle_sync.c PROPERTIES LANGUAGE C) if(WIN32) target_link_libraries(chatgui PRIVATE ${QT_LIBS} rlottie::rlottie ZLIB::ZLIB OpenSSL::Crypto ZXing::ZXing utun pthread) else() diff --git a/tools/chatgui/libutun/CMakeLists.txt b/tools/chatgui/libutun/CMakeLists.txt index d40e834c..64718f27 100644 --- a/tools/chatgui/libutun/CMakeLists.txt +++ b/tools/chatgui/libutun/CMakeLists.txt @@ -62,7 +62,7 @@ set(UTUN_COMMON_SOURCES ${SRC_DIR}/topo_node.c ${SRC_DIR}/topo_node_sqlite.c ${SRC_DIR}/route_connectivity.c - # db_sync.c replaced by db_sync_stub.c (no P2P sync in chatgui) + ${SRC_DIR}/db_sync.c ${SRC_DIR}/conn_mgr.c ${TRANSPORT_DIR}/chat_sync.c ${TRANSPORT_DIR}/merkle_sync.c diff --git a/tools/chatgui/transport/chat_core.c b/tools/chatgui/transport/chat_core.c index 5e697fe4..5da3aad2 100644 --- a/tools/chatgui/transport/chat_core.c +++ b/tools/chatgui/transport/chat_core.c @@ -6,7 +6,6 @@ */ #include "chat_core.h" -#include "chat_sync.h" #include "db_sync.h" #include "gui_bridge.h" #include "topo_node_sqlite.h" @@ -46,9 +45,6 @@ static struct chat_core_ctx { char** si_ch_id; int si_count, si_capacity; - /* курсоры (макс 16) */ - sqlite3_stmt* cursors[16]; - uint32_t next_cursor_id; uint8_t initialized; } g_cc; @@ -265,9 +261,6 @@ void chat_core_destroy(struct UTUN_INSTANCE* inst) { if (!g_cc.initialized) return; g_cc.initialized = 0; - for (int i = 0; i < 16; i++) - if (g_cc.cursors[i]) sqlite3_finalize(g_cc.cursors[i]); - if (g_cc.db) { sqlite3_close(g_cc.db); g_cc.db = NULL; } g_cc.inst = NULL; DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: destroyed", CC_ID); @@ -430,138 +423,6 @@ int chat_core_chain_hash_at(const char* ch_id, uint32_t pos, uint8_t* hash_out) return 0; } -int chat_core_insert_record(const char* ch_id, const uint8_t* rec, size_t len) { - if (!g_cc.initialized || !rec || len < 29) { DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: insert_record malformed (len=%zu < 29) ch=%s", CC_ID, len, ch_id); return -1; } - - const uint8_t* p = rec; - size_t rem = len; - int64_t ts; uint64_t dh, nid; - memcpy(&ts, p, 8); p += 8; rem -= 8; - memcpy(&dh, p, 8); p += 8; rem -= 8; - memcpy(&nid, p, 8); p += 8; rem -= 8; - if (rem < 1) { DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: insert_record truncated at ct_len ch=%s", CC_ID, ch_id); return -1; } - uint8_t ct_len = *p; p++; rem--; - if (rem < ct_len) { DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: insert_record truncated at ct ch=%s", CC_ID, ch_id); return -1; } - char ct_buf[64]; memcpy(ct_buf, p, ct_len); ct_buf[ct_len] = '\0'; - p += ct_len; rem -= ct_len; - if (rem < 4) { DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: insert_record truncated at dlen ch=%s", CC_ID, ch_id); return -1; } - uint32_t dlen; memcpy(&dlen, p, 4); p += 4; rem -= 4; - if (rem < dlen) { DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: insert_record truncated at data ch=%s", CC_ID, ch_id); return -1; } - const uint8_t* rdata = p; p += dlen; rem -= dlen; - if (rem < 32) { DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: insert_record truncated at chain_hash ch=%s", CC_ID, ch_id); return -1; } - - uint8_t peer_ch[32]; memcpy(peer_ch, p, 32); - - /* verify chain_hash */ - uint8_t prev[32]; get_last_chain_hash(ch_id, prev); - uint8_t expected[32]; - compute_chain_hash(prev, ts, dh, expected); - if (memcmp(expected, peer_ch, 32) != 0) { - DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: chain_hash mismatch ch=%s ts=%lld dh=0x%016llx", - CC_ID, ch_id, (long long)ts, (unsigned long long)dh); - return -1; - } - - char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl)); - char sql[256]; snprintf(sql, sizeof(sql), - "INSERT OR IGNORE INTO \"%s\"" - " (node_id, content_type, data, timestamp, datahash, chain_hash," - " signature, is_outgoing, is_read)" - " VALUES(?,?,?,?,?,?,?,0,1)", tbl); - - sqlite3_stmt* stmt = NULL; - if (sqlite3_prepare_v2(g_cc.db, sql, -1, &stmt, NULL) != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: insert_record prepare failed ch=%s", CC_ID, ch_id); return -1; } - sqlite3_bind_int64(stmt, 1, (sqlite3_int64)nid); - sqlite3_bind_text(stmt, 2, ct_buf, ct_len, SQLITE_STATIC); - sqlite3_bind_blob(stmt, 3, rdata, (int)dlen, SQLITE_STATIC); - sqlite3_bind_int64(stmt, 4, ts); - sqlite3_bind_int64(stmt, 5, (sqlite3_int64)dh); - sqlite3_bind_blob(stmt, 6, peer_ch, 32, SQLITE_STATIC); - { - static const char zero64[64] = {0}; - sqlite3_bind_blob(stmt, 7, zero64, 64, SQLITE_STATIC); - } - int rc = sqlite3_step(stmt); - sqlite3_finalize(stmt); - if (rc == SQLITE_DONE) { - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: insert_record OK ch=%s ts=%lld dh=0x%016llx", - CC_ID, ch_id, (long long)ts, (unsigned long long)dh); - return 0; - } - if (rc == SQLITE_CONSTRAINT) return 1; - DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: insert_record step failed rc=%d ch=%s ts=%lld dh=0x%016llx", - CC_ID, rc, ch_id, (long long)ts, (unsigned long long)dh); - return -1; -} - -uint32_t chat_core_cursor_open(const char* ch_id) { - if (!g_cc.initialized) return 0; - for (int i = 0; i < 16; i++) { - if (g_cc.cursors[i]) continue; - char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl)); - char sql[256]; snprintf(sql, sizeof(sql), - "SELECT timestamp, datahash, node_id, content_type, data, chain_hash" - " FROM \"%s\" ORDER BY timestamp, datahash ASC", tbl); - sqlite3_stmt* stmt = NULL; - if (sqlite3_prepare_v2(g_cc.db, sql, -1, &stmt, NULL) != SQLITE_OK) return 0; - uint32_t id = ++g_cc.next_cursor_id; - g_cc.cursors[i] = stmt; - return id; - } - return 0; -} - -int chat_core_cursor_next(uint32_t cursor_id, uint8_t* buf, size_t buf_size, - size_t* out_len) { - if (!g_cc.initialized || !buf || !out_len) return -1; - for (int i = 0; i < 16; i++) { - if (!g_cc.cursors[i]) continue; - /* проверяем cursor_id = i+1 */ - if ((uint32_t)(i + 1) != cursor_id) continue; - sqlite3_stmt* stmt = g_cc.cursors[i]; - if (sqlite3_step(stmt) != SQLITE_ROW) { *out_len = 0; return 0; } - - int64_t ts = sqlite3_column_int64(stmt, 0); - uint64_t dh = (uint64_t)sqlite3_column_int64(stmt, 1); - uint64_t nid= (uint64_t)sqlite3_column_int64(stmt, 2); - const char* ct = (const char*)sqlite3_column_text(stmt, 3); - int ct_len = sqlite3_column_bytes(stmt, 3); - const uint8_t* dptr = (const uint8_t*)sqlite3_column_blob(stmt, 4); - int dlen = sqlite3_column_bytes(stmt, 4); - const uint8_t* ch_ptr = (const uint8_t*)sqlite3_column_blob(stmt, 5); - int ch_len = sqlite3_column_bytes(stmt, 5); - - size_t need = 8 + 8 + 8 + 1 + (size_t)ct_len + 4 + (size_t)dlen + 32; - if (need > buf_size) return -1; - - uint8_t* out = buf; - memcpy(out, &ts, 8); out += 8; - memcpy(out, &dh, 8); out += 8; - memcpy(out, &nid, 8); out += 8; - *out++ = (uint8_t)ct_len; - if (ct_len) { memcpy(out, ct, (size_t)ct_len); out += ct_len; } - uint32_t dl32 = (uint32_t)dlen; - memcpy(out, &dl32, 4); out += 4; - if (dlen) { memcpy(out, dptr, (size_t)dlen); out += dlen; } - if (ch_ptr && ch_len >= 32) memcpy(out, ch_ptr, 32); else memset(out, 0, 32); - out += 32; - *out_len = (size_t)(out - buf); - return 0; - } - return -1; -} - -void chat_core_cursor_close(uint32_t cursor_id) { - if (!g_cc.initialized) return; - for (int i = 0; i < 16; i++) { - if (g_cc.cursors[i] && (uint32_t)(i + 1) == cursor_id) { - sqlite3_finalize(g_cc.cursors[i]); - g_cc.cursors[i] = NULL; - return; - } - } -} - int chat_core_list_channels(uint8_t* buf, size_t buf_size, size_t* out_len) { if (!g_cc.initialized || !buf || !out_len) return -1; sqlite3_stmt* stmt = NULL; @@ -1296,7 +1157,6 @@ static void on_db_sync_insert(struct DB_SYNC_INSTANCE* si, const char* data, siz if (rc==SQLITE_DONE) { uint8_t evt[65]; uint8_t cl=(uint8_t)strlen(ch_id); evt[0]=cl; memcpy(evt+1,ch_id,cl); gui_bridge_post(GUI_EVT_MSG_RECEIVED, evt, 1+cl); - chat_sync_push(g_cc.inst, ch_id, jn, jct, (const uint8_t*)jd_start, (uint32_t)jd_len, ts, dh, chain_h); } else if (rc == SQLITE_CONSTRAINT) { DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: on_db_sync_insert duplicate ch=%s ts=%lld dh=0x%016llx", CC_ID, ch_id, (long long)ts, (unsigned long long)dh); @@ -1305,8 +1165,3 @@ static void on_db_sync_insert(struct DB_SYNC_INSTANCE* si, const char* data, siz CC_ID, rc, ch_id); } } - -/* ─── stubs for chat_sync compatibility (db_sync handles this now) ─── */ - -void chat_core_mark_sent(const char* ch_id, uint64_t ts, uint64_t dh, uint64_t peer_id) { (void)ch_id; (void)ts; (void)dh; (void)peer_id; } -void chat_core_ttl_delete(const char* ch_id, uint64_t node_id, uint64_t cutoff_us) { (void)ch_id; (void)node_id; (void)cutoff_us; } diff --git a/tools/chatgui/transport/chat_core.h b/tools/chatgui/transport/chat_core.h index abc1c730..a957a218 100644 --- a/tools/chatgui/transport/chat_core.h +++ b/tools/chatgui/transport/chat_core.h @@ -33,13 +33,10 @@ struct chat_msg_submit { void chat_core_submit_message(struct chat_msg_submit* req); -/* ── Stubs for chat_sync compatibility (todo: remove when chat_sync is fully replaced) ── */ - -void chat_core_mark_sent(const char* ch_id, uint64_t ts, uint64_t dh, uint64_t peer_id); -void chat_core_ttl_delete(const char* ch_id, uint64_t node_id, uint64_t cutoff_us); - /* ── DB-операции (оставлены для интроспекции) ── */ +uint32_t chat_core_count(const char* ch_id); +int chat_core_chain_hash_at(const char* ch_id, uint32_t pos, uint8_t* hash_out); int chat_core_list_channels(uint8_t* buf, size_t buf_size, size_t* out_len); int chat_core_list_peers(const char* ch_id, uint8_t* buf, size_t buf_size, size_t* out_len); int chat_core_load_nodeinfo(uint64_t node_id, uint8_t* buf, size_t buf_size, size_t* out_len); diff --git a/tools/chatgui/transport/chat_sync.c b/tools/chatgui/transport/chat_sync.c index 8f13d08b..607d17b7 100644 --- a/tools/chatgui/transport/chat_sync.c +++ b/tools/chatgui/transport/chat_sync.c @@ -267,7 +267,6 @@ struct chat_sync { struct channel_cache* channels; int channel_count; void* refresh_timer; - void* ttl_timer; void* info_req_timer; void* join_timer; uint8_t initialized; @@ -339,222 +338,6 @@ static struct channel_cache* cs_find(struct chat_sync* cs, const char* ch_id) { return NULL; } -/* ── INIT_SYNC handler ── */ - -static void cs_handle_init_sync(struct chat_sync* cs, uint64_t peer, - const char* ch_id, const uint8_t* pl, size_t len) { - if (len < 36) { DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: INIT_SYNC too short len=%zu peer=%016llx", CS_ID, len, (unsigned long long)peer); return; } - uint32_t peer_count; memcpy(&peer_count, pl, 4); - uint8_t peer_last_ch[32]; memcpy(peer_last_ch, pl + 4, 32); - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: RECV INIT_SYNC peer=%016llx ch=%s peer_count=%u my_count=%u", - CS_ID, (unsigned long long)peer, ch_id, peer_count, chat_core_count(ch_id)); - - uint32_t my_count = chat_core_count(ch_id); - uint32_t tp = peer_count < my_count ? peer_count : my_count; - if (tp > 0) tp--; - uint8_t my_ch[32]; - chat_core_chain_hash_at(ch_id, tp, my_ch); - - if (memcmp(my_ch, peer_last_ch, 32) == 0) { - /* synced */ - uint8_t resp[64]; resp[0] = CS_MSG_INIT_RESP; - memcpy(resp + 1, &my_count, 4); resp[5] = 0; - cs_send(cs, ch_id, peer, resp, 6); - struct channel_cache* ch = cs_find(cs, ch_id); - if (ch) ch->synced = CS_SYNC_DONE; - return; - } - - /* divergence — send INIT_RESP with 0 sparse for now, peer will handle */ - uint8_t resp[64]; resp[0] = CS_MSG_INIT_RESP; - memcpy(resp + 1, &my_count, 4); - resp[5] = 0; /* sparse_count=0 */ - cs_send(cs, ch_id, peer, resp, 6); -} - -/* ── INIT_RESP handler ── */ - -static void cs_handle_init_resp(struct chat_sync* cs, uint64_t peer, - const char* ch_id, const uint8_t* pl, size_t len) { - if (len < 5) { - DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: INIT_RESP too short len=%zu peer=%016llx ch=%s", - CS_ID, len, (unsigned long long)peer, ch_id); - return; - } - uint32_t peer_count = *(const uint32_t*)pl; - uint8_t sparse_count = pl[4]; - - struct channel_cache* ch = cs_find(cs, ch_id); - uint32_t my_count = chat_core_count(ch_id); - - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: RECV INIT_RESP peer=%016llx ch=%s peer_count=%u my_count=%u len=%zu", - CS_ID, (unsigned long long)peer, ch_id, peer_count, my_count, len); - - if (len >= 41) { - /* extended format: peer_count(4) + test_pos(4) + peer_ch(32) + sparse_count(1) */ - uint32_t test_pos = *(const uint32_t*)(pl + 4); - uint8_t peer_ch[32]; memcpy(peer_ch, pl + 8, 32); - sparse_count = pl[40]; - uint8_t my_ch[32]; chat_core_chain_hash_at(ch_id, test_pos, my_ch); - if (memcmp(my_ch, peer_ch, 32) == 0) { - if (ch && peer_count > test_pos + 1) { - uint32_t from = test_pos + 1; - uint8_t snd[7]; snd[0] = CS_MSG_SEND_DATA; - memcpy(snd + 1, &from, 4); uint16_t z = 0; memcpy(snd + 5, &z, 2); - cs_send(cs, ch_id, peer, snd, 7); - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: INIT_RESP chain_match, request data from=%u", CS_ID, from); - } else { - if (ch) ch->synced = CS_SYNC_DONE; - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: INIT_RESP chain_match, synced", CS_ID); - } - return; - } - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: INIT_RESP chain_mismatch, request from start", CS_ID); - uint32_t from = 0; - uint8_t snd[7]; snd[0] = CS_MSG_SEND_DATA; - memcpy(snd + 1, &from, 4); uint16_t z = 0; memcpy(snd + 5, &z, 2); - cs_send(cs, ch_id, peer, snd, 7); - return; - } - - /* short format (synced): peer_count(4) + sparse_count(1) */ - if (peer_count > my_count) { - uint32_t from = my_count; - uint8_t snd[7]; snd[0] = CS_MSG_SEND_DATA; - memcpy(snd + 1, &from, 4); uint16_t z = 0; memcpy(snd + 5, &z, 2); - cs_send(cs, ch_id, peer, snd, 7); - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: INIT_RESP peer ahead (peer=%u my=%u), request from=%u", - CS_ID, peer_count, my_count, from); - } else if (peer_count == my_count) { - if (ch) ch->synced = CS_SYNC_DONE; - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: INIT_RESP synced (peer=%u == my=%u)", - CS_ID, peer_count, my_count); - } else { - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: INIT_RESP peer behind (peer=%u < my=%u), waiting for peer SEND_DATA", - CS_ID, peer_count, my_count); - } -} - -/* ── SEND_DATA handler ── */ - -static void cs_handle_send_data(struct chat_sync* cs, uint64_t peer, - const char* ch_id, const uint8_t* pl, size_t len) { - if (len < 6) { DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: SEND_DATA too short len=%zu peer=%016llx", CS_ID, len, (unsigned long long)peer); return; } - uint32_t from = *(const uint32_t*)pl; - uint16_t count = *(const uint16_t*)(pl + 4); - - if (count == 0) { - /* peer requests OUR data starting from 'from' */ - uint32_t cur = chat_core_cursor_open(ch_id); - if (!cur) { DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: cursor_open failed ch=%s", CS_ID, ch_id); return; } - - uint16_t sent = 0; - uint8_t rec_buf[4096]; size_t rec_len; - - while (chat_core_cursor_next(cur, rec_buf, sizeof(rec_buf), &rec_len) == 0 - && rec_len > 0) { - uint8_t snd[8192]; uint32_t off = 0; - snd[off++] = CS_MSG_SEND_DATA; - uint32_t pos = from + sent; - memcpy(snd + off, &pos, 4); off += 4; - uint16_t c = 1; memcpy(snd + off, &c, 2); off += 2; - if (off + rec_len <= (uint32_t)sizeof(snd)) - { memcpy(snd + off, rec_buf, rec_len); off += (uint32_t)rec_len; } - cs_send(cs, ch_id, peer, snd, off); - sent++; - } - chat_core_cursor_close(cur); - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: SEND_DATA sent=%u records to=%016llx ch=%s from=%u", - CS_ID, sent, (unsigned long long)peer, ch_id, from); - return; - } - - /* peer sent US data — insert each record */ - uint16_t inserted = 0; - const uint8_t* ptr = pl + 6; - size_t remain = len - 6; - - for (uint16_t i = 0; i < count && remain > 0; i++) { - if (remain < 29) { DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: SEND_DATA record truncated (remain=%zu < 29) i=%u/%u ch=%s", CS_ID, remain, i, count, ch_id); break; } - const uint8_t* rec_start = ptr; - ptr += 8 + 8 + 8; remain -= 24; /* ts, dh, nid */ - if (remain < 1) break; - uint8_t ct_len = *ptr; ptr++; remain--; - if (remain < ct_len) break; - ptr += ct_len; remain -= ct_len; - if (remain < 4) break; - uint32_t dlen; memcpy(&dlen, ptr, 4); ptr += 4; remain -= 4; - if (remain < dlen + 32) break; - ptr += dlen; remain -= dlen; - ptr += 32; remain -= 32; /* chain_hash */ - - size_t reclen = (size_t)(ptr - rec_start); - if (chat_core_insert_record(ch_id, rec_start, reclen) == 0) inserted++; - } - - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: SEND_DATA recv=%u/%u from=%u to=%016llx ch=%s", - CS_ID, inserted, count, from, (unsigned long long)peer, ch_id); - - uint32_t next = from + count; - uint8_t snd[7]; snd[0] = CS_MSG_SEND_DATA; - memcpy(snd + 1, &next, 4); uint16_t z = 0; memcpy(snd + 5, &z, 2); - cs_send(cs, ch_id, peer, snd, 7); - - struct channel_cache* ch = cs_find(cs, ch_id); - if (ch) ch->synced = CS_SYNC_DONE; -} - -/* ── PUSH handler ── */ - -static void cs_handle_push(struct chat_sync* cs, uint64_t peer, - const char* ch_id, const uint8_t* pl, size_t len) { - if (len < 24 + 1 + 4) { DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: PUSH too short len=%zu peer=%016llx", CS_ID, len, (unsigned long long)peer); return; } - uint64_t ts, dh; - memcpy(&ts, pl, 8); memcpy(&dh, pl + 8, 8); - - const uint8_t* rec = pl + 24; - int reclen = (int)(len - 24); - int rc = chat_core_insert_record(ch_id, rec, (size_t)reclen); - - if (rc == 0) { - uint8_t ack[17]; ack[0] = CS_MSG_ACK_PUSH; - memcpy(ack + 1, &dh, 8); memcpy(ack + 9, &ts, 8); - cs_send(cs, ch_id, peer, ack, 17); - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: RECV PUSH OK ch=%s ts=%lld dh=0x%016llx, sending ACK", - CS_ID, ch_id, (long long)ts, (unsigned long long)dh); - - uint8_t evt[65]; evt[0] = (uint8_t)strlen(ch_id); - memcpy(evt + 1, ch_id, evt[0]); - gui_bridge_post(GUI_EVT_MSG_RECEIVED, evt, 1 + evt[0]); - - struct channel_cache* ch = cs_find(cs, ch_id); - if (ch) ch->msg_count++; - } else { - DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: PUSH insert_record failed rc=%d ch=%s ts=%lld dh=0x%016llx", - CS_ID, rc, ch_id, (long long)ts, (unsigned long long)dh); - } -} - -/* ── ACK_PUSH handler ── */ - -static void cs_handle_ack_push(struct chat_sync* cs, uint64_t peer, - const char* ch_id, const uint8_t* pl, size_t len) { - (void)peer; - if (len < 16) { DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: ACK_PUSH too short len=%zu peer=%016llx", CS_ID, len, (unsigned long long)peer); return; } - uint64_t ts, dh; memcpy(&dh, pl, 8); memcpy(&ts, pl + 8, 8); - chat_core_mark_sent(ch_id, ts, dh, cs->inst->node_id); -} - -/* ── SYNC_DONE handler ── */ - -static void cs_handle_sync_done(struct chat_sync* cs, uint64_t peer, - const char* ch_id, const uint8_t* pl, size_t len) { - (void)peer; - if (len < 4) { DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: SYNC_DONE too short len=%zu peer=%016llx", CS_ID, len, (unsigned long long)peer); return; } - uint32_t pc; memcpy(&pc, pl, 4); - struct channel_cache* ch = cs_find(cs, ch_id); - if (ch) { ch->msg_count = pc; ch->synced = CS_SYNC_DONE; } -} /* ── Forward declarations for new message handlers ── */ @@ -598,12 +381,6 @@ static void chat_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { CS_ID, cs_msg_name(type), (unsigned long long)peer, ch_id, plen); switch (type) { - case CS_MSG_INIT_SYNC: cs_handle_init_sync(g_cs, peer, ch_id, pl, plen); break; - case CS_MSG_INIT_RESP: cs_handle_init_resp(g_cs, peer, ch_id, pl, plen); break; - case CS_MSG_SEND_DATA: cs_handle_send_data(g_cs, peer, ch_id, pl, plen); break; - case CS_MSG_PUSH: cs_handle_push(g_cs, peer, ch_id, pl, plen); break; - case CS_MSG_ACK_PUSH: cs_handle_ack_push(g_cs, peer, ch_id, pl, plen); break; - case CS_MSG_SYNC_DONE: cs_handle_sync_done(g_cs, peer, ch_id, pl, plen); break; case CS_MSG_CHANNEL_INFO_REQ: cs_handle_channel_info_req(g_cs, peer, ch_id, pl, plen); break; case CS_MSG_CHANNEL_INFO_RESP:cs_handle_channel_info_resp(g_cs, peer, ch_id, pl, plen); break; case CS_MSG_CHANNEL_JOIN: cs_handle_channel_join(g_cs, peer, ch_id, pl, plen); break; @@ -725,29 +502,6 @@ static void cs_on_conn_up(struct ETCP_CONN* conn, void* arg) { return; } - /* sync path: initiate INIT_SYNC for shared channels */ - int synced_cnt = 0; - for (int i = 0; i < g_cs->channel_count; i++) { - struct channel_cache* ch = &g_cs->channels[i]; - int found = 0; - for (int j = 0; j < ch->peer_count; j++) - if (ch->peer_ids[j] == peer) { found = 1; break; } - if (!found) continue; - ch->synced = CS_SYNC_IN_PROGRESS; - synced_cnt++; - uint8_t msg[37]; msg[0] = CS_MSG_INIT_SYNC; - chat_core_chain_hash_at(ch->channel_id, - ch->msg_count > 0 ? ch->msg_count - 1 : 0, ch->last_chain_hash); - memcpy(msg + 1, &ch->msg_count, 4); - memcpy(msg + 5, ch->last_chain_hash, 32); - cs_send(g_cs, ch->channel_id, peer, msg, 37); - member_sync_start(g_cs->inst, peer, ch->channel_id, _on_member_sync_done, ch); - cs_post_channel_online(g_cs, ch->channel_id); - } - if (synced_cnt > 0) { - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: conn_up sync path peer=%016llx, syncing %d channels", - CS_ID, (unsigned long long)peer, synced_cnt); - } member_sync_set_online(g_cs->inst, peer, 1); } @@ -843,16 +597,6 @@ static void refresh_timer_cb(void* arg) { /* ── TTL cleanup ── */ -static void ttl_timer_cb(void* arg) { - (void)arg; - if (!g_cs || !g_cs->initialized) return; - uint64_t cutoff = get_time_us() - 86400000000ULL; - for (int i = 0; i < g_cs->channel_count; i++) - chat_core_ttl_delete(g_cs->channels[i].channel_id, - g_cs->inst->node_id, cutoff); - g_cs->ttl_timer = uasync_set_timeout(g_cs->inst->ua, - 3600u * 10000u, g_cs, ttl_timer_cb, "cs_ttl"); -} /* ── Public API ── */ @@ -880,8 +624,6 @@ int chat_sync_init(struct UTUN_INSTANCE* inst, cs->refresh_timer = uasync_set_timeout(inst->ua, 50000u, cs, refresh_timer_cb, "cs_refresh"); - cs->ttl_timer = uasync_set_timeout(inst->ua, 3600u * 10000u, cs, - ttl_timer_cb, "cs_ttl"); cs->info_req_timer = NULL; cs->join_timer = NULL; @@ -900,7 +642,6 @@ void chat_sync_destroy(struct UTUN_INSTANCE* inst) { etcp_unbind(inst, ETCP_RT_ID_CHAT_SYNC); if (cs->refresh_timer) { uasync_cancel_timeout(inst->ua, cs->refresh_timer); cs->refresh_timer = NULL; } - if (cs->ttl_timer) { uasync_cancel_timeout(inst->ua, cs->ttl_timer); cs->ttl_timer = NULL; } cs_cancel_proto_timers(cs); struct ETCP_CONN* c = inst->connections; @@ -916,55 +657,6 @@ void chat_sync_destroy(struct UTUN_INSTANCE* inst) { u_free(cs); } -int chat_sync_push(struct UTUN_INSTANCE* inst, - const char* ch_id, uint64_t node_id, - const char* content_type, const uint8_t* data, uint32_t data_len, - uint64_t timestamp, uint64_t datahash, - const uint8_t* chain_hash) { - if (!g_cs || !g_cs->initialized) return -1; - struct channel_cache* ch = cs_find(g_cs, ch_id); - - /* build record in insert_record format: ts(8)+dh(8)+nid(8)+ct(1)+ct+dlen(4)+data+chain_hash(32) */ - uint8_t ct_len = content_type ? (uint8_t)strlen(content_type) : 0; - if (ct_len > 63) ct_len = 63; - size_t rec_len = 24 + 1 + ct_len + 4 + data_len + 32; - /* PUSH wrapper: CS_MSG_PUSH(1) + header_ts(8) + header_dh(8) + pad(8) + record */ - size_t total = 1 + 24 + rec_len; - uint8_t* buf = u_malloc(total); - if (!buf) return -1; - - uint8_t* p = buf; - *p++ = CS_MSG_PUSH; - memcpy(p, ×tamp, 8); p += 8; - memcpy(p, &datahash, 8); p += 8; - uint64_t z = 0; memcpy(p, &z, 8); p += 8; /* pad */ - /* record: ts + dh + nid + ct + dlen + data + chain_hash */ - memcpy(p, ×tamp, 8); p += 8; - memcpy(p, &datahash, 8); p += 8; - memcpy(p, &node_id, 8); p += 8; - *p++ = ct_len; - if (ct_len) { memcpy(p, content_type, ct_len); p += ct_len; } - uint32_t dlen = data_len; - memcpy(p, &dlen, 4); p += 4; - if (data_len) { memcpy(p, data, data_len); p += data_len; } - memcpy(p, chain_hash ? chain_hash : (const uint8_t*)"\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0\0", 32); p += 32; - - int sent = 0; - uint64_t myid = g_cs->inst->node_id; - if (ch) { - for (int i = 0; i < ch->peer_count; i++) { - if (ch->peer_ids[i] == node_id || ch->peer_ids[i] == myid) continue; - cs_send(g_cs, ch_id, ch->peer_ids[i], buf, (size_t)(p - buf)); - sent++; - } - } - u_free(buf); - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: PUSH sent to %d peers ch=%s ts=%lld dh=0x%016llx", - CS_ID, sent, ch_id, (long long)timestamp, (unsigned long long)datahash); - (void)inst; - return sent; -} - void chat_sync_connect_node(struct UTUN_INSTANCE* inst, uint64_t node_id) { uint8_t buf[2048]; size_t buf_len; if (chat_core_load_nodeinfo(node_id, buf, sizeof(buf), &buf_len) == 0 && buf_len > 0) { @@ -1363,27 +1055,6 @@ static void cs_handle_welcome(struct chat_sync* cs, uint64_t peer, cs_refresh_channels(cs); - /* initiate message sync with the inviter after join */ - for (int i = 0; i < cs->channel_count; i++) { - struct channel_cache* ch2 = &cs->channels[i]; - if (strcmp(ch2->channel_id, ch_id) != 0) continue; - int found = 0; - for (int j = 0; j < ch2->peer_count; j++) - if (ch2->peer_ids[j] == peer) { found = 1; break; } - if (found) { - ch2->synced = CS_SYNC_IN_PROGRESS; - uint8_t msg[37]; msg[0] = CS_MSG_INIT_SYNC; - chat_core_chain_hash_at(ch2->channel_id, - ch2->msg_count > 0 ? ch2->msg_count - 1 : 0, ch2->last_chain_hash); - memcpy(msg + 1, &ch2->msg_count, 4); - memcpy(msg + 5, ch2->last_chain_hash, 32); - cs_send(cs, ch2->channel_id, peer, msg, 37); - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: WELCOME starting member_sync with inviter node=0x%016llx ch=%s", - CS_ID, (unsigned long long)peer, ch_id); - member_sync_start(cs->inst, peer, ch2->channel_id, _on_member_sync_done, ch2); - } - break; - } uint8_t evt[65]; uint8_t ch_id_len = (uint8_t)strlen(ch_id); evt[0] = ch_id_len; memcpy(evt + 1, ch_id, ch_id_len); diff --git a/tools/chatgui/transport/chat_sync.h b/tools/chatgui/transport/chat_sync.h index f3088705..70d20acb 100644 --- a/tools/chatgui/transport/chat_sync.h +++ b/tools/chatgui/transport/chat_sync.h @@ -46,13 +46,6 @@ int chat_sync_init(struct UTUN_INSTANCE* inst, void (*gui_result_cb)(void*, int, const uint8_t*, int)); void chat_sync_destroy(struct UTUN_INSTANCE* inst); -/* Вызывается из GUI (через gui_bridge_post_uasync) при отправке */ -int chat_sync_push(struct UTUN_INSTANCE* inst, - const char* channel_id, uint64_t node_id, - const char* content_type, const uint8_t* data, uint32_t data_len, - uint64_t timestamp, uint64_t datahash, - const uint8_t* chain_hash); - /* Прослойка conn_mgr: загрузить nodeinfo из БД и подключиться */ void chat_sync_connect_node(struct UTUN_INSTANCE* inst, uint64_t node_id); diff --git a/tools/chatgui/transport/db_sync_stub.c b/tools/chatgui/transport/db_sync_stub.c deleted file mode 100644 index 2fc680a2..00000000 --- a/tools/chatgui/transport/db_sync_stub.c +++ /dev/null @@ -1,162 +0,0 @@ -// db_sync_stub.c — Minimal db_sync replacement for chatgui (no P2P sync, local only) -// -// Назначение: заменяет полноценный db_sync.c когда синхронизация через ETCP не нужна. -// Сохраняет API-совместимость: chat_core.c использует db_sync_instance_add / db_sync_set_insert_cb / -// db_sync_insert_signed — стаб обрабатывает их локально, без сетевого обмена. -// -// utun_instance.c вызывает db_sync_init/db_sync_destroy — стаб создаёт минимальную инфраструктуру -// без SQLite и без регистрации ETCP-сервисов. - -#include "db_sync.h" -#include "utun_instance.h" -#include "../lib/mem.h" -#include "../lib/debug_config.h" -#include -#include - -#define DEBUG_CATEGORY_DB_SYNC DEBUG_CATEGORY_DEBUG - -/* ── Internal structs (match db_sync.h opaque types) ── */ - -struct DB_SYNC_INSTANCE { - struct DB_SYNC* db_sync; - uint8_t enabled; - uint64_t last_timestamp_ms; - db_sync_insert_cb on_insert; - void* on_insert_arg; -}; - -struct DB_SYNC { - struct UTUN_INSTANCE* inst; - struct DB_SYNC_INSTANCE* instances; - int instance_count, instance_capacity; - uint8_t enabled; -}; - -/* ── Init / Destroy ── */ - -int db_sync_init(struct UTUN_INSTANCE* inst) { - if (!inst) { - DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync stub: NULL instance"); - return -1; - } - if (inst->db_sync) { - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync stub: already initialized"); - return 0; - } - struct DB_SYNC* db = u_calloc(1, sizeof(struct DB_SYNC)); - if (!db) { - DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync stub: u_calloc failed"); - return -1; - } - db->inst = inst; - db->enabled = 1; - db->instance_count = 0; - db->instance_capacity = 0; - db->instances = NULL; - inst->db_sync = db; - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: stub initialized (local-only, no P2P sync)"); - return 0; -} - -void db_sync_destroy(struct UTUN_INSTANCE* inst) { - if (!inst || !inst->db_sync) return; - struct DB_SYNC* db = inst->db_sync; - inst->db_sync = NULL; - for (int i = 0; i < db->instance_count; i++) { - struct DB_SYNC_INSTANCE* si = &db->instances[i]; - if (si->on_insert_arg) u_free(si->on_insert_arg); - } - if (db->instances) u_free(db->instances); - u_free(db); - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: stub destroyed"); -} - -/* ── Instance management ── */ - -struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, - const char* name, uint64_t id) { - if (!inst || !inst->db_sync || !name || !name[0]) { - DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "db_sync stub: instance_add invalid args"); - return NULL; - } - struct DB_SYNC* db = inst->db_sync; - - if (db->instance_count >= db->instance_capacity) { - int nc = db->instance_capacity ? db->instance_capacity * 2 : 4; - struct DB_SYNC_INSTANCE* np = u_realloc(db->instances, nc * sizeof(*db->instances)); - if (!np) return NULL; - db->instances = np; - db->instance_capacity = nc; - } - struct DB_SYNC_INSTANCE* si = &db->instances[db->instance_count++]; - memset(si, 0, sizeof(*si)); - si->db_sync = db; - si->enabled = 1; - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync stub: instance added name=%s id=0x%016llx idx=%d", - name, (unsigned long long)id, db->instance_count - 1); - return si; -} - -void db_sync_instance_remove(struct DB_SYNC_INSTANCE* si) { - if (!si) return; - struct DB_SYNC* db = si->db_sync; - si->enabled = 0; - if (si->on_insert_arg) { u_free(si->on_insert_arg); si->on_insert_arg = NULL; } - int idx = (int)(si - db->instances); - if (idx >= 0 && idx < db->instance_count) { - memmove(&db->instances[idx], &db->instances[idx + 1], - (db->instance_count - idx - 1) * sizeof(*db->instances)); - db->instance_count--; - } - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync stub: instance removed idx=%d", idx); -} - -/* ── Insert ── */ - -int db_sync_insert_signed(struct DB_SYNC_INSTANCE* si, const char* json_data, - size_t len, const uint8_t* sig, size_t sig_len) { - (void)sig; (void)sig_len; - if (!si || !si->enabled || !json_data || len == 0) return -1; - struct timeval tv; utun_gettimeofday(&tv, NULL); - si->last_timestamp_ms = (uint64_t)tv.tv_sec * 1000ULL + (uint64_t)tv.tv_usec / 1000ULL; - if (si->on_insert) - si->on_insert(si, json_data, len, si->db_sync->inst->node_id, si->on_insert_arg); - return 0; -} - -int db_sync_insert_len(struct DB_SYNC_INSTANCE* si, const char* json_data, size_t len) { - return db_sync_insert_signed(si, json_data, len, NULL, 0); -} - -int db_sync_insert(struct DB_SYNC_INSTANCE* si, const char* json_data) { - if (!json_data) return -1; - return db_sync_insert_len(si, json_data, strlen(json_data)); -} - -/* ── Count / Select ── */ - -uint32_t db_sync_count(struct DB_SYNC_INSTANCE* si) { - (void)si; - return 0; -} - -int db_sync_select(struct DB_SYNC_INSTANCE* si, uint32_t offset, uint32_t limit, - db_sync_select_cb cb, void* arg) { - (void)si; (void)offset; (void)limit; (void)cb; (void)arg; - return 0; -} - -/* ── Callbacks ── */ - -void db_sync_set_insert_cb(struct DB_SYNC_INSTANCE* si, db_sync_insert_cb cb, void* arg) { - if (!si) return; - si->on_insert = cb; - if (si->on_insert_arg) u_free(si->on_insert_arg); - si->on_insert_arg = arg; -} - -uint64_t db_sync_get_last_timestamp(struct DB_SYNC_INSTANCE* si) { - if (!si) return 0; - return si->last_timestamp_ms; -} diff --git a/tools/chatgui/transport/utun_node.cpp b/tools/chatgui/transport/utun_node.cpp index d9991f24..958e87e4 100644 --- a/tools/chatgui/transport/utun_node.cpp +++ b/tools/chatgui/transport/utun_node.cpp @@ -224,6 +224,12 @@ void UtunNode::runLoop() { m_instance->test_user_ptr = this; + /* enable db_sync for chat message P2P sync */ + m_instance->config->global.db_sync_enabled = 1; + if (!m_dbPath.isEmpty()) + snprintf(m_instance->config->global.db_path, sizeof(m_instance->config->global.db_path), + "%s", m_dbPath.toUtf8().constData()); + /* set ua for gui_bridge before init, so GUI can post */ gui_bridge_set_uasync(ua);