Browse Source

chatgui: migrate message sync from chat_sync to db_sync

- Replace db_sync_stub.c with real src/db_sync.c in libutun
- Remove message sync protocol from chat_sync (INIT_SYNC/INIT_RESP/SEND_DATA/PUSH/ACK_PUSH/SYNC_DONE)
- Keep chat_sync auxiliary: join protocol, peer management, auto-connect
- Remove chat_core_insert_record, cursor_*, mark_sent/ttl_delete stubs
- chat_core on_db_sync_insert now only does SQLite INSERT + GUI notify
- Move db_sync_init() from instance_init_common to utun_instance_init for config override timing
- Force db_sync_enabled=1 in utun_node.cpp
- db_sync handles all P2P message sync via service 0x20
topo_upd
Evgeny 3 months ago
parent
commit
6252556a02
  1. 6
      src/utun_instance.c
  2. 3
      tools/chatgui/CMakeLists.txt
  3. 2
      tools/chatgui/libutun/CMakeLists.txt
  4. 145
      tools/chatgui/transport/chat_core.c
  5. 7
      tools/chatgui/transport/chat_core.h
  6. 329
      tools/chatgui/transport/chat_sync.c
  7. 7
      tools/chatgui/transport/chat_sync.h
  8. 162
      tools/chatgui/transport/db_sync_stub.c
  9. 6
      tools/chatgui/transport/utun_node.cpp

6
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");
@ -550,6 +547,9 @@ 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);

3
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()

2
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

145
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; }

7
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);

329
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, &timestamp, 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, &timestamp, 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);

7
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);

162
tools/chatgui/transport/db_sync_stub.c

@ -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 <string.h>
#include <stdlib.h>
#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;
}

6
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);

Loading…
Cancel
Save