diff --git a/src/db_sync.c b/src/db_sync.c index b3354fbc..873cd4e1 100644 --- a/src/db_sync.c +++ b/src/db_sync.c @@ -324,7 +324,7 @@ static void db_cascade_from(struct DB_SYNC_INSTANCE* si, uint32_t from_pos) sqlite3_stmt* sel, *upd; if (si_prep(si, &sel, - "SELECT id,timestamp,author,author_signature FROM \"%s\"" + "SELECT id,timestamp,node_id,author_signature FROM \"%s\"" " ORDER BY timestamp, author_signature" " LIMIT -1 OFFSET ?") != SQLITE_OK) { sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return; } @@ -445,7 +445,7 @@ static int db_record_insert(struct DB_SYNC_INSTANCE* si, // Insert new record if (si_prep(si, &stmt, "INSERT INTO \"%s\"" - " (timestamp,author,id,chain_hash,flags,data," + " (timestamp,node_id,id,chain_hash,flags,data," " author_signature,delivered_peers,delivery_chain)" " VALUES (?,?,?,?,0,?,?,0,'')") != SQLITE_OK) { sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return -1; } @@ -464,7 +464,7 @@ static int db_record_insert(struct DB_SYNC_INSTANCE* si, if (do_cascade) { sqlite3_stmt* sel; if (si_prep(si, &sel, - "SELECT timestamp,id,author,author_signature FROM \"%s\"" + "SELECT timestamp,id,node_id,author_signature FROM \"%s\"" " WHERE timestamp>?1" " OR (timestamp=?1 AND author_signature>?2)" " ORDER BY timestamp, author_signature") != SQLITE_OK) @@ -551,7 +551,7 @@ static void si_delivery_update(struct DB_SYNC_INSTANCE* si, uint64_t ts, uint64_ sqlite3_stmt* stmt; if (si_prep(si, &stmt, "SELECT delivery_chain FROM \"%s\"" - " WHERE timestamp=? AND author=?") == SQLITE_OK) + " WHERE timestamp=? AND node_id=?") == SQLITE_OK) { sqlite3_bind_int64(stmt, 1, (sqlite3_int64)ts); sqlite3_bind_int64(stmt, 2, (sqlite3_int64)author); @@ -661,7 +661,7 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, uint64_t src, const sqlite3_stmt* stmt; si_prep(si, &stmt, - "SELECT id,timestamp,author,data,author_signature" + "SELECT id,timestamp,node_id,data,author_signature" " FROM \"%s\" ORDER BY timestamp, author_signature" " LIMIT ? OFFSET ?"); sqlite3_bind_int64(stmt, 1, (sqlite3_int64)b); @@ -713,7 +713,7 @@ static void db_handle_init_resp(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint16_t* rcp = (uint16_t*)(sdbuf + off); off += 2; sqlite3_stmt* stmt; si_prep(si, &stmt, - "SELECT id,timestamp,author,data,author_signature" + "SELECT id,timestamp,node_id,data,author_signature" " FROM \"%s\" ORDER BY timestamp, author_signature" " LIMIT ? OFFSET ?"); sqlite3_bind_int64(stmt, 1, (sqlite3_int64)b); @@ -996,8 +996,8 @@ static void db_handle_ack_push(struct DB_SYNC_INSTANCE* si, uint64_t src, const sqlite3_stmt* stmt; if (si_prep(si, &stmt, "UPDATE \"%s\" SET flags = flags | 1" - " WHERE timestamp=? AND author=?" - " AND author=? AND (flags & 1) = 0") == SQLITE_OK) + " WHERE timestamp=? AND node_id=?" + " AND node_id=? AND (flags & 1) = 0") == SQLITE_OK) { sqlite3_bind_int64(stmt, 1, (sqlite3_int64)ts); sqlite3_bind_int64(stmt, 2, (sqlite3_int64)author); @@ -1188,7 +1188,7 @@ static void db_verify_chain(struct DB_SYNC_INSTANCE* si) sqlite3_stmt* stmt; if (si_prep(si, &stmt, - "SELECT id,timestamp,author,author_signature,chain_hash FROM \"%s\"" + "SELECT id,timestamp,node_id,author_signature,chain_hash FROM \"%s\"" " ORDER BY timestamp, author_signature") != SQLITE_OK) return; for (uint32_t pos = 0; sqlite3_step(stmt) == SQLITE_ROW; pos++) { @@ -1305,7 +1305,7 @@ static void db_sync_instance_ttl_cb(void* arg) sqlite3_stmt* stmt; if (si_prep(si, &stmt, "DELETE FROM \"%s\"" - " WHERE author=? AND (flags & 1) = 0" + " WHERE node_id=? AND (flags & 1) = 0" " AND timestampttl_timer = uasync_set_timeout(si->db_sync->inst->ua, interval, si, db_sync_instance_ttl_cb, "db_sync_ttl"); @@ -1346,7 +1346,7 @@ int db_sync_init(struct UTUN_INSTANCE* inst) const char* dp = inst->config->global.db_path; char sp[512]; - if (dp[0]) snprintf(sp, sizeof(sp), "%s/sync", dp); + if (dp[0]) snprintf(sp, sizeof(sp), "%s/chats.db", dp); else snprintf(sp, sizeof(sp), "/tmp/utun_db_sync"); if (db_sqlite_open(db, sp) != 0) { @@ -1412,31 +1412,28 @@ void db_sync_destroy(struct UTUN_INSTANCE* inst) DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "db_sync: destroyed"); } -struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const char* name, uint64_t id) +struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const char* table_name, uint64_t hash) { - if (!inst || !inst->db_sync || !name || !name[0]) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "instance_add invalid args"); return NULL; } - DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "name=%s id=%016llx", name, (unsigned long long)id); + if (!inst || !inst->db_sync || !table_name || !table_name[0]) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "instance_add invalid args"); return NULL; } + DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "table=%s hash=%016llx", table_name, (unsigned long long)hash); struct DB_SYNC* db = inst->db_sync; if (!db->enabled || !db->db) return NULL; - if (!si_name_valid(name)) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "invalid instance name '%s'", name); return NULL; } - - uint64_t hash = db_instance_hash_compute(name, id); struct DB_SYNC_INSTANCE* si = db_instance_find(db, hash); if (si) { DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "instance already exists hash=%016llx", (unsigned long long)hash); return si; } si = db_instance_alloc(db); if (!si) return NULL; si->hash = hash; - snprintf(si->table_name, sizeof(si->table_name), "db_sync_%s_%llx", name, (unsigned long long)id); + snprintf(si->table_name, sizeof(si->table_name), "%s", table_name); - // Create table — PK is (timestamp, author), author_signature NOT NULL + // Create table — PK is (timestamp, node_id), author_signature NOT NULL { char sql[512]; snprintf(sql, sizeof(sql), "CREATE TABLE IF NOT EXISTS \"%s\" (" " timestamp INTEGER NOT NULL," - " author INTEGER NOT NULL," + " node_id INTEGER NOT NULL," " id INTEGER NOT NULL," " chain_hash BLOB NOT NULL," " flags INTEGER NOT NULL DEFAULT 0," @@ -1450,7 +1447,7 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const if (rc != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "CREATE TABLE %s: %s", si->table_name, sqlite3_errmsg(db->db)); si->enabled = 0; return si; } snprintf(sql, sizeof(sql), "CREATE INDEX IF NOT EXISTS \"idx_%s_ttl\"" - " ON \"%s\" (author, timestamp)", + " ON \"%s\" (node_id, timestamp)", si->table_name, si->table_name); rc = sqlite3_exec(db->db, sql, NULL, NULL, NULL); if (rc != SQLITE_OK) DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "CREATE INDEX %s: %s", si->table_name, sqlite3_errmsg(db->db)); @@ -1469,6 +1466,8 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const db_verify_chain(si); + DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "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); + si->ttl_timer = uasync_set_timeout(inst->ua, DB_SYNC_TTL_INTERVAL * 10000u, si, db_sync_instance_ttl_cb, "db_sync_ttl"); int peers_found = 0, peers_synced = 0; @@ -1487,8 +1486,8 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const } DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, - "instance added name=%s id=%llx tbl=%s hash=%016llx next_id=%llu peers_found=%d peers_synced=%d mc=%u", - name, (unsigned long long)id, si->table_name, (unsigned long long)hash, + "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, (unsigned long long)si->next_id, peers_found, peers_synced, db_count(si)); return si; } @@ -1592,7 +1591,7 @@ int db_sync_select(struct DB_SYNC_INSTANCE* si, uint32_t offset, uint32_t limit, sqlite3_stmt* stmt; if (si_prep(si, &stmt, - "SELECT id,timestamp,author,data," + "SELECT id,timestamp,node_id,data," "author_signature,delivered_peers,delivery_chain" " FROM \"%s\" ORDER BY timestamp, author_signature" " LIMIT ? OFFSET ?") != SQLITE_OK) diff --git a/src/db_sync.h b/src/db_sync.h index c620db54..09fd05a2 100644 --- a/src/db_sync.h +++ b/src/db_sync.h @@ -90,7 +90,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* name, uint64_t id); +struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const char* table_name, uint64_t hash); void db_sync_instance_remove(struct DB_SYNC_INSTANCE* si); // Data operations (per-instance) diff --git a/tools/chatgui/db/db_manager.cpp b/tools/chatgui/db/db_manager.cpp index c1142934..0a3c3a10 100644 --- a/tools/chatgui/db/db_manager.cpp +++ b/tools/chatgui/db/db_manager.cpp @@ -46,7 +46,7 @@ QByteArray DbManager::computeChainHash(const QByteArray& prevChain, QByteArray DbManager::getLastChainHash(const QString& chId) const { QString tbl = msgTableName(chId); QByteArray sql = QStringLiteral( - "SELECT chain_hash FROM \"%1\" ORDER BY timestamp, datahash DESC LIMIT 1") + "SELECT chain_hash FROM \"%1\" ORDER BY timestamp, author_signature DESC LIMIT 1") .arg(tbl).toUtf8(); sqlite3_stmt* stmt = prepareOrNull(sql.constData()); if (!stmt) return QByteArray(32, '\0'); @@ -125,24 +125,39 @@ QList DbManager::getChannels() const { QList DbManager::getMessages(const QString& chId, int limit) const { QList list; QString sql = QStringLiteral( - "SELECT id, node_id, content_type, data, timestamp, datahash, chain_hash," - " signature, is_outgoing, is_read" - " FROM \"%1\" ORDER BY timestamp, datahash ASC LIMIT ?").arg(msgTableName(chId)); + "SELECT id, node_id, data, timestamp, chain_hash, author_signature" + " FROM \"%1\" ORDER BY timestamp, author_signature ASC LIMIT ?").arg(msgTableName(chId)); sqlite3_stmt* stmt = prepareOrNull(sql.toUtf8().constData()); if (!stmt) return list; sqlite3_bind_int(stmt, 1, limit); + /* helper: parse JSON field from data column: {"n":,"ch":"","ct":"<>","d":"<>"} */ + auto parseJsonField = [](const QByteArray& json, const char* key) -> QByteArray { + QByteArray pat = "\"" + QByteArray(key) + "\":\""; + int idx = json.indexOf(pat); + if (idx < 0) return {}; + idx += pat.size(); + int end = idx; + while (end < json.size() && !(json[end] == '"' && (end == idx || json[end - 1] != '\\'))) end++; + if (end > idx) return QByteArray(json.constData() + idx, end - idx); + return {}; + }; + uint64_t myId = 0; + { sqlite3_stmt* ms = prepareOrNull("SELECT node_id FROM local_identity WHERE id=1"); + if (ms) { if (sqlite3_step(ms) == SQLITE_ROW) myId = (uint64_t)sqlite3_column_int64(ms, 0); sqlite3_finalize(ms); } } while (sqlite3_step(stmt) == SQLITE_ROW) { MessageRow m; m.id = sqlite3_column_int64(stmt, 0); m.authorNodeId = (quint64)sqlite3_column_int64(stmt, 1); - m.contentType = colText(stmt, 2); - m.data = colBlob(stmt, 3); - m.timestamp = sqlite3_column_int64(stmt, 4); - m.datahash = (uint64_t)sqlite3_column_int64(stmt, 5); - m.chainHash = colBlob(stmt, 6); - m.signature = colBlob(stmt, 7); - m.isOutgoing = sqlite3_column_int(stmt, 8) != 0; - m.isRead = sqlite3_column_int(stmt, 9) != 0; + QByteArray jdata = colBlob(stmt, 2); + m.contentType = QString::fromUtf8(parseJsonField(jdata, "ct")); + if (m.contentType.isEmpty()) m.contentType = "text/plain"; + m.data = parseJsonField(jdata, "d"); + m.timestamp = sqlite3_column_int64(stmt, 3); + m.datahash = 0; + m.chainHash = colBlob(stmt, 4); + m.signature = colBlob(stmt, 5); + m.isOutgoing = (m.authorNodeId == myId); + m.isRead = true; list.append(m); } sqlite3_finalize(stmt); @@ -376,7 +391,7 @@ QByteArray DbManager::syncCount(const QString& chId) const { QByteArray DbManager::syncChainHashAt(const QString& chId, uint32_t pos) const { QString sql = QStringLiteral( - "SELECT chain_hash FROM \"%1\" ORDER BY timestamp, datahash ASC LIMIT 1 OFFSET ?") + "SELECT chain_hash FROM \"%1\" ORDER BY timestamp, author_signature ASC LIMIT 1 OFFSET ?") .arg(msgTableName(chId)); sqlite3_stmt* stmt = prepareOrNull(sql.toUtf8().constData()); if (!stmt) return QByteArray(32, '\0'); @@ -389,8 +404,8 @@ QByteArray DbManager::syncChainHashAt(const QString& chId, uint32_t pos) const { QByteArray DbManager::syncCursorOpen(const QString& chId) { QString sql = QStringLiteral( - "SELECT timestamp, datahash, node_id, content_type, data" - " FROM \"%1\" ORDER BY timestamp, datahash ASC").arg(msgTableName(chId)); + "SELECT timestamp, node_id, data, author_signature" + " FROM \"%1\" ORDER BY timestamp, author_signature ASC").arg(msgTableName(chId)); sqlite3_stmt* stmt = prepareOrNull(sql.toUtf8().constData()); if (!stmt) { QByteArray r(4, '\0'); return r; } uint32_t id = m_nextCursorId++; @@ -405,10 +420,23 @@ QByteArray DbManager::syncCursorNext(uint32_t cursorId) { /* return: [ts:8][dh:8][nid:8][ct_len:1][ct:var][data_len:4][data:var] */ qint64 ts = sqlite3_column_int64(it.value(), 0); - uint64_t dh = (uint64_t)sqlite3_column_int64(it.value(), 1); - uint64_t nid = (uint64_t)sqlite3_column_int64(it.value(), 2); - QString ct = colText(it.value(), 3); - QByteArray data = colBlob(it.value(), 4); + uint64_t nid = (uint64_t)sqlite3_column_int64(it.value(), 1); + QByteArray jdata = colBlob(it.value(), 2); + /* parse JSON from data column */ + auto jField = [](const QByteArray& j, const char* k) -> QByteArray { + QByteArray p = "\"" + QByteArray(k) + "\":\""; + int i = j.indexOf(p); if (i < 0) return {}; + i += p.size(); int e = i; + while (e < j.size() && !(j[e] == '"' && (e == i || j[e - 1] != '\\'))) e++; + return (e > i) ? QByteArray(j.constData() + i, e - i) : QByteArray(); + }; + QByteArray ct_raw = jField(jdata, "ct"); + QString ct = QString::fromUtf8(ct_raw); + if (ct.isEmpty()) ct = "text/plain"; + QByteArray ddata = jField(jdata, "d"); + /* compute datahash for wire format compat */ + uint8_t dh_buf[32]; SHA256((const uint8_t*)ddata.constData(), ddata.size(), dh_buf); + uint64_t dh; memcpy(&dh, dh_buf, 8); QByteArray out; out.append((const char*)&ts, 8); @@ -417,9 +445,9 @@ QByteArray DbManager::syncCursorNext(uint32_t cursorId) { QByteArray ctb = ct.toUtf8(); out.append((char)(uint8_t)ctb.size()); out.append(ctb); - uint32_t dLen = (uint32_t)data.size(); + uint32_t dLen = (uint32_t)ddata.size(); out.append((const char*)&dLen, 4); - out.append(data); + out.append(ddata); return out; } diff --git a/tools/chatgui/transport/chat_core.c b/tools/chatgui/transport/chat_core.c index 237e7b52..0b3a240b 100644 --- a/tools/chatgui/transport/chat_core.c +++ b/tools/chatgui/transport/chat_core.c @@ -49,7 +49,7 @@ static struct chat_core_ctx { static struct DB_SYNC_INSTANCE* si_find(const char* ch_id); static void si_register(struct DB_SYNC_INSTANCE* si, const char* ch_id); -static void on_db_sync_insert(struct DB_SYNC_INSTANCE* si, uint64_t record_ts, const char* data, size_t len, uint64_t author, void* arg); +static void on_msg_inserted(struct DB_SYNC_INSTANCE* si, uint64_t record_ts, const char* data, size_t len, uint64_t author, void* arg); /* ─── утилиты ─── */ @@ -399,34 +399,11 @@ void chat_core_submit_message(struct chat_msg_submit* req) { int ret = db_sync_insert_signed(si, json, strlen(json), sig, 64, ts); if (ret != 0) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: db_sync_insert_signed failed ret=%d", CC_ID, ret); - /* Insert directly to msg_ table as fallback (dual-write not triggered by callback) */ - char tbl[80]; msg_table_name(req->channel_id, tbl, sizeof(tbl)); - - /* compute datahash and chain_hash */ - uint8_t dh_buf2[32]; SHA256(req->data, req->data_len, dh_buf2); - uint64_t datahash2; memcpy(&datahash2, dh_buf2, 8); - uint8_t prev_ch2[32]; msg_get_prev_chain_hash(tbl, prev_ch2); - uint8_t ch_buf2[48]; memcpy(ch_buf2, prev_ch2, 32); memcpy(ch_buf2+32, &req->timestamp, 8); memcpy(ch_buf2+40, &datahash2, 8); - uint8_t ch_hash2[32]; SHA256(ch_buf2, 48, ch_hash2); - - char sql[512]; snprintf(sql, sizeof(sql), "INSERT OR IGNORE INTO \"%s\" (node_id,content_type,data,timestamp,datahash,chain_hash,signature,is_outgoing,is_read) VALUES(?,?,?,?,?,?,?,?,?)", tbl); - sqlite3_stmt* st=NULL; - if (sqlite3_prepare_v2(g_cc.db,sql,-1,&st,NULL)==SQLITE_OK) { - sqlite3_bind_int64(st,1,(sqlite3_int64)g_cc.my_node_id); - sqlite3_bind_text(st,2,req->content_type,-1,SQLITE_STATIC); - sqlite3_bind_blob(st,3,req->data,(int)req->data_len,SQLITE_STATIC); - sqlite3_bind_int64(st,4,(sqlite3_int64)req->timestamp); - sqlite3_bind_int64(st,5,(sqlite3_int64)datahash2); - sqlite3_bind_blob(st,6,ch_hash2,32,SQLITE_STATIC); - sqlite3_bind_blob(st,7,sig,64,SQLITE_STATIC); - sqlite3_bind_int(st,8,1); - sqlite3_bind_int(st,9,1); - sqlite3_step(st); sqlite3_finalize(st); - } - /* notify GUI anyway */ - uint8_t evt[65]; uint8_t cl=(uint8_t)strlen(req->channel_id); evt[0]=cl; memcpy(evt+1,req->channel_id,cl); - gui_bridge_post(GUI_EVT_MSG_RECEIVED, evt, 1+cl); + return; } + /* notify GUI */ + { uint8_t evt[65]; uint8_t cl=(uint8_t)strlen(req->channel_id); evt[0]=cl; memcpy(evt+1,req->channel_id,cl); + gui_bridge_post(GUI_EVT_MSG_RECEIVED, evt, 1+cl); } } void chat_core_submit_trampoline(void* arg) { chat_core_submit_message((struct chat_msg_submit*)arg); u_free(arg); } @@ -877,40 +854,16 @@ void chat_core_connect_auto(uint64_t node_id, } } -/* ─── подготовка инфраструктуры канала (msg-таблица + db_sync instance) ─── */ +/* ─── подготовка инфраструктуры канала (db_sync instance) ─── */ void chat_core_ensure_channel_ready(const char* ch_id) { if (!g_cc.initialized || !ch_id || !ch_id[0]) return; if (si_find(ch_id)) return; char tbl_msg[80]; msg_table_name(ch_id, tbl_msg, sizeof(tbl_msg)); - char sql[512]; - - snprintf(sql, sizeof(sql), - "CREATE TABLE IF NOT EXISTS \"%s\" (" - " id INTEGER PRIMARY KEY AUTOINCREMENT," - " node_id INTEGER NOT NULL," - " content_type TEXT NOT NULL," - " data BLOB NOT NULL," - " timestamp INTEGER NOT NULL," - " datahash INTEGER DEFAULT 0," - " chain_hash BLOB," - " signature BLOB," - " is_outgoing INTEGER DEFAULT 0," - " is_read INTEGER DEFAULT 0," - " UNIQUE(timestamp, node_id))", tbl_msg); - db_exec(sql); - /* migration for existing tables without datahash/chain_hash columns */ - snprintf(sql, sizeof(sql), "ALTER TABLE \"%s\" ADD COLUMN datahash INTEGER DEFAULT 0", tbl_msg); - sqlite3_exec(g_cc.db, sql, NULL, NULL, NULL); - snprintf(sql, sizeof(sql), "ALTER TABLE \"%s\" ADD COLUMN chain_hash BLOB", tbl_msg); - sqlite3_exec(g_cc.db, sql, NULL, NULL, NULL); - snprintf(sql, sizeof(sql), - "CREATE INDEX IF NOT EXISTS \"idx_%s_ts_node\" ON \"%s\"(timestamp, node_id)", - tbl_msg, tbl_msg); - db_exec(sql); char tbl_peers[80]; peers_table_name(ch_id, tbl_peers, sizeof(tbl_peers)); + char sql[512]; snprintf(sql, sizeof(sql), "CREATE TABLE IF NOT EXISTS \"%s\" (" " node_id INTEGER PRIMARY KEY," @@ -929,12 +882,12 @@ void chat_core_ensure_channel_ready(const char* ch_id) { 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, "chats", ch_hash); + struct DB_SYNC_INSTANCE* si = db_sync_instance_add(g_cc.inst, tbl_msg, ch_hash); if (si) { si_register(si, ch_id); - db_sync_set_insert_cb(si, on_db_sync_insert, u_strdup(ch_id)); - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: channel ready ch=%s hash=0x%016llx", - CC_ID, ch_id, (unsigned long long)ch_hash); + db_sync_set_insert_cb(si, on_msg_inserted, u_strdup(ch_id)); + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: channel ready ch=%s tbl=%s hash=0x%016llx", + CC_ID, ch_id, tbl_msg, (unsigned long long)ch_hash); } else { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: db_sync_instance_add failed for ch=%s", CC_ID, ch_id); @@ -1043,6 +996,13 @@ void chat_core_create_channel_trampoline(void* arg) { /* ─── db_sync helpers ─── */ +static void on_msg_inserted(struct DB_SYNC_INSTANCE* si, uint64_t record_ts, const char* data, size_t len, uint64_t author, void* arg) { + (void)si; (void)record_ts; (void)data; (void)len; (void)author; + const char* ch_id = (const char*)arg; + 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); +} + static struct DB_SYNC_INSTANCE* si_find(const char* ch_id) { for (int i=0; i=g_cc.si_capacity) { int nc=g_cc.si_capacity?g_cc.si_capacity*2:8; g_cc.si=u_realloc(g_cc.si,nc*sizeof(void*)); g_cc.si_ch_id=u_realloc(g_cc.si_ch_id,nc*sizeof(char*)); g_cc.si_capacity=nc; } g_cc.si[g_cc.si_count]=si; g_cc.si_ch_id[g_cc.si_count]=u_strdup(ch_id); g_cc.si_count++; -} - - -static void on_db_sync_insert(struct DB_SYNC_INSTANCE* si, uint64_t record_ts, const char* data, size_t len, uint64_t author, void* arg) { - const char* ch_id = (const char*)arg; - static int insert_count = 0; - /* parse JSON: {"n":,"ch":"","ct":"","d":""} */ - const char* p = data; const char* end = data + len; uint64_t jn = 0; - { /* "n" */ - const char* n = strstr(p, "\"n\":"); if (!n) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: on_db_sync_insert JSON parse fail: no '\"n\"' field ch=%s", CC_ID, ch_id); return; } - n += 4; jn = strtoull(n, NULL, 0); p = n; - } - const char* jct = "text/plain"; size_t jct_len = 10; - { /* "ct" */ - const char* ct = strstr(p, "\"ct\":\""); if (ct) { ct += 6; const char* ce = strchr(ct, '"'); if (ce) { jct = ct; jct_len = (size_t)(ce - ct); p = ce; } } - } - const char* jd_start = NULL; size_t jd_len = 0; - { /* "d" */ - const char* d = strstr(p, "\"d\":\""); if (d) { d += 5; const char* de = d; while (de < end) { if (*de == '"' && (de == d || *(de - 1) != '\\')) break; de++; } if (de < end) { jd_start = d; jd_len = (size_t)(de - d); } } - } - if (!jd_start) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: on_db_sync_insert JSON parse fail: no '\"d\"' field ch=%s", CC_ID, ch_id); return; } - - char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl)); - /* compute datahash and chain_hash */ - uint8_t dh_buf[32]; SHA256((const uint8_t*)jd_start, jd_len, dh_buf); - uint64_t datahash; memcpy(&datahash, dh_buf, 8); - uint8_t prev_ch[32]; msg_get_prev_chain_hash(tbl, prev_ch); - uint8_t ch_buf[48]; memcpy(ch_buf, prev_ch, 32); memcpy(ch_buf+32, &record_ts, 8); memcpy(ch_buf+40, &datahash, 8); - uint8_t ch_hash[32]; SHA256(ch_buf, 48, ch_hash); - - char sql[512]; snprintf(sql, sizeof(sql), - "INSERT OR IGNORE INTO \"%s\" (node_id,content_type,data,timestamp,datahash,chain_hash,signature,is_outgoing,is_read)" - " VALUES(?,?,?,?,?,?,?,?,?)", tbl); - sqlite3_stmt* st=NULL; - if (sqlite3_prepare_v2(g_cc.db,sql,-1,&st,NULL)!=SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: on_db_sync_insert prepare failed ch=%s", CC_ID, ch_id); return; } - sqlite3_bind_int64(st,1,(sqlite3_int64)jn); - sqlite3_bind_text(st,2,jct,(int)jct_len,SQLITE_STATIC); - sqlite3_bind_blob(st,3,jd_start,(int)jd_len,SQLITE_STATIC); - sqlite3_bind_int64(st,4,(sqlite3_int64)record_ts); - sqlite3_bind_int64(st,5,(sqlite3_int64)datahash); - sqlite3_bind_blob(st,6,ch_hash,32,SQLITE_STATIC); - { static const uint8_t z64[64]={0}; sqlite3_bind_blob(st,7,z64,64,SQLITE_STATIC); } - sqlite3_bind_int(st,8,(uint64_t)jn==g_cc.my_node_id?1:0); - sqlite3_bind_int(st,9,1); - int rc=sqlite3_step(st); sqlite3_finalize(st); - if (rc==SQLITE_DONE) { - insert_count++; - if (insert_count <= 3 || insert_count % 10 == 0) - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: on_db_sync_insert #%d ch=%s n=%llu ts=%llu ct=%.*s", - CC_ID, insert_count, ch_id, (unsigned long long)jn, (unsigned long long)record_ts, - (int)jct_len, jct); - 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); - } else if (rc == SQLITE_CONSTRAINT) { - DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: on_db_sync_insert duplicate ch=%s ts=%lld", - CC_ID, ch_id, (long long)record_ts); - } else { - DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: on_db_sync_insert step failed rc=%d ch=%s", - CC_ID, rc, ch_id); - } + uint32_t msg_cnt = db_sync_count(si); + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: si_register ch=%s mc=%u", CC_ID, ch_id, msg_cnt); } /* ─── сбор статуса (NTP + connections) и отправка в GUI ─── */