Browse Source

db_sync: unified msg table — one DB, one table per channel

- db_sync: rename column author→node_id, open chats.db instead of sync file,
  db_sync_instance_add takes table_name directly
- chat_core: remove CREATE TABLE msg_ch_*, remove on_db_sync_insert (74 lines),
  remove fallback INSERT in submit_message, add minimal on_msg_inserted
  for GUI notification only
- db_manager: update getMessages for new schema (parse JSON from data column),
  fix ORDER BY to use author_signature instead of datahash,
  update syncCursorNext to parse JSON
- si_register simplified: db_count from same table, no more msg vs sync gap
topo_upd
Evgeny 3 months ago
parent
commit
e71cae16c1
  1. 47
      src/db_sync.c
  2. 2
      src/db_sync.h
  3. 70
      tools/chatgui/db/db_manager.cpp
  4. 138
      tools/chatgui/transport/chat_core.c

47
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; sqlite3_stmt* sel, *upd;
if (si_prep(si, &sel, 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" " ORDER BY timestamp, author_signature"
" LIMIT -1 OFFSET ?") != SQLITE_OK) " LIMIT -1 OFFSET ?") != SQLITE_OK)
{ sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return; } { 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 // Insert new record
if (si_prep(si, &stmt, if (si_prep(si, &stmt,
"INSERT INTO \"%s\"" "INSERT INTO \"%s\""
" (timestamp,author,id,chain_hash,flags,data," " (timestamp,node_id,id,chain_hash,flags,data,"
" author_signature,delivered_peers,delivery_chain)" " author_signature,delivered_peers,delivery_chain)"
" VALUES (?,?,?,?,0,?,?,0,'')") != SQLITE_OK) " VALUES (?,?,?,?,0,?,?,0,'')") != SQLITE_OK)
{ sqlite3_exec(db, "ROLLBACK", NULL, NULL, NULL); return -1; } { 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) { if (do_cascade) {
sqlite3_stmt* sel; sqlite3_stmt* sel;
if (si_prep(si, &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" " WHERE timestamp>?1"
" OR (timestamp=?1 AND author_signature>?2)" " OR (timestamp=?1 AND author_signature>?2)"
" ORDER BY timestamp, author_signature") != SQLITE_OK) " 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; sqlite3_stmt* stmt;
if (si_prep(si, &stmt, if (si_prep(si, &stmt,
"SELECT delivery_chain FROM \"%s\"" "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, 1, (sqlite3_int64)ts);
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)author); 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; sqlite3_stmt* stmt;
si_prep(si, &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" " FROM \"%s\" ORDER BY timestamp, author_signature"
" LIMIT ? OFFSET ?"); " LIMIT ? OFFSET ?");
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)b); 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; uint16_t* rcp = (uint16_t*)(sdbuf + off); off += 2;
sqlite3_stmt* stmt; sqlite3_stmt* stmt;
si_prep(si, &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" " FROM \"%s\" ORDER BY timestamp, author_signature"
" LIMIT ? OFFSET ?"); " LIMIT ? OFFSET ?");
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)b); 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; sqlite3_stmt* stmt;
if (si_prep(si, &stmt, if (si_prep(si, &stmt,
"UPDATE \"%s\" SET flags = flags | 1" "UPDATE \"%s\" SET flags = flags | 1"
" WHERE timestamp=? AND author=?" " WHERE timestamp=? AND node_id=?"
" AND author=? AND (flags & 1) = 0") == SQLITE_OK) " AND node_id=? AND (flags & 1) = 0") == SQLITE_OK)
{ {
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)ts); sqlite3_bind_int64(stmt, 1, (sqlite3_int64)ts);
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)author); 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; sqlite3_stmt* stmt;
if (si_prep(si, &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) " ORDER BY timestamp, author_signature") != SQLITE_OK)
return; return;
for (uint32_t pos = 0; sqlite3_step(stmt) == SQLITE_ROW; pos++) { 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; sqlite3_stmt* stmt;
if (si_prep(si, &stmt, if (si_prep(si, &stmt,
"DELETE FROM \"%s\"" "DELETE FROM \"%s\""
" WHERE author=? AND (flags & 1) = 0" " WHERE node_id=? AND (flags & 1) = 0"
" AND timestamp<?") != SQLITE_OK) " AND timestamp<?") != SQLITE_OK)
{ {
si->ttl_timer = uasync_set_timeout(si->db_sync->inst->ua, interval, si, db_sync_instance_ttl_cb, "db_sync_ttl"); si->ttl_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; const char* dp = inst->config->global.db_path;
char sp[512]; 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"); else snprintf(sp, sizeof(sp), "/tmp/utun_db_sync");
if (db_sqlite_open(db, sp) != 0) { 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"); 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; } 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, "name=%s id=%016llx", name, (unsigned long long)id); DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "table=%s hash=%016llx", table_name, (unsigned long long)hash);
struct DB_SYNC* db = inst->db_sync; struct DB_SYNC* db = inst->db_sync;
if (!db->enabled || !db->db) return NULL; if (!db->enabled || !db->db) return NULL;
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); 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; } if (si) { DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "instance already exists hash=%016llx", (unsigned long long)hash); return si; }
si = db_instance_alloc(db); si = db_instance_alloc(db);
if (!si) return NULL; if (!si) return NULL;
si->hash = hash; 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]; char sql[512];
snprintf(sql, sizeof(sql), snprintf(sql, sizeof(sql),
"CREATE TABLE IF NOT EXISTS \"%s\" (" "CREATE TABLE IF NOT EXISTS \"%s\" ("
" timestamp INTEGER NOT NULL," " timestamp INTEGER NOT NULL,"
" author INTEGER NOT NULL," " node_id INTEGER NOT NULL,"
" id INTEGER NOT NULL," " id INTEGER NOT NULL,"
" chain_hash BLOB NOT NULL," " chain_hash BLOB NOT NULL,"
" flags INTEGER NOT NULL DEFAULT 0," " 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; } 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), snprintf(sql, sizeof(sql),
"CREATE INDEX IF NOT EXISTS \"idx_%s_ttl\"" "CREATE INDEX IF NOT EXISTS \"idx_%s_ttl\""
" ON \"%s\" (author, timestamp)", " ON \"%s\" (node_id, timestamp)",
si->table_name, si->table_name); si->table_name, si->table_name);
rc = sqlite3_exec(db->db, sql, NULL, NULL, NULL); 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)); 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); 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"); 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; 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, 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", "instance added table=%s 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, table_name, si->table_name, (unsigned long long)hash,
(unsigned long long)si->next_id, peers_found, peers_synced, db_count(si)); (unsigned long long)si->next_id, peers_found, peers_synced, db_count(si));
return si; return si;
} }
@ -1592,7 +1591,7 @@ int db_sync_select(struct DB_SYNC_INSTANCE* si, uint32_t offset, uint32_t limit,
sqlite3_stmt* stmt; sqlite3_stmt* stmt;
if (si_prep(si, &stmt, if (si_prep(si, &stmt,
"SELECT id,timestamp,author,data," "SELECT id,timestamp,node_id,data,"
"author_signature,delivered_peers,delivery_chain" "author_signature,delivered_peers,delivery_chain"
" FROM \"%s\" ORDER BY timestamp, author_signature" " FROM \"%s\" ORDER BY timestamp, author_signature"
" LIMIT ? OFFSET ?") != SQLITE_OK) " LIMIT ? OFFSET ?") != SQLITE_OK)

2
src/db_sync.h

@ -90,7 +90,7 @@ int db_sync_init(struct UTUN_INSTANCE* inst);
void db_sync_destroy(struct UTUN_INSTANCE* inst); void db_sync_destroy(struct UTUN_INSTANCE* inst);
// Instance management // Instance management
struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const char* 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); void db_sync_instance_remove(struct DB_SYNC_INSTANCE* si);
// Data operations (per-instance) // Data operations (per-instance)

70
tools/chatgui/db/db_manager.cpp

@ -46,7 +46,7 @@ QByteArray DbManager::computeChainHash(const QByteArray& prevChain,
QByteArray DbManager::getLastChainHash(const QString& chId) const { QByteArray DbManager::getLastChainHash(const QString& chId) const {
QString tbl = msgTableName(chId); QString tbl = msgTableName(chId);
QByteArray sql = QStringLiteral( 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(); .arg(tbl).toUtf8();
sqlite3_stmt* stmt = prepareOrNull(sql.constData()); sqlite3_stmt* stmt = prepareOrNull(sql.constData());
if (!stmt) return QByteArray(32, '\0'); if (!stmt) return QByteArray(32, '\0');
@ -125,24 +125,39 @@ QList<ChannelRow> DbManager::getChannels() const {
QList<MessageRow> DbManager::getMessages(const QString& chId, int limit) const { QList<MessageRow> DbManager::getMessages(const QString& chId, int limit) const {
QList<MessageRow> list; QList<MessageRow> list;
QString sql = QStringLiteral( QString sql = QStringLiteral(
"SELECT id, node_id, content_type, data, timestamp, datahash, chain_hash," "SELECT id, node_id, data, timestamp, chain_hash, author_signature"
" signature, is_outgoing, is_read" " FROM \"%1\" ORDER BY timestamp, author_signature ASC LIMIT ?").arg(msgTableName(chId));
" FROM \"%1\" ORDER BY timestamp, datahash ASC LIMIT ?").arg(msgTableName(chId));
sqlite3_stmt* stmt = prepareOrNull(sql.toUtf8().constData()); sqlite3_stmt* stmt = prepareOrNull(sql.toUtf8().constData());
if (!stmt) return list; if (!stmt) return list;
sqlite3_bind_int(stmt, 1, limit); sqlite3_bind_int(stmt, 1, limit);
/* helper: parse JSON field from data column: {"n":<id>,"ch":"<id>","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) { while (sqlite3_step(stmt) == SQLITE_ROW) {
MessageRow m; MessageRow m;
m.id = sqlite3_column_int64(stmt, 0); m.id = sqlite3_column_int64(stmt, 0);
m.authorNodeId = (quint64)sqlite3_column_int64(stmt, 1); m.authorNodeId = (quint64)sqlite3_column_int64(stmt, 1);
m.contentType = colText(stmt, 2); QByteArray jdata = colBlob(stmt, 2);
m.data = colBlob(stmt, 3); m.contentType = QString::fromUtf8(parseJsonField(jdata, "ct"));
m.timestamp = sqlite3_column_int64(stmt, 4); if (m.contentType.isEmpty()) m.contentType = "text/plain";
m.datahash = (uint64_t)sqlite3_column_int64(stmt, 5); m.data = parseJsonField(jdata, "d");
m.chainHash = colBlob(stmt, 6); m.timestamp = sqlite3_column_int64(stmt, 3);
m.signature = colBlob(stmt, 7); m.datahash = 0;
m.isOutgoing = sqlite3_column_int(stmt, 8) != 0; m.chainHash = colBlob(stmt, 4);
m.isRead = sqlite3_column_int(stmt, 9) != 0; m.signature = colBlob(stmt, 5);
m.isOutgoing = (m.authorNodeId == myId);
m.isRead = true;
list.append(m); list.append(m);
} }
sqlite3_finalize(stmt); sqlite3_finalize(stmt);
@ -376,7 +391,7 @@ QByteArray DbManager::syncCount(const QString& chId) const {
QByteArray DbManager::syncChainHashAt(const QString& chId, uint32_t pos) const { QByteArray DbManager::syncChainHashAt(const QString& chId, uint32_t pos) const {
QString sql = QStringLiteral( 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)); .arg(msgTableName(chId));
sqlite3_stmt* stmt = prepareOrNull(sql.toUtf8().constData()); sqlite3_stmt* stmt = prepareOrNull(sql.toUtf8().constData());
if (!stmt) return QByteArray(32, '\0'); 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) { QByteArray DbManager::syncCursorOpen(const QString& chId) {
QString sql = QStringLiteral( QString sql = QStringLiteral(
"SELECT timestamp, datahash, node_id, content_type, data" "SELECT timestamp, node_id, data, author_signature"
" FROM \"%1\" ORDER BY timestamp, datahash ASC").arg(msgTableName(chId)); " FROM \"%1\" ORDER BY timestamp, author_signature ASC").arg(msgTableName(chId));
sqlite3_stmt* stmt = prepareOrNull(sql.toUtf8().constData()); sqlite3_stmt* stmt = prepareOrNull(sql.toUtf8().constData());
if (!stmt) { QByteArray r(4, '\0'); return r; } if (!stmt) { QByteArray r(4, '\0'); return r; }
uint32_t id = m_nextCursorId++; 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] */ /* 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); 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(), 1);
uint64_t nid = (uint64_t)sqlite3_column_int64(it.value(), 2); QByteArray jdata = colBlob(it.value(), 2);
QString ct = colText(it.value(), 3); /* parse JSON from data column */
QByteArray data = colBlob(it.value(), 4); 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; QByteArray out;
out.append((const char*)&ts, 8); out.append((const char*)&ts, 8);
@ -417,9 +445,9 @@ QByteArray DbManager::syncCursorNext(uint32_t cursorId) {
QByteArray ctb = ct.toUtf8(); QByteArray ctb = ct.toUtf8();
out.append((char)(uint8_t)ctb.size()); out.append((char)(uint8_t)ctb.size());
out.append(ctb); 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((const char*)&dLen, 4);
out.append(data); out.append(ddata);
return out; return out;
} }

138
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 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 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); int ret = db_sync_insert_signed(si, json, strlen(json), sig, 64, ts);
if (ret != 0) { if (ret != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: db_sync_insert_signed failed ret=%d", CC_ID, ret); 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) */ return;
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);
} }
/* 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); } 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) { void chat_core_ensure_channel_ready(const char* ch_id) {
if (!g_cc.initialized || !ch_id || !ch_id[0]) return; if (!g_cc.initialized || !ch_id || !ch_id[0]) return;
if (si_find(ch_id)) return; if (si_find(ch_id)) return;
char tbl_msg[80]; msg_table_name(ch_id, tbl_msg, sizeof(tbl_msg)); char tbl_msg[80]; msg_table_name(ch_id, tbl_msg, sizeof(tbl_msg));
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 tbl_peers[80]; peers_table_name(ch_id, tbl_peers, sizeof(tbl_peers));
char sql[512];
snprintf(sql, sizeof(sql), snprintf(sql, sizeof(sql),
"CREATE TABLE IF NOT EXISTS \"%s\" (" "CREATE TABLE IF NOT EXISTS \"%s\" ("
" node_id INTEGER PRIMARY KEY," " node_id INTEGER PRIMARY KEY,"
@ -929,12 +882,12 @@ void chat_core_ensure_channel_ready(const char* ch_id) {
uint64_t ch_hash = 0; 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); } { 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) { if (si) {
si_register(si, ch_id); si_register(si, ch_id);
db_sync_set_insert_cb(si, on_db_sync_insert, u_strdup(ch_id)); db_sync_set_insert_cb(si, on_msg_inserted, u_strdup(ch_id));
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: channel ready ch=%s hash=0x%016llx", DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: channel ready ch=%s tbl=%s hash=0x%016llx",
CC_ID, ch_id, (unsigned long long)ch_hash); CC_ID, ch_id, tbl_msg, (unsigned long long)ch_hash);
} else { } else {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: db_sync_instance_add failed for ch=%s", DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: db_sync_instance_add failed for ch=%s",
CC_ID, ch_id); CC_ID, ch_id);
@ -1043,6 +996,13 @@ void chat_core_create_channel_trampoline(void* arg) {
/* ─── db_sync helpers ─── */ /* ─── 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) { static struct DB_SYNC_INSTANCE* si_find(const char* ch_id) {
for (int i=0; i<g_cc.si_count; i++) if (strcmp(g_cc.si_ch_id[i], ch_id)==0) return g_cc.si[i]; for (int i=0; i<g_cc.si_count; i++) if (strcmp(g_cc.si_ch_id[i], ch_id)==0) return g_cc.si[i];
return NULL; return NULL;
@ -1050,67 +1010,9 @@ 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 si_register(struct DB_SYNC_INSTANCE* si, const char* ch_id) {
if (g_cc.si_count>=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; } if (g_cc.si_count>=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++; 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":<uint64>,"ch":"<str>","ct":"<str>","d":"<str>"} */
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 */ uint32_t msg_cnt = db_sync_count(si);
uint8_t dh_buf[32]; SHA256((const uint8_t*)jd_start, jd_len, dh_buf); DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: si_register ch=%s mc=%u", CC_ID, ch_id, msg_cnt);
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);
}
} }
/* ─── сбор статуса (NTP + connections) и отправка в GUI ─── */ /* ─── сбор статуса (NTP + connections) и отправка в GUI ─── */

Loading…
Cancel
Save