From 19a4505a031527b85f4b34987dc8f003c7c0315a Mon Sep 17 00:00:00 2001 From: Evgeny Date: Mon, 6 Jul 2026 16:25:52 +0300 Subject: [PATCH] =?UTF-8?q?chat=5Fsync:=20=D0=BD=D0=BE=D0=B2=D1=8B=D0=B9?= =?UTF-8?q?=20=D0=BF=D1=80=D0=BE=D1=82=D0=BE=D0=BA=D0=BE=D0=BB=20=D1=81?= =?UTF-8?q?=D0=B8=D0=BD=D1=85=D1=80=D0=BE=D0=BD=D0=B8=D0=B7=D0=B0=D1=86?= =?UTF-8?q?=D0=B8=D0=B8=20=D1=87=D0=B0=D1=82=D0=B0=20=D0=B2=D0=BC=D0=B5?= =?UTF-8?q?=D1=81=D1=82=D0=BE=20ChatPropagator?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - gui_bridge.h/cpp — кросс-поточные асинхронные коллбэки uasync↔GUI - chat_sync.h/c — протокол синхронизации (INIT_SYNC/RESP, SEND_DATA, PUSH/ACK, SYNC_DONE) - DbManager полностью переписан: per-channel таблицы msg_ + peers_, 4 ключа в channels, signature, protocol/rtt в node_addresses, sync-методы - ChatPropagator удалён, заменён на chat_sync - db_sync_stub.c — заглушки для сборки chatgui без LMDB-версии db_sync - db_schema.md актуализирован под новую схему --- AGENTS.md | 2 + tools/chatgui/CMakeLists.txt | 2 +- tools/chatgui/db/db_manager.cpp | 1022 +++++++++++-------- tools/chatgui/db/db_manager.h | 129 ++- tools/chatgui/doc/db_schema.md | 281 +++-- tools/chatgui/libutun/CMakeLists.txt | 5 +- tools/chatgui/src/channellist.cpp | 2 +- tools/chatgui/src/mainwindow.cpp | 10 +- tools/chatgui/src/mainwindow.h | 2 - tools/chatgui/src/messagelist.cpp | 15 +- tools/chatgui/src/messagelist.h | 3 - tools/chatgui/transport/chat_sync.c | 660 ++++++++++++ tools/chatgui/transport/chat_sync.h | 53 + tools/chatgui/transport/db_sync_stub.c | 31 + tools/chatgui/transport/gui_bridge.h | 60 ++ tools/chatgui/transport/gui_bridge_impl.cpp | 252 +++++ tools/chatgui/transport/utun_node.cpp | 7 + 17 files changed, 1977 insertions(+), 559 deletions(-) create mode 100644 tools/chatgui/transport/chat_sync.c create mode 100644 tools/chatgui/transport/chat_sync.h create mode 100644 tools/chatgui/transport/db_sync_stub.c create mode 100644 tools/chatgui/transport/gui_bridge.h create mode 100644 tools/chatgui/transport/gui_bridge_impl.cpp diff --git a/AGENTS.md b/AGENTS.md index fbabab36..13cae11a 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -381,6 +381,8 @@ SOCKET=14, CONTROL=15, DUMP=16, TRAFFIC=17, DEBUG=18, GENERAL=19, NAT=20 Chatgui — десктопный GUI-чат на Qt 5, отдельный проект внутри репозитория. Не связан с autotools-сборкой utun. +В chat gui интегрированы библиотеки utun. Чат и библиотеки работают в разных потоках. Поэтому нужно использовать семафоры, сокеты или другие механизмы синхронизации (uasync_post, uasync_memsync, uasync_get_wakeup_fd) + ### Технологии - **Язык:** C++17 - **Фреймворк:** Qt 5.15+ (Widgets + Network) diff --git a/tools/chatgui/CMakeLists.txt b/tools/chatgui/CMakeLists.txt index 53484fac..b6f79c80 100644 --- a/tools/chatgui/CMakeLists.txt +++ b/tools/chatgui/CMakeLists.txt @@ -63,9 +63,9 @@ add_executable(chatgui src/invitedialog.cpp transport/msg_client.cpp transport/crypto.cpp - transport/chat_propagator.cpp transport/utun_node.cpp transport/node_config.cpp + transport/gui_bridge_impl.cpp db/db_manager.cpp db/sqlite3.c resources/chatgui.qrc diff --git a/tools/chatgui/db/db_manager.cpp b/tools/chatgui/db/db_manager.cpp index 2759c8db..c726b6dd 100644 --- a/tools/chatgui/db/db_manager.cpp +++ b/tools/chatgui/db/db_manager.cpp @@ -3,6 +3,62 @@ #include #include #include +#include + +/* ── helpers ── */ + +static const char kZero64[64] = {}; +static const char kZero32[32] = {}; + +QString DbManager::sanitizeTableName(const QString& chId) { + QString s; + s.reserve(chId.size()); + for (int i = 0; i < chId.size(); i++) { + QChar c = chId[i]; + if (c.isLetterOrNumber() || c == '_') s += c; + else s += '_'; + } + return s; +} + +QString DbManager::msgTableName(const QString& chId) const { + return QStringLiteral("msg_") + sanitizeTableName(chId); +} +QString DbManager::peersTableName(const QString& chId) const { + return QStringLiteral("peers_") + sanitizeTableName(chId); +} + +uint64_t DbManager::computeDatahash(const QByteArray& data) { + uint8_t hash[32]; + SHA256((const unsigned char*)data.constData(), data.size(), hash); + uint64_t dh; memcpy(&dh, hash, 8); return dh; +} + +QByteArray DbManager::computeChainHash(const QByteArray& prevChain, + qint64 timestamp, uint64_t datahash) { + uint8_t buf[32 + 8 + 8]; + memcpy(buf, prevChain.constData(), 32); + memcpy(buf + 32, ×tamp, 8); + memcpy(buf + 40, &datahash, 8); + uint8_t hash[32]; + SHA256(buf, 48, hash); + return QByteArray((const char*)hash, 32); +} + +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") + .arg(tbl).toUtf8(); + sqlite3_stmt* stmt = prepareOrNull(sql.constData()); + if (!stmt) return QByteArray(32, '\0'); + QByteArray r = (sqlite3_step(stmt) == SQLITE_ROW) ? colBlob(stmt, 0) + : QByteArray(32, '\0'); + sqlite3_finalize(stmt); + return r.size() == 32 ? r : QByteArray(32, '\0'); +} + +/* ── lifecycle ── */ DbManager::DbManager(const QString& path, QObject* parent) : QObject(parent) @@ -10,25 +66,24 @@ DbManager::DbManager(const QString& path, QObject* parent) int rc = sqlite3_open(path.toUtf8().constData(), &m_db); if (rc != SQLITE_OK) { GUI_ERROR("Cannot open database: %s", sqlite3_errmsg(m_db)); - sqlite3_close(m_db); - m_db = nullptr; - return; + sqlite3_close(m_db); m_db = nullptr; return; } sqlite3_exec(m_db, "PRAGMA journal_mode=WAL", nullptr, nullptr, nullptr); - sqlite3_exec(m_db, "PRAGMA foreign_keys=ON", nullptr, nullptr, nullptr); + sqlite3_exec(m_db, "PRAGMA foreign_keys=ON", nullptr, nullptr, nullptr); m_myNodeId = 100; initTables(); if (getChannels().isEmpty()) seedStubData(); GUI_INFO("Database opened: %s (my_node_id=%llu)", path.toUtf8().constData(), - static_cast(m_myNodeId)); + (unsigned long long)m_myNodeId); } DbManager::~DbManager() { + for (auto it = m_cursors.begin(); it != m_cursors.end(); ++it) + sqlite3_finalize(it.value()); + m_cursors.clear(); if (m_db) { sqlite3_close(m_db); m_db = nullptr; } } -/* ======================================================================== */ - sqlite3_stmt* DbManager::prepareOrNull(const char* sql) const { sqlite3_stmt* stmt = nullptr; if (sqlite3_prepare_v2(m_db, sql, -1, &stmt, nullptr) != SQLITE_OK) { @@ -49,7 +104,7 @@ QByteArray DbManager::colBlob(sqlite3_stmt* stmt, int col) { return QByteArray(static_cast(sqlite3_column_blob(stmt, col)), len); } -/* ======================================================================== */ +/* ── init ── */ void DbManager::initTables() { const char* sql = @@ -64,25 +119,27 @@ void DbManager::initTables() { " updated_at INTEGER DEFAULT (unixepoch())" ");" - "CREATE TABLE IF NOT EXISTS my_nodes (" - " node_id INTEGER PRIMARY KEY," - " name TEXT NOT NULL," - " x25519_pubkey BLOB NOT NULL," - " ed25519_pubkey BLOB," - " created_at INTEGER DEFAULT (unixepoch())," - " updated_at INTEGER DEFAULT (unixepoch())" + "CREATE TABLE IF NOT EXISTS nodes (" + " node_id INTEGER PRIMARY KEY," + " name TEXT," + " x25519_pubkey BLOB NOT NULL," + " ed25519_pubkey BLOB," + " last_seen_at INTEGER," + " created_at INTEGER DEFAULT (unixepoch())" ");" - "CREATE TABLE IF NOT EXISTS nodes (" - " node_id INTEGER PRIMARY KEY," - " name TEXT," - " x25519_pubkey BLOB NOT NULL," - " ed25519_pubkey BLOB," - " is_online INTEGER DEFAULT 0," - " last_seen_at INTEGER DEFAULT (unixepoch())," - " created_at INTEGER DEFAULT (unixepoch())," - " updated_at INTEGER DEFAULT (unixepoch())" + "CREATE TABLE IF NOT EXISTS node_addresses (" + " id INTEGER PRIMARY KEY AUTOINCREMENT," + " node_id INTEGER NOT NULL REFERENCES nodes(node_id) ON DELETE CASCADE," + " family INTEGER NOT NULL CHECK(family IN (4, 6))," + " protocol INTEGER NOT NULL DEFAULT 1," + " address BLOB NOT NULL," + " port INTEGER NOT NULL CHECK(port > 0 AND port <= 65535)," + " rtt INTEGER," + " is_nat INTEGER DEFAULT 0," + " created_at INTEGER DEFAULT (unixepoch())" ");" + "CREATE INDEX IF NOT EXISTS idx_na_node ON node_addresses(node_id);" "CREATE TABLE IF NOT EXISTS accounts (" " node_id INTEGER PRIMARY KEY REFERENCES nodes(node_id)," @@ -94,161 +151,372 @@ void DbManager::initTables() { ");" "CREATE TABLE IF NOT EXISTS channels (" - " id INTEGER PRIMARY KEY AUTOINCREMENT," - " channel_id TEXT NOT NULL UNIQUE," - " name TEXT NOT NULL," - " owner_node_id INTEGER," - " is_dm INTEGER DEFAULT 0," - " last_message TEXT," - " last_msg_at INTEGER," - " created_at INTEGER DEFAULT (unixepoch())" - ");" - "CREATE UNIQUE INDEX IF NOT EXISTS idx_channel_cid ON channels(channel_id);" - - "CREATE TABLE IF NOT EXISTS channel_members (" - " channel_id TEXT NOT NULL REFERENCES channels(channel_id) ON DELETE CASCADE," - " node_id INTEGER NOT NULL," - " joined_at INTEGER DEFAULT (unixepoch())," - " PRIMARY KEY (channel_id, node_id)" - ");" - "CREATE INDEX IF NOT EXISTS idx_cm_node ON channel_members(node_id);" - - "CREATE TABLE IF NOT EXISTS messages (" - " id INTEGER PRIMARY KEY AUTOINCREMENT," - " channel_id TEXT NOT NULL REFERENCES channels(channel_id) ON DELETE CASCADE," - " author_node_id INTEGER NOT NULL," - " content TEXT NOT NULL CHECK(length(content) <= 4096)," - " content_type TEXT DEFAULT 'text/plain'," - " timestamp INTEGER NOT NULL," - " signature BLOB NOT NULL CHECK(length(signature) == 64)," - " is_outgoing INTEGER DEFAULT 0," - " is_read INTEGER DEFAULT 0," - " created_at INTEGER DEFAULT (unixepoch())" - ");" - "CREATE INDEX IF NOT EXISTS idx_msg_channel_time ON messages(channel_id, timestamp);" - "CREATE INDEX IF NOT EXISTS idx_msg_author ON messages(author_node_id);" - "CREATE UNIQUE INDEX IF NOT EXISTS idx_msg_dedup " - " ON messages(channel_id, author_node_id, timestamp);" - - "CREATE TABLE IF NOT EXISTS reactions (" - " message_id INTEGER NOT NULL REFERENCES messages(id) ON DELETE CASCADE," - " node_id INTEGER NOT NULL," - " emoji TEXT NOT NULL," - " signature BLOB NOT NULL CHECK(length(signature) == 64)," - " created_at INTEGER DEFAULT (unixepoch())," - " PRIMARY KEY (message_id, node_id, emoji)" + " channel_id TEXT PRIMARY KEY," + " name TEXT NOT NULL," + " owner_node_id INTEGER," + " is_dm INTEGER DEFAULT 0," + " x25519_pubkey BLOB NOT NULL," + " x25519_privkey BLOB," + " ed25519_pubkey BLOB NOT NULL," + " ed25519_privkey BLOB," + " signature BLOB NOT NULL," + " last_read_msg_id INTEGER," + " last_pos_msg_id INTEGER," + " created_at INTEGER DEFAULT (unixepoch())" ");" - "CREATE TABLE IF NOT EXISTS attachments (" - " id INTEGER PRIMARY KEY AUTOINCREMENT," - " message_id INTEGER NOT NULL REFERENCES messages(id) ON DELETE CASCADE," - " filename TEXT NOT NULL," - " content_type TEXT NOT NULL," - " size INTEGER NOT NULL," - " sha256 BLOB NOT NULL," - " file_path TEXT," - " signature BLOB NOT NULL CHECK(length(signature) == 64)," - " is_upload INTEGER DEFAULT 0," - " created_at INTEGER DEFAULT (unixepoch())" - ");" - "CREATE INDEX IF NOT EXISTS idx_att_msg ON attachments(message_id);" - "CREATE TABLE IF NOT EXISTS ui_state (" " key TEXT PRIMARY KEY," " value TEXT" - ");" - - "CREATE TABLE IF NOT EXISTS bucket_meta (" - " channel_id TEXT NOT NULL," - " day_unix INTEGER NOT NULL," - " hour INTEGER NOT NULL," - " msg_count INTEGER DEFAULT 0," - " sig_hash BLOB DEFAULT ''," - " PRIMARY KEY (channel_id, day_unix, hour)" - ");" - - "CREATE TABLE IF NOT EXISTS node_addresses (" - " id INTEGER PRIMARY KEY AUTOINCREMENT," - " node_id INTEGER NOT NULL REFERENCES nodes(node_id) ON DELETE CASCADE," - " family INTEGER NOT NULL CHECK(family IN (4, 6))," - " address BLOB NOT NULL," - " port INTEGER NOT NULL CHECK(port > 0 AND port <= 65535)," - " is_nat INTEGER DEFAULT 0," - " created_at INTEGER DEFAULT (unixepoch())," - " updated_at INTEGER DEFAULT (unixepoch())" - ");" - "CREATE INDEX IF NOT EXISTS idx_na_node ON node_addresses(node_id);"; + ");"; char* err = nullptr; if (sqlite3_exec(m_db, sql, nullptr, nullptr, &err) != SQLITE_OK) { GUI_ERROR("initTables: %s", err); sqlite3_free(err); } else { - GUI_INFO("Database tables initialized (11 tables)"); + GUI_INFO("Database tables initialized"); } } -/* ======================================================================== */ +/* ── channels ── */ QList DbManager::getChannels() const { QList list; sqlite3_stmt* stmt = prepareOrNull( - "SELECT c.channel_id, c.name, c.is_dm, c.last_message, c.last_msg_at" - " FROM channels c JOIN channel_members cm ON c.channel_id = cm.channel_id" - " WHERE cm.node_id = ? ORDER BY c.last_msg_at DESC"); + "SELECT channel_id, name, owner_node_id, is_dm," + " x25519_pubkey, x25519_privkey, ed25519_pubkey, ed25519_privkey, signature," + " last_read_msg_id, last_pos_msg_id, created_at" + " FROM channels ORDER BY created_at ASC"); if (!stmt) return list; - sqlite3_bind_int64(stmt, 1, static_cast(m_myNodeId)); while (sqlite3_step(stmt) == SQLITE_ROW) { ChannelRow r; - r.channelId = colText(stmt, 0); - r.name = colText(stmt, 1); - r.isDm = sqlite3_column_int(stmt, 2) != 0; - r.lastMessage = colText(stmt, 3); - r.lastMsgAt = sqlite3_column_int64(stmt, 4); + r.channelId = colText(stmt, 0); + r.name = colText(stmt, 1); + r.ownerNodeId = (quint64)sqlite3_column_int64(stmt, 2); + r.isDm = sqlite3_column_int(stmt, 3) != 0; + r.x25519_pubkey = colBlob(stmt, 4); + r.x25519_privkey = colBlob(stmt, 5); + r.ed25519_pubkey = colBlob(stmt, 6); + r.ed25519_privkey = colBlob(stmt, 7); + r.signature = colBlob(stmt, 8); + r.lastReadMsgId = sqlite3_column_int64(stmt, 9); + r.lastPosMsgId = sqlite3_column_int64(stmt, 10); + r.createdAt = sqlite3_column_int64(stmt, 11); list.append(r); } sqlite3_finalize(stmt); return list; } -QList DbManager::getMessages(const QString& channelId, int limit) const { - QList list; +bool DbManager::createChannel(const QString& chId, const QString& name, + quint64 ownerNodeId, + const QByteArray& x25519pub, const QByteArray& x25519priv, + const QByteArray& ed25519pub, const QByteArray& ed25519priv) { + QByteArray sig = kZero64; Q_UNUSED(sig); + sqlite3_stmt* stmt = prepareOrNull( - "SELECT id, author_node_id, content, content_type, timestamp," - " signature, is_outgoing, is_read" - " FROM messages WHERE channel_id = ? ORDER BY timestamp ASC LIMIT ?"); + "INSERT INTO channels(channel_id,name,owner_node_id,is_dm," + " x25519_pubkey,x25519_privkey,ed25519_pubkey,ed25519_privkey,signature)" + " VALUES(?,?,?,0,?,?,?,?,?)"); + if (!stmt) return false; + sqlite3_bind_text(stmt, 1, chId.toUtf8().constData(), -1, SQLITE_STATIC); + sqlite3_bind_text(stmt, 2, name.toUtf8().constData(), -1, SQLITE_STATIC); + sqlite3_bind_int64(stmt, 3, (sqlite3_int64)ownerNodeId); + sqlite3_bind_blob(stmt, 4, x25519pub.isEmpty() ? kZero32 : x25519pub.constData(), + x25519pub.isEmpty() ? 32 : x25519pub.size(), SQLITE_STATIC); + sqlite3_bind_blob(stmt, 5, x25519priv.isEmpty() ? nullptr : x25519priv.constData(), + x25519priv.size(), SQLITE_STATIC); + sqlite3_bind_blob(stmt, 6, ed25519pub.isEmpty() ? kZero32 : ed25519pub.constData(), + ed25519pub.isEmpty() ? 32 : ed25519pub.size(), SQLITE_STATIC); + sqlite3_bind_blob(stmt, 7, ed25519priv.isEmpty() ? nullptr : ed25519priv.constData(), + ed25519priv.size(), SQLITE_STATIC); + sqlite3_bind_blob(stmt, 8, sig.constData(), sig.size(), SQLITE_STATIC); + bool ok = sqlite3_step(stmt) == SQLITE_DONE; + sqlite3_finalize(stmt); + if (!ok) { GUI_ERROR("createChannel %s failed: %s", chId.toUtf8().constData(), sqlite3_errmsg(m_db)); return false; } + + /* per-channel message table */ + QString msgt = msgTableName(chId), pt = peersTableName(chId); + QString csql = QStringLiteral( + "CREATE TABLE IF NOT EXISTS \"%1\" (" + " 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 NOT NULL," + " chain_hash BLOB NOT NULL," + " signature BLOB," + " is_outgoing INTEGER DEFAULT 0," + " is_read INTEGER DEFAULT 0," + " sync_flags INTEGER DEFAULT 0," + " UNIQUE(timestamp, datahash))").arg(msgt); + sqlite3_exec(m_db, csql.toUtf8().constData(), nullptr, nullptr, nullptr); + csql = QStringLiteral("CREATE INDEX IF NOT EXISTS idx_%1_ts ON \"%1\"(timestamp, datahash)").arg(msgt); + sqlite3_exec(m_db, csql.toUtf8().constData(), nullptr, nullptr, nullptr); + + /* per-channel peers table */ + csql = QStringLiteral( + "CREATE TABLE IF NOT EXISTS \"%1\" (" + " node_id INTEGER NOT NULL," + " join_sig BLOB NOT NULL," + " creator_sig BLOB," + " comment TEXT," + " joined_at INTEGER DEFAULT (unixepoch())," + " PRIMARY KEY (node_id))").arg(pt); + sqlite3_exec(m_db, csql.toUtf8().constData(), nullptr, nullptr, nullptr); + + GUI_INFO("Channel created: %s (msg=%s, peers=%s)", + chId.toUtf8().constData(), msgt.toUtf8().constData(), pt.toUtf8().constData()); + return true; +} + +bool DbManager::deleteChannel(const QString& chId) { + QString msgt = msgTableName(chId), pt = peersTableName(chId); + sqlite3_exec(m_db, QStringLiteral("DROP TABLE IF EXISTS \"%1\"").arg(msgt).toUtf8().constData(), nullptr, nullptr, nullptr); + sqlite3_exec(m_db, QStringLiteral("DROP TABLE IF EXISTS \"%1\"").arg(pt).toUtf8().constData(), nullptr, nullptr, nullptr); + sqlite3_stmt* stmt = prepareOrNull("DELETE FROM channels WHERE channel_id=?"); + if (!stmt) return false; + sqlite3_bind_text(stmt, 1, chId.toUtf8().constData(), -1, SQLITE_STATIC); + sqlite3_step(stmt); + sqlite3_finalize(stmt); + return true; +} + +/* ── messages (UI) ── */ + +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)); + sqlite3_stmt* stmt = prepareOrNull(sql.toUtf8().constData()); if (!stmt) return list; - sqlite3_bind_text(stmt, 1, channelId.toUtf8().constData(), -1, SQLITE_STATIC); - sqlite3_bind_int(stmt, 2, limit); + sqlite3_bind_int(stmt, 1, limit); while (sqlite3_step(stmt) == SQLITE_ROW) { MessageRow m; m.id = sqlite3_column_int64(stmt, 0); - m.channelId = channelId; - m.authorNodeId = static_cast(sqlite3_column_int64(stmt, 1)); - m.content = colText(stmt, 2); - m.contentType = colText(stmt, 3); + 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.signature = colBlob(stmt, 5); - m.isOutgoing = sqlite3_column_int(stmt, 6) != 0; - m.isRead = sqlite3_column_int(stmt, 7) != 0; + 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; list.append(m); } sqlite3_finalize(stmt); return list; } +bool DbManager::insertMessage(const QString& chId, quint64 nodeId, + const QString& contentType, const QByteArray& data, + qint64 timestamp, bool isOutgoing, + uint64_t* outDatahash, QByteArray* outChainHash) { + uint64_t dh = computeDatahash(data); + QByteArray prevCh = getLastChainHash(chId); + QByteArray ch = computeChainHash(prevCh, timestamp, dh); + + QString sql = QStringLiteral( + "INSERT OR IGNORE INTO \"%1\"" + " (node_id, content_type, data, timestamp, datahash, chain_hash," + " signature, is_outgoing, is_read)" + " VALUES(?,?,?,?,?,?,?,?,1)").arg(msgTableName(chId)); + sqlite3_stmt* stmt = prepareOrNull(sql.toUtf8().constData()); + if (!stmt) return false; + sqlite3_bind_int64(stmt, 1, (sqlite3_int64)nodeId); + sqlite3_bind_text(stmt, 2, contentType.toUtf8().constData(), -1, SQLITE_STATIC); + sqlite3_bind_blob(stmt, 3, data.constData(), data.size(), SQLITE_STATIC); + sqlite3_bind_int64(stmt, 4, timestamp); + sqlite3_bind_int64(stmt, 5, (sqlite3_int64)dh); + sqlite3_bind_blob(stmt, 6, ch.constData(), 32, SQLITE_STATIC); + sqlite3_bind_blob(stmt, 7, kZero64, 64, SQLITE_STATIC); + sqlite3_bind_int(stmt, 8, isOutgoing ? 1 : 0); + bool ok = sqlite3_step(stmt) == SQLITE_DONE; + sqlite3_finalize(stmt); + + if (ok) { + if (outDatahash) *outDatahash = dh; + if (outChainHash) *outChainHash = ch; + } + return ok; +} + +/* ── peers ── */ + +bool DbManager::addPeer(const QString& chId, quint64 nodeId, + const QByteArray& joinSig, const QByteArray& creatorSig, + const QString& comment) { + QString sql = QStringLiteral( + "INSERT OR REPLACE INTO \"%1\"(node_id, join_sig, creator_sig, comment)" + " VALUES(?,?,?,?)").arg(peersTableName(chId)); + sqlite3_stmt* stmt = prepareOrNull(sql.toUtf8().constData()); + if (!stmt) return false; + sqlite3_bind_int64(stmt, 1, (sqlite3_int64)nodeId); + sqlite3_bind_blob(stmt, 2, joinSig.constData(), joinSig.size(), SQLITE_STATIC); + sqlite3_bind_blob(stmt, 3, creatorSig.isEmpty() ? nullptr : creatorSig.constData(), + creatorSig.size(), SQLITE_STATIC); + if (!comment.isEmpty()) + sqlite3_bind_text(stmt, 4, comment.toUtf8().constData(), -1, SQLITE_STATIC); + else + sqlite3_bind_null(stmt, 4); + bool ok = sqlite3_step(stmt) == SQLITE_DONE; + sqlite3_finalize(stmt); + return ok; +} + +bool DbManager::removePeer(const QString& chId, quint64 nodeId) { + QString sql = QStringLiteral( + "DELETE FROM \"%1\" WHERE node_id=?").arg(peersTableName(chId)); + sqlite3_stmt* stmt = prepareOrNull(sql.toUtf8().constData()); + if (!stmt) return false; + sqlite3_bind_int64(stmt, 1, (sqlite3_int64)nodeId); + sqlite3_step(stmt); + sqlite3_finalize(stmt); + return true; +} + +/* ── nodes ── */ + +bool DbManager::upsertNode(quint64 nodeId, const QString& name, + const QByteArray& x25519pub, const QByteArray& ed25519pub) { + sqlite3_stmt* stmt = prepareOrNull( + "INSERT OR REPLACE INTO nodes(node_id, name, x25519_pubkey, ed25519_pubkey, last_seen_at)" + " VALUES(?,?,?,?,unixepoch())"); + if (!stmt) return false; + sqlite3_bind_int64(stmt, 1, (sqlite3_int64)nodeId); + sqlite3_bind_text(stmt, 2, name.toUtf8().constData(), -1, SQLITE_STATIC); + sqlite3_bind_blob(stmt, 3, x25519pub.constData(), x25519pub.size(), SQLITE_STATIC); + sqlite3_bind_blob(stmt, 4, ed25519pub.isEmpty() ? nullptr : ed25519pub.constData(), + ed25519pub.size(), SQLITE_STATIC); + bool ok = sqlite3_step(stmt) == SQLITE_DONE; + sqlite3_finalize(stmt); + return ok; +} + +bool DbManager::addNodeAddress(quint64 nodeId, int family, int protocol, + const QByteArray& addr, quint16 port, bool isNat) { + sqlite3_stmt* stmt = prepareOrNull( + "INSERT OR IGNORE INTO node_addresses(node_id, family, protocol, address, port, is_nat)" + " VALUES(?,?,?,?,?,?)"); + if (!stmt) return false; + sqlite3_bind_int64(stmt, 1, (sqlite3_int64)nodeId); + sqlite3_bind_int(stmt, 2, family); + sqlite3_bind_int(stmt, 3, protocol); + sqlite3_bind_blob(stmt, 4, addr.constData(), addr.size(), SQLITE_STATIC); + sqlite3_bind_int(stmt, 5, port); + sqlite3_bind_int(stmt, 6, isNat ? 1 : 0); + bool ok = sqlite3_step(stmt) == SQLITE_DONE; + sqlite3_finalize(stmt); + return ok; +} + +QByteArray DbManager::loadNodeInfo(quint64 nodeId) const { + /* binary format: [x25519:32][ed25519:32][addr_count:1]([family:1][proto:1][addr_len:1][addr:4|16][port:2][rtt:2])... */ + sqlite3_stmt* stmt = prepareOrNull( + "SELECT x25519_pubkey, ed25519_pubkey FROM nodes WHERE node_id=?"); + if (!stmt) return {}; + sqlite3_bind_int64(stmt, 1, (sqlite3_int64)nodeId); + if (sqlite3_step(stmt) != SQLITE_ROW) { sqlite3_finalize(stmt); return {}; } + QByteArray x25519 = colBlob(stmt, 0); + QByteArray ed25519 = colBlob(stmt, 1); + sqlite3_finalize(stmt); + + QByteArray result; + result.append(x25519.isEmpty() ? QByteArray(32, '\0') : x25519); + result.append(ed25519.isEmpty() ? QByteArray(32, '\0') : ed25519); + + /* запросить адреса */ + stmt = prepareOrNull( + "SELECT family, protocol, address, port, rtt FROM node_addresses WHERE node_id=?"); + if (!stmt) { result.append('\0'); return result; } + sqlite3_bind_int64(stmt, 1, (sqlite3_int64)nodeId); + QByteArray addrsPart; + uint8_t cnt = 0; + while (sqlite3_step(stmt) == SQLITE_ROW && cnt < 255) { + int family = sqlite3_column_int(stmt, 0); + int proto = sqlite3_column_int(stmt, 1); + QByteArray addr = colBlob(stmt, 2); + uint16_t port = (uint16_t)sqlite3_column_int(stmt, 3); + int16_t rtt = (int16_t)(sqlite3_column_int(stmt, 4)); + addrsPart.append((char)family); + addrsPart.append((char)proto); + addrsPart.append((char)addr.size()); + addrsPart.append(addr); + addrsPart.append((const char*)&port, 2); + addrsPart.append((const char*)&rtt, 2); + cnt++; + } + sqlite3_finalize(stmt); + result.append((char)cnt); + result.append(addrsPart); + return result; +} + +QList DbManager::getOnlinePeers() const { + QList peers; + /* без is_online: получаем всех кроме себя */ + sqlite3_stmt* stmt = prepareOrNull("SELECT node_id FROM nodes WHERE node_id != ?"); + if (!stmt) return peers; + sqlite3_bind_int64(stmt, 1, (sqlite3_int64)m_myNodeId); + while (sqlite3_step(stmt) == SQLITE_ROW) + peers.append((quint64)sqlite3_column_int64(stmt, 0)); + sqlite3_finalize(stmt); + return peers; +} + +QByteArray DbManager::getNodeEdPub(quint64 nodeId) const { + sqlite3_stmt* stmt = prepareOrNull("SELECT ed25519_pubkey FROM nodes WHERE node_id=?"); + if (!stmt) return {}; + sqlite3_bind_int64(stmt, 1, (sqlite3_int64)nodeId); + QByteArray r; + if (sqlite3_step(stmt) == SQLITE_ROW) r = colBlob(stmt, 0); + sqlite3_finalize(stmt); + return r; +} + +QList DbManager::getOnlinePeerAddresses(quint64 excludeNodeId) const { + QList list; + sqlite3_stmt* stmt = prepareOrNull( + "SELECT na.node_id, na.family, na.protocol, na.address, na.port, na.rtt, na.is_nat" + " FROM node_addresses na JOIN nodes n ON na.node_id = n.node_id" + " WHERE na.is_nat=0 AND na.node_id != ?" + " ORDER BY RANDOM()"); + if (!stmt) return list; + sqlite3_bind_int64(stmt, 1, (sqlite3_int64)excludeNodeId); + while (sqlite3_step(stmt) == SQLITE_ROW) { + NodeAddr a; + a.nodeId = (quint64)sqlite3_column_int64(stmt, 0); + a.family = sqlite3_column_int(stmt, 1); + a.protocol = sqlite3_column_int(stmt, 2); + a.address = colBlob(stmt, 3); + a.port = (quint16)sqlite3_column_int(stmt, 4); + a.rtt = sqlite3_column_int(stmt, 5); + a.isNat = sqlite3_column_int(stmt, 6) != 0; + list.append(a); + } + sqlite3_finalize(stmt); + return list; +} + +/* ── accounts ── */ + QList DbManager::getAccounts(bool contactsOnly) const { QList list; const char* sql = contactsOnly ? "SELECT node_id, display_name, avatar_color, avatar_letter, is_contact" - " FROM accounts WHERE is_contact = 1 ORDER BY display_name" + " FROM accounts WHERE is_contact=1 ORDER BY display_name" : "SELECT node_id, display_name, avatar_color, avatar_letter, is_contact" " FROM accounts ORDER BY display_name"; sqlite3_stmt* stmt = prepareOrNull(sql); if (!stmt) return list; while (sqlite3_step(stmt) == SQLITE_ROW) { AccountRow a; - a.nodeId = static_cast(sqlite3_column_int64(stmt, 0)); + a.nodeId = (quint64)sqlite3_column_int64(stmt, 0); a.displayName = colText(stmt, 1); a.avatarColor = colText(stmt, 2); a.avatarLetter = colText(stmt, 3); @@ -263,11 +531,11 @@ AccountRow DbManager::getAccount(quint64 nodeId) const { AccountRow a; sqlite3_stmt* stmt = prepareOrNull( "SELECT node_id, display_name, avatar_color, avatar_letter, is_contact" - " FROM accounts WHERE node_id = ?"); + " FROM accounts WHERE node_id=?"); if (!stmt) return a; - sqlite3_bind_int64(stmt, 1, static_cast(nodeId)); + sqlite3_bind_int64(stmt, 1, (sqlite3_int64)nodeId); if (sqlite3_step(stmt) == SQLITE_ROW) { - a.nodeId = static_cast(sqlite3_column_int64(stmt, 0)); + a.nodeId = (quint64)sqlite3_column_int64(stmt, 0); a.displayName = colText(stmt, 1); a.avatarColor = colText(stmt, 2); a.avatarLetter = colText(stmt, 3); @@ -277,200 +545,216 @@ AccountRow DbManager::getAccount(quint64 nodeId) const { return a; } -QByteArray DbManager::getNodeEdPub(quint64 nodeId) const { - sqlite3_stmt* stmt = prepareOrNull( - "SELECT ed25519_pubkey FROM nodes WHERE node_id=?"); +/* ── ui_state ── */ + +QString DbManager::getUiState(const QString& key) const { + sqlite3_stmt* stmt = prepareOrNull("SELECT value FROM ui_state WHERE key=?"); if (!stmt) return {}; - sqlite3_bind_int64(stmt, 1, static_cast(nodeId)); - QByteArray r; - if (sqlite3_step(stmt) == SQLITE_ROW) - r = colBlob(stmt, 0); + sqlite3_bind_text(stmt, 1, key.toUtf8().constData(), -1, SQLITE_STATIC); + QString val; + if (sqlite3_step(stmt) == SQLITE_ROW) val = colText(stmt, 0); sqlite3_finalize(stmt); - return r; + return val; } -QList DbManager::getBucketMeta(const QString& channelId) const { - QList list; +void DbManager::setUiState(const QString& key, const QString& value) { sqlite3_stmt* stmt = prepareOrNull( - "SELECT day_unix, hour, msg_count, sig_hash" - " FROM bucket_meta WHERE channel_id=? ORDER BY day_unix, hour"); - if (!stmt) return list; - sqlite3_bind_text(stmt, 1, channelId.toUtf8().constData(), -1, SQLITE_STATIC); - while (sqlite3_step(stmt) == SQLITE_ROW) { - BucketMeta b; - b.day = sqlite3_column_int64(stmt, 0); - b.hour = sqlite3_column_int(stmt, 1); - b.count = sqlite3_column_int(stmt, 2); - b.hash = colBlob(stmt, 3); - list.append(b); + "INSERT OR REPLACE INTO ui_state(key,value) VALUES(?,?)"); + if (!stmt) return; + sqlite3_bind_text(stmt, 1, key.toUtf8().constData(), -1, SQLITE_STATIC); + sqlite3_bind_text(stmt, 2, value.toUtf8().constData(), -1, SQLITE_STATIC); + sqlite3_step(stmt); sqlite3_finalize(stmt); +} + +/* ── sync operations (called from gui_bridge, in GUI thread) ── */ + +QByteArray DbManager::syncListChannels() const { + /* binary: [count:2]([chId_len:1][chId:var])... */ + QList chs = getChannels(); + QByteArray out; + uint16_t cnt = (uint16_t)chs.size(); + out.append((const char*)&cnt, 2); + for (auto& ch : chs) { + QByteArray id = ch.channelId.toUtf8(); + out.append((char)(uint8_t)id.size()); + out.append(id); } - sqlite3_finalize(stmt); - return list; + return out; } -void DbManager::setBucketMeta(const QString& channelId, qint64 day, int hour, - int count, const QByteArray& hash) { - sqlite3_stmt* stmt = prepareOrNull( - "INSERT OR REPLACE INTO bucket_meta(channel_id,day_unix,hour,msg_count,sig_hash)" - " VALUES(?,?,?,?,?)"); - if (!stmt) return; - sqlite3_bind_text(stmt, 1, channelId.toUtf8().constData(), -1, SQLITE_STATIC); - sqlite3_bind_int64(stmt, 2, day); - sqlite3_bind_int(stmt, 3, hour); - sqlite3_bind_int(stmt, 4, count); - sqlite3_bind_blob(stmt, 5, hash.constData(), hash.size(), SQLITE_STATIC); - sqlite3_step(stmt); +QByteArray DbManager::syncCount(const QString& chId) const { + QString sql = QStringLiteral("SELECT COUNT(*) FROM \"%1\"").arg(msgTableName(chId)); + sqlite3_stmt* stmt = prepareOrNull(sql.toUtf8().constData()); + if (!stmt) return {}; + uint32_t cnt = 0; + if (sqlite3_step(stmt) == SQLITE_ROW) cnt = (uint32_t)sqlite3_column_int64(stmt, 0); sqlite3_finalize(stmt); + QByteArray out; out.append((const char*)&cnt, 4); return out; } -QList DbManager::getMessageSigs(const QString& channelId, - qint64 day, int hour) const { - QList sigs; - qint64 dayStart = day + hour * 3600; - qint64 dayEnd = dayStart + 3600; - sqlite3_stmt* stmt = prepareOrNull( - "SELECT signature FROM messages WHERE channel_id=?" - " AND timestamp/1000 >= ? AND timestamp/1000 < ? ORDER BY timestamp"); - if (!stmt) return sigs; - sqlite3_bind_text(stmt, 1, channelId.toUtf8().constData(), -1, SQLITE_STATIC); - sqlite3_bind_int64(stmt, 2, dayStart); - sqlite3_bind_int64(stmt, 3, dayEnd); - while (sqlite3_step(stmt) == SQLITE_ROW) - sigs.append(colBlob(stmt, 0)); - sqlite3_finalize(stmt); - return sigs; -} - -QList DbManager::getMessagesBySigs(const QList& sigs) const { - QList msgs; - if (sigs.isEmpty()) return msgs; - QString sql = "SELECT id, channel_id, author_node_id, content, content_type," - " timestamp, signature, is_outgoing, is_read FROM messages" - " WHERE signature IN ("; - for (int i = 0; i < sigs.size(); i++) - sql += (i > 0 ? ",?" : "?"); - sql += ") ORDER BY timestamp"; +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 ?") + .arg(msgTableName(chId)); sqlite3_stmt* stmt = prepareOrNull(sql.toUtf8().constData()); - if (!stmt) return msgs; - for (int i = 0; i < sigs.size(); i++) - sqlite3_bind_blob(stmt, i + 1, sigs[i].constData(), sigs[i].size(), SQLITE_STATIC); - while (sqlite3_step(stmt) == SQLITE_ROW) { - MessageRow m; - m.id = sqlite3_column_int64(stmt, 0); - m.channelId = colText(stmt, 1); - m.authorNodeId = static_cast(sqlite3_column_int64(stmt, 2)); - m.content = colText(stmt, 3); - m.contentType = colText(stmt, 4); - m.timestamp = sqlite3_column_int64(stmt, 5); - m.signature = colBlob(stmt, 6); - m.isOutgoing = sqlite3_column_int(stmt, 7) != 0; - m.isRead = sqlite3_column_int(stmt, 8) != 0; - msgs.append(m); - } + if (!stmt) return QByteArray(32, '\0'); + sqlite3_bind_int64(stmt, 1, pos); + QByteArray r; + if (sqlite3_step(stmt) == SQLITE_ROW) r = colBlob(stmt, 0); sqlite3_finalize(stmt); - return msgs; + return r.size() == 32 ? r : QByteArray(32, '\0'); } -QList DbManager::getOnlinePeers() const { - QList peers; - sqlite3_stmt* stmt = prepareOrNull( - "SELECT node_id FROM nodes WHERE is_online=1 AND node_id != ?"); - if (!stmt) return peers; - sqlite3_bind_int64(stmt, 1, static_cast(m_myNodeId)); - while (sqlite3_step(stmt) == SQLITE_ROW) - peers.append(static_cast(sqlite3_column_int64(stmt, 0))); +QByteArray DbManager::syncInsert(const QString& chId, const QByteArray& recData) { + /* recData: [timestamp:8][datahash:8][node_id:8][ct_len:1][ct:var][data_len:4][data:var][chain_hash:32] */ + if (recData.size() < 29) return QByteArray(1, (char)-1); + + const uint8_t* p = (const uint8_t*)recData.constData(); + int rem = recData.size(); + qint64 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) return QByteArray(1, (char)-1); + uint8_t ctLen = *p; p++; rem--; + if (rem < ctLen) return QByteArray(1, (char)-1); + QString ct = QString::fromUtf8((const char*)p, ctLen); p += ctLen; rem -= ctLen; + if (rem < 4) return QByteArray(1, (char)-1); + uint32_t dLen; memcpy(&dLen, p, 4); p += 4; rem -= 4; + if (rem < (int)dLen) return QByteArray(1, (char)-1); + QByteArray data((const char*)p, dLen); p += dLen; rem -= dLen; + if (rem < 32) return QByteArray(1, (char)-1); + QByteArray peerCH((const char*)p, 32); + + /* verify chain_hash */ + QByteArray prev = getLastChainHash(chId); + QByteArray expected = computeChainHash(prev, ts, dh); + if (expected != peerCH) { + GUI_ERROR("syncInsert: chain_hash mismatch for %s ts=%lld dh=%llx", chId.toUtf8().constData(), (long long)ts, (unsigned long long)dh); + return QByteArray(1, (char)-1); + } + + QString sql = QStringLiteral( + "INSERT OR IGNORE INTO \"%1\"" + " (node_id, content_type, data, timestamp, datahash, chain_hash, signature, is_outgoing, is_read)" + " VALUES(?,?,?,?,?,?,?,0,1)").arg(msgTableName(chId)); + sqlite3_stmt* stmt = prepareOrNull(sql.toUtf8().constData()); + if (!stmt) return QByteArray(1, (char)-1); + sqlite3_bind_int64(stmt, 1, (sqlite3_int64)nid); + sqlite3_bind_text(stmt, 2, ct.toUtf8().constData(), -1, SQLITE_STATIC); + sqlite3_bind_blob(stmt, 3, data.constData(), data.size(), SQLITE_STATIC); + sqlite3_bind_int64(stmt, 4, ts); + sqlite3_bind_int64(stmt, 5, (sqlite3_int64)dh); + sqlite3_bind_blob(stmt, 6, peerCH.constData(), 32, SQLITE_STATIC); + sqlite3_bind_blob(stmt, 7, kZero64, 64, SQLITE_STATIC); + int rc = sqlite3_step(stmt); sqlite3_finalize(stmt); - return peers; + if (rc == SQLITE_DONE) return QByteArray(1, '\0'); + if (rc == SQLITE_CONSTRAINT) return QByteArray(1, (char)1); /* duplicate */ + return QByteArray(1, (char)-1); } -bool DbManager::insertMessage(const QString& channelId, quint64 authorNodeId, - const QString& content, qint64 timestamp, int isOutgoing) { - static const unsigned char sig64[64] = {}; - sqlite3_stmt* stmt = prepareOrNull( - "INSERT INTO messages" - " (channel_id,author_node_id,content,timestamp,signature,is_outgoing,is_read)" - " VALUES(?,?,?,?,?,?,1)"); - if (!stmt) return false; - sqlite3_bind_text(stmt, 1, channelId.toUtf8().constData(), -1, SQLITE_STATIC); - sqlite3_bind_int64(stmt, 2, static_cast(authorNodeId)); - sqlite3_bind_text(stmt, 3, content.toUtf8().constData(), -1, SQLITE_STATIC); - sqlite3_bind_int64(stmt, 4, timestamp); - sqlite3_bind_blob(stmt, 5, sig64, 64, SQLITE_STATIC); - sqlite3_bind_int(stmt, 6, isOutgoing); - bool ok = (sqlite3_step(stmt) == SQLITE_DONE); - sqlite3_finalize(stmt); +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)); + sqlite3_stmt* stmt = prepareOrNull(sql.toUtf8().constData()); + if (!stmt) { QByteArray r(4, '\0'); return r; } + uint32_t id = m_nextCursorId++; + m_cursors[id] = stmt; + QByteArray r; r.append((const char*)&id, 4); return r; +} - /* update bucket_meta */ - qint64 day = timestamp / 86400000; /* ms per day */ - int hour = static_cast((timestamp % 86400000) / 3600000); - auto sigs = getMessageSigs(channelId, day, hour); - int cnt = sigs.size(); - QByteArray hash(32, '\0'); - if (cnt > 0) { - EVP_MD_CTX* md_ctx = EVP_MD_CTX_new(); - EVP_DigestInit_ex(md_ctx, EVP_sha256(), NULL); - for (auto& s : sigs) - EVP_DigestUpdate(md_ctx, s.constData(), s.size()); - unsigned int len = SHA256_DIGEST_LENGTH; - uint8_t buf[SHA256_DIGEST_LENGTH]; - EVP_DigestFinal_ex(md_ctx, buf, &len); - EVP_MD_CTX_free(md_ctx); - hash = QByteArray(reinterpret_cast(buf), static_cast(len)); - } - setBucketMeta(channelId, day, hour, cnt, hash); - return ok; +QByteArray DbManager::syncCursorNext(uint32_t cursorId) { + auto it = m_cursors.find(cursorId); + if (it == m_cursors.end()) return {}; + if (sqlite3_step(it.value()) != SQLITE_ROW) return {}; + + /* 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); + + QByteArray out; + out.append((const char*)&ts, 8); + out.append((const char*)&dh, 8); + out.append((const char*)&nid, 8); + QByteArray ctb = ct.toUtf8(); + out.append((char)(uint8_t)ctb.size()); + out.append(ctb); + uint32_t dLen = (uint32_t)data.size(); + out.append((const char*)&dLen, 4); + out.append(data); + return out; } -QByteArray DbManager::getLocalPubkey() const { - sqlite3_stmt* stmt = prepareOrNull( - "SELECT x25519_pubkey FROM local_identity WHERE id=1"); +void DbManager::syncCursorClose(uint32_t cursorId) { + auto it = m_cursors.find(cursorId); + if (it != m_cursors.end()) { sqlite3_finalize(it.value()); m_cursors.erase(it); } +} + +QByteArray DbManager::syncMarkSent(const QString& chId, quint64 ts, quint64 dh, quint64 myNodeId) { + QString sql = QStringLiteral( + "UPDATE \"%1\" SET sync_flags=sync_flags|1" + " WHERE timestamp=? AND datahash=? AND node_id=? AND is_outgoing=1") + .arg(msgTableName(chId)); + sqlite3_stmt* stmt = prepareOrNull(sql.toUtf8().constData()); if (!stmt) return {}; - QByteArray r; - if (sqlite3_step(stmt) == SQLITE_ROW) - r = colBlob(stmt, 0); - sqlite3_finalize(stmt); - return r; + sqlite3_bind_int64(stmt, 1, (sqlite3_int64)ts); + sqlite3_bind_int64(stmt, 2, (sqlite3_int64)dh); + sqlite3_bind_int64(stmt, 3, (sqlite3_int64)myNodeId); + sqlite3_step(stmt); sqlite3_finalize(stmt); + return {}; } -QList DbManager::getOnlinePeerAddresses(quint64 excludeNodeId) const { - QList list; - sqlite3_stmt* stmt = prepareOrNull( - "SELECT na.node_id, na.family, na.address, na.port" - " FROM node_addresses na JOIN nodes n ON na.node_id = n.node_id" - " WHERE n.is_online=1 AND na.is_nat=0 AND na.node_id != ?" - " ORDER BY RANDOM()"); - if (!stmt) return list; - sqlite3_bind_int64(stmt, 1, static_cast(excludeNodeId)); - while (sqlite3_step(stmt) == SQLITE_ROW) { - NodeAddr a; - a.nodeId = static_cast(sqlite3_column_int64(stmt, 0)); - a.family = sqlite3_column_int(stmt, 1); - a.address = colBlob(stmt, 2); - a.port = static_cast(sqlite3_column_int(stmt, 3)); - list.append(a); - } +QByteArray DbManager::syncTtlDelete(const QString& chId, quint64 myNodeId, quint64 cutoffUs) { + QString sql = QStringLiteral( + "DELETE FROM \"%1\" WHERE node_id=? AND is_outgoing=1 AND (sync_flags&1)=0 AND timestamp ids; + while (sqlite3_step(stmt) == SQLITE_ROW) + ids.append((uint64_t)sqlite3_column_int64(stmt, 0)); sqlite3_finalize(stmt); - return list; + QByteArray out; + uint16_t cnt = (uint16_t)ids.size(); + out.append((const char*)&cnt, 2); + for (auto id : ids) out.append((const char*)&id, 8); + return out; +} + +/* ── gui notifications ── */ + +void DbManager::onSyncMessageReceived(const QString& chId) { + Q_UNUSED(chId) + /* Обновление last_message в channels не нужно (убран). Вызывается emitter'ом из gui_bridge_impl */ } -/* ======================================================================== */ +void DbManager::setNodeOnline(quint64 nodeId, bool online) { + Q_UNUSED(nodeId); Q_UNUSED(online) + /* is_online убран из nodes. Обновление статуса через conn_mgr — позже */ +} -static constexpr unsigned char kDummySig[64] = { - 0x00,0x01,0x02,0x03,0x04,0x05,0x06,0x07, - 0x08,0x09,0x0A,0x0B,0x0C,0x0D,0x0E,0x0F, - 0x10,0x11,0x12,0x13,0x14,0x15,0x16,0x17, - 0x18,0x19,0x1A,0x1B,0x1C,0x1D,0x1E,0x1F, - 0x20,0x21,0x22,0x23,0x24,0x25,0x26,0x27, - 0x28,0x29,0x2A,0x2B,0x2C,0x2D,0x2E,0x2F, - 0x30,0x31,0x32,0x33,0x34,0x35,0x36,0x37, - 0x38,0x39,0x3A,0x3B,0x3C,0x3D,0x3E,0x3F -}; +/* ── seed stub data ── */ void DbManager::seedStubData() { GUI_INFO("Seeding stub data..."); sqlite3_exec(m_db, "BEGIN", nullptr, nullptr, nullptr); - /* my_node dummy pubkey */ unsigned char pub_dummy[32] = {}; sqlite3_stmt* ins = prepareOrNull( "INSERT OR IGNORE INTO local_identity(id,node_id,name,x25519_pubkey,ed25519_pubkey)" @@ -480,16 +764,15 @@ void DbManager::seedStubData() { sqlite3_bind_blob(ins, 3, pub_dummy, 32, SQLITE_STATIC); sqlite3_step(ins); sqlite3_finalize(ins); - /* nodes + accounts */ struct { quint64 id; const char* name; const char* letter; const char* color; } accs[] = { {1, "Alice", "A", "#4A90E2"}, {2, "Bob", "B", "#7ED321"}, {3, "Carol", "C", "#F5A623"}, {4, "Dave", "D", "#9013FE"}, {5, "Eve", "E", "#00B894"}, }; for (auto& a : accs) { - ins = prepareOrNull("INSERT OR IGNORE INTO nodes(node_id,name,x25519_pubkey,ed25519_pubkey)" - " VALUES(?,?,?,?)"); - sqlite3_bind_int64(ins, 1, static_cast(a.id)); + ins = prepareOrNull("INSERT OR IGNORE INTO nodes(node_id,name,x25519_pubkey,ed25519_pubkey,last_seen_at)" + " VALUES(?,?,?,?,unixepoch())"); + sqlite3_bind_int64(ins, 1, (sqlite3_int64)a.id); sqlite3_bind_text(ins, 2, a.name, -1, SQLITE_STATIC); sqlite3_bind_blob(ins, 3, pub_dummy, 32, SQLITE_STATIC); sqlite3_bind_blob(ins, 4, pub_dummy, 32, SQLITE_STATIC); @@ -497,149 +780,24 @@ void DbManager::seedStubData() { ins = prepareOrNull("INSERT OR IGNORE INTO accounts(node_id,display_name,avatar_color,avatar_letter)" " VALUES(?,?,?,?)"); - sqlite3_bind_int64(ins, 1, static_cast(a.id)); + sqlite3_bind_int64(ins, 1, (sqlite3_int64)a.id); sqlite3_bind_text(ins, 2, a.name, -1, SQLITE_STATIC); sqlite3_bind_text(ins, 3, a.color, -1, SQLITE_STATIC); sqlite3_bind_text(ins, 4, a.letter, -1, SQLITE_STATIC); sqlite3_step(ins); sqlite3_finalize(ins); } - /* set some nodes online for demo */ - sqlite3_exec(m_db, - "UPDATE nodes SET is_online=1 WHERE node_id IN (1,2,3)", nullptr, nullptr, nullptr); - - /* node addresses for online peers (stub, for invite link demo) */ - struct { quint64 id; int fam; unsigned char addr[16]; int port; int nat; } addrs[] = { - {1, 4, {192,168,1,10}, 8001, 0}, - {1, 6, {0xfd,0x12,0,0,0,0,0,0,0,0,0,0,0,0,0,1}, 8001, 0}, - {2, 4, {10,0,0,5}, 9001, 0}, - {2, 4, {172,16,0,100}, 9001, 0}, - {3, 4, {192,168,5,20}, 10001, 1}, /* Carol behind NAT */ - }; - for (auto& ad : addrs) { - ins = prepareOrNull("INSERT OR IGNORE INTO node_addresses(node_id,family,address,port,is_nat)" - " VALUES(?,?,?,?,?)"); - sqlite3_bind_int64(ins, 1, static_cast(ad.id)); - sqlite3_bind_int(ins, 2, ad.fam); - sqlite3_bind_blob(ins, 3, ad.addr, ad.fam == 4 ? 4 : 16, SQLITE_STATIC); - sqlite3_bind_int(ins, 4, ad.port); - sqlite3_bind_int(ins, 5, ad.nat); - sqlite3_step(ins); sqlite3_finalize(ins); - } - - /* channels (5) */ - struct { const char* cid; const char* name; } chs[] = { - {"ch:general", "#general"}, - {"ch:random", "#random"}, - {"ch:dev-team", "#dev-team"}, - {"ch:announce", "#announcements"}, - {"ch:offtopic", "#offtopic"}, - }; - for (auto& ch : chs) { - ins = prepareOrNull("INSERT OR IGNORE INTO channels(channel_id,name) VALUES(?,?)"); - sqlite3_bind_text(ins, 1, ch.cid, -1, SQLITE_STATIC); - sqlite3_bind_text(ins, 2, ch.name, -1, SQLITE_STATIC); - sqlite3_step(ins); sqlite3_finalize(ins); - - /* my node member */ - ins = prepareOrNull("INSERT OR IGNORE INTO channel_members(channel_id,node_id) VALUES(?,?)"); - sqlite3_bind_text(ins, 1, ch.cid, -1, SQLITE_STATIC); - sqlite3_bind_int64(ins, 2, static_cast(m_myNodeId)); - sqlite3_step(ins); sqlite3_finalize(ins); - - /* all 5 accounts as members (everyone in every channel for demo) */ - for (auto& a : accs) { - ins = prepareOrNull("INSERT OR IGNORE INTO channel_members(channel_id,node_id) VALUES(?,?)"); - sqlite3_bind_text(ins, 1, ch.cid, -1, SQLITE_STATIC); - sqlite3_bind_int64(ins, 2, static_cast(a.id)); - sqlite3_step(ins); sqlite3_finalize(ins); - } - } - - /* messages (8) */ - struct { const char* cid; quint64 author; qint64 ts; const char* text; quint8 o; } msgs[] = { - {"ch:general", 1, 1719200100000, "Hey team! How's the deployment going?", 0}, - {"ch:general", 2, 1719200220000, "All green on staging. Running the final integration tests now.", 0}, - {"ch:general", 3, 1719200520000, "I've pushed the ETCP-42 fix to staging. BGP handshake should be stable now. Please review the patch when you get a chance.", 0}, - {"ch:general", 4, 1719200700000, "Wait — the BGP session is still flapping after the handshake. I'm seeing the keepalive timer fire but no route updates come through. Logs attached.", 0}, - {"ch:general", 1, 1719200880000, "Reviewed the patch. LGTM. Let's merge and deploy to production tomorrow morning.", 0}, - {"ch:general", 2, 1719201000000, "Can you share the full logs for the BGP session? I'll dig into the keepalive timer logic.", 0}, - {"ch:general", 3, 1719201180000, "Sounds good. Let's coordinate in #dev-team. I'll prepare the release notes.", 0}, - {"ch:general", 1, 1719201300000, "Deploy scheduled for tomorrow 10:00 UTC. All features are frozen — no more merges to main.", 0}, - }; - for (auto& m : msgs) { - ins = prepareOrNull( - "INSERT OR IGNORE INTO messages" - " (channel_id,author_node_id,content,timestamp,signature,is_outgoing,is_read)" - " VALUES(?,?,?,?,?,?,1)"); - sqlite3_bind_text(ins, 1, m.cid, -1, SQLITE_STATIC); - sqlite3_bind_int64(ins, 2, static_cast(m.author)); - sqlite3_bind_text(ins, 3, m.text, -1, SQLITE_STATIC); - sqlite3_bind_int64(ins, 4, m.ts); - sqlite3_bind_blob(ins, 5, kDummySig, 64, SQLITE_STATIC); - sqlite3_bind_int(ins, 6, m.o); - sqlite3_step(ins); sqlite3_finalize(ins); - } + /* channels (1 демо-канал) */ + QByteArray dummyPub(32, '\0'), dummySig(64, '\0'), dummyCH(32, '\0'); + createChannel("ch:general", "#general", 0, dummyPub, QByteArray(), dummyPub, QByteArray()); - /* reactions */ - struct { qint64 msgRowId; quint64 node; const char* emoji; } reacts[] = { - {2, 1, "\xF0\x9F\x91\x8D"}, {2, 3, "\xF0\x9F\x91\x8D"}, - {2, 1, "\xF0\x9F\x91\x8B"}, {2, 3, "\xF0\x9F\x91\x8B"}, {2, 4, "\xF0\x9F\x91\x8B"}, - {3, 1, "\xF0\x9F\x91\x8D"}, {3, 2, "\xF0\x9F\x91\x8D"}, {3, 4, "\xF0\x9F\x91\x8D"}, - {3, 1, "\xE2\x9D\xA4"}, - {4, 2, "\xF0\x9F\x98\xAE"}, - {5, 1, "\xF0\x9F\x91\x8D"}, {5, 2, "\xF0\x9F\x91\x8D"}, {5, 3, "\xF0\x9F\x91\x8D"}, - {5, 4, "\xF0\x9F\x91\x8D"}, {5, 5, "\xF0\x9F\x91\x8D"}, - {5, 2, "\xF0\x9F\x8E\x89"}, {5, 5, "\xF0\x9F\x8E\x89"}, - {7, 1, "\xF0\x9F\xA4\x96"}, {7, 2, "\xF0\x9F\xA4\x96"}, - {8, 2, "\xF0\x9F\x9A\x80"}, {8, 3, "\xF0\x9F\x9A\x80"}, {8, 4, "\xF0\x9F\x9A\x80"}, - {8, 1, "\xE2\x9D\xA4"}, {8, 2, "\xE2\x9D\xA4"}, {8, 3, "\xE2\x9D\xA4"}, - {8, 4, "\xE2\x9D\xA4"}, {8, 5, "\xE2\x9D\xA4"}, {8, 1, "\xE2\x9D\xA4"}, {8, 1, "\xE2\x9D\xA4"}, - }; - std::string lastMsg; - qint64 lastTs = 0; - for (auto& r : reacts) { - ins = prepareOrNull( - "INSERT OR IGNORE INTO reactions(message_id,node_id,emoji,signature)" - " VALUES(?,?,?,?)"); - sqlite3_bind_int64(ins, 1, r.msgRowId); - sqlite3_bind_int64(ins, 2, static_cast(r.node)); - sqlite3_bind_text(ins, 3, r.emoji, -1, SQLITE_STATIC); - sqlite3_bind_blob(ins, 4, kDummySig, 64, SQLITE_STATIC); - sqlite3_step(ins); sqlite3_finalize(ins); + /* peers */ + for (auto& a : accs) { + addPeer("ch:general", a.id, dummySig, dummySig, a.name); } - lastMsg = "Deploy scheduled for tomorrow 10:00 UTC. All features are frozen — no more merges to main."; - lastTs = 1719201300000; - /* update channel last_msg */ - ins = prepareOrNull("UPDATE channels SET last_message=?, last_msg_at=? WHERE channel_id='ch:general'"); - sqlite3_bind_text(ins, 1, lastMsg.c_str(), -1, SQLITE_STATIC); - sqlite3_bind_int64(ins, 2, lastTs); - sqlite3_step(ins); sqlite3_finalize(ins); + /* messages — 0 для простоты, добавим позже в seed v2 */ sqlite3_exec(m_db, "COMMIT", nullptr, nullptr, nullptr); - GUI_INFO("Seed data inserted (5 accounts, 5 channels, 8 messages, %d reactions, %d addresses)", - static_cast(sizeof(reacts)/sizeof(reacts[0])), - static_cast(sizeof(addrs)/sizeof(addrs[0]))); -} - -QString DbManager::getUiState(const QString& key) const { - sqlite3_stmt* stmt = prepareOrNull("SELECT value FROM ui_state WHERE key=?"); - if (!stmt) return {}; - sqlite3_bind_text(stmt, 1, key.toUtf8().constData(), -1, SQLITE_STATIC); - QString val; - if (sqlite3_step(stmt) == SQLITE_ROW) - val = colText(stmt, 0); - sqlite3_finalize(stmt); - return val; -} - -void DbManager::setUiState(const QString& key, const QString& value) { - sqlite3_stmt* stmt = prepareOrNull( - "INSERT OR REPLACE INTO ui_state(key,value) VALUES(?,?)"); - if (!stmt) return; - sqlite3_bind_text(stmt, 1, key.toUtf8().constData(), -1, SQLITE_STATIC); - sqlite3_bind_text(stmt, 2, value.toUtf8().constData(), -1, SQLITE_STATIC); - sqlite3_step(stmt); - sqlite3_finalize(stmt); + GUI_INFO("Seed data inserted (5 accounts, 1 channel)"); } diff --git a/tools/chatgui/db/db_manager.h b/tools/chatgui/db/db_manager.h index 89bbec0f..c1b3a1cb 100644 --- a/tools/chatgui/db/db_manager.h +++ b/tools/chatgui/db/db_manager.h @@ -4,42 +4,56 @@ #include #include #include +#include #include #include struct ChannelRow { - QString channelId; - QString name; - bool isDm = false; - QString lastMessage; - qint64 lastMsgAt = 0; + QString channelId; + QString name; + quint64 ownerNodeId = 0; + bool isDm = false; + + QByteArray x25519_pubkey; + QByteArray x25519_privkey; + QByteArray ed25519_pubkey; + QByteArray ed25519_privkey; + QByteArray signature; // Ed25519 подпись всех полей записи privkey канала + + qint64 lastReadMsgId = 0; + qint64 lastPosMsgId = 0; + qint64 createdAt = 0; }; struct MessageRow { qint64 id = 0; - QString channelId; quint64 authorNodeId = 0; - QString content; QString contentType; + QByteArray data; qint64 timestamp = 0; + uint64_t datahash = 0; + QByteArray chainHash; QByteArray signature; bool isOutgoing = false; bool isRead = false; }; struct AccountRow { - quint64 nodeId = 0; - QString displayName; - QString avatarColor; - QString avatarLetter; - bool isContact = true; + quint64 nodeId = 0; + QString displayName; + QString avatarColor; + QString avatarLetter; + bool isContact = true; }; struct NodeAddr { quint64 nodeId = 0; - int family = 0; // 4=IPv4, 6=IPv6 + int family = 0; + int protocol = 0; // битмаска: 1=UDP, 2=TCP, 3=TCP+UDP QByteArray address; quint16 port = 0; + int rtt = 0; + bool isNat = false; }; class DbManager : public QObject { @@ -48,41 +62,90 @@ public: explicit DbManager(const QString& path, QObject* parent = nullptr); ~DbManager() override; - bool isOpen() const { return m_db != nullptr; } + bool isOpen() const { return m_db != nullptr; } quint64 myNodeId() const { return m_myNodeId; } - QList getChannels() const; - QList getMessages(const QString& channelId, int limit = 50) const; - QList getAccounts(bool contactsOnly = false) const; - AccountRow getAccount(quint64 nodeId) const; - QByteArray getNodeEdPub(quint64 nodeId) const; - - struct BucketMeta { qint64 day; int hour; int count; QByteArray hash; }; - QList getBucketMeta(const QString& channelId) const; - void setBucketMeta(const QString& channelId, qint64 day, int hour, - int count, const QByteArray& hash); - QList getMessageSigs(const QString& channelId, qint64 day, int hour) const; - QList getMessagesBySigs(const QList& sigs) const; - QList getOnlinePeers() const; - QByteArray getLocalPubkey() const; - QList getOnlinePeerAddresses(quint64 excludeNodeId) const; + /* ── Каналы ── */ + QList getChannels() const; + bool createChannel(const QString& chId, const QString& name, + quint64 ownerNodeId = 0, + const QByteArray& x25519pub = {}, const QByteArray& x25519priv = {}, + const QByteArray& ed25519pub = {}, const QByteArray& ed25519priv = {}); + bool deleteChannel(const QString& chId); - void seedStubData(); + /* ── Сообщения (UI) ── */ + QList getMessages(const QString& chId, int limit = 50) const; + bool insertMessage(const QString& chId, quint64 nodeId, + const QString& contentType, const QByteArray& data, + qint64 timestamp, bool isOutgoing, + uint64_t* outDatahash, QByteArray* outChainHash); + + /* ── Участники канала ── */ + bool addPeer(const QString& chId, quint64 nodeId, + const QByteArray& joinSig, const QByteArray& creatorSig, + const QString& comment = {}); + bool removePeer(const QString& chId, quint64 nodeId); - bool insertMessage(const QString& channelId, quint64 authorNodeId, - const QString& content, qint64 timestamp, int isOutgoing); + /* ── Узлы ── */ + bool upsertNode(quint64 nodeId, const QString& name, + const QByteArray& x25519pub, const QByteArray& ed25519pub); + bool addNodeAddress(quint64 nodeId, int family, int protocol, + const QByteArray& addr, quint16 port, bool isNat); + QByteArray loadNodeInfo(quint64 nodeId) const; + QList getOnlinePeers() const; + QByteArray getNodeEdPub(quint64 nodeId) const; + QList getOnlinePeerAddresses(quint64 excludeNodeId) const; + /* ── Аккаунты ── */ + QList getAccounts(bool contactsOnly = false) const; + AccountRow getAccount(quint64 nodeId) const; + + /* ── UI state ── */ QString getUiState(const QString& key) const; void setUiState(const QString& key, const QString& value); + /* ── Sync-операции (вызываются gui_bridge, в GUI-потоке) ── */ + QByteArray syncListChannels() const; + QByteArray syncCount(const QString& chId) const; + QByteArray syncChainHashAt(const QString& chId, uint32_t pos) const; + QByteArray syncInsert(const QString& chId, const QByteArray& recData); + QByteArray syncCursorOpen(const QString& chId); + QByteArray syncCursorNext(uint32_t cursorId); + void syncCursorClose(uint32_t cursorId); + QByteArray syncMarkSent(const QString& chId, quint64 ts, quint64 dh, quint64 myNodeId); + QByteArray syncTtlDelete(const QString& chId, quint64 myNodeId, quint64 cutoffUs); + QByteArray syncListPeers(const QString& chId) const; + + /* ── Уведомления из uasync ── */ + void onSyncMessageReceived(const QString& chId); + void setNodeOnline(quint64 nodeId, bool online); + void setUasyncInstance(void* ua) { m_ua = ua; } + void* uasyncInstance() const { return m_ua; } + + /* ── seed ── */ + void seedStubData(); + private: void initTables(); sqlite3_stmt* prepareOrNull(const char* sql) const; - static QString colText(sqlite3_stmt* stmt, int col); + static QString colText(sqlite3_stmt* stmt, int col); static QByteArray colBlob(sqlite3_stmt* stmt, int col); + QString msgTableName(const QString& chId) const; + QString peersTableName(const QString& chId) const; + static QString sanitizeTableName(const QString& chId); + static uint64_t computeDatahash(const QByteArray& data); + static QByteArray computeChainHash(const QByteArray& prevChain, + qint64 timestamp, uint64_t datahash); + QByteArray getLastChainHash(const QString& chId) const; + sqlite3* m_db = nullptr; quint64 m_myNodeId = 0; + void* m_ua = nullptr; + + /* Кэш открытых курсоров для sync-операций */ + mutable QMap m_cursors; + mutable uint32_t m_nextCursorId = 1; }; #endif diff --git a/tools/chatgui/doc/db_schema.md b/tools/chatgui/doc/db_schema.md index adfc8bac..5b195025 100644 --- a/tools/chatgui/doc/db_schema.md +++ b/tools/chatgui/doc/db_schema.md @@ -1,107 +1,244 @@ # База данных децентрализованного чата (chatgui) -Сторона сервиса (gui): -- +Чат децентрализованный — нет центрального сервера. Каждый узел хранит локальную SQLite БД +(`chats.db`, WAL mode). Синхронизация сообщений между узлами — через протокол `chat_sync` +поверх `etcp_router` (svc_id 0x30). Данные о других узлах поступают из NODEINFO +через conn_mgr. -## Архитектура хранения +Вся работа с БД — в GUI-потоке (класс `DbManager`). Протокол синхронизации работает +в uasync-потоке (`chat_sync.c`) и общается с БД асинхронно через `gui_bridge`. -Чат **децентрализованный** — нет центрального сервера. Каждый узел (нода uTun) хранит свою -**локальную SQLite БД**. В ней: +## Обзор таблиц -- **Информация о узлах (и о себе)** (как подключиться: адреса, порты, ключи — всё из NODEINFO). плюс nickname. - и список узлов = список участников. каждый узел имеет parent node (дерево). -- **История сообщений** — +| Таблица | Назначение | +|---------|-----------| +| `local_identity` | Наш собственный узел (id, ключи) | +| `nodes` | Все известные узлы (node_id, ключи, last_seen) | +| `node_addresses` | Адреса узлов (family, protocol, address, port, rtt, is_nat) | +| `accounts` | Профили узлов (display_name, avatar) | +| `channels` | Метаданные каналов/сетей (ключи, подпись, позиции скролла) | +| `msg_` | Per-channel таблица сообщений | +| `peers_` | Per-channel таблица участников (подписи входа, создателя, comment) | +| `ui_state` | Локальное состояние UI (key-value) | -Данные о других узлах поступают из NODEINFO gossip-протокола uTun и синхронизируются в фоне. -Сообщения приходят через `msg_transport` (TCP, локальный IPC). -аттачи - загружают и кешируют узлы в соответствии со своими настройками. +--- -## Структура таблиц SQLite +## Таблицы ### `local_identity` — наш собственный узел +```sql +CREATE TABLE IF NOT EXISTS local_identity ( + id INTEGER PRIMARY KEY CHECK (id = 1), + node_id INTEGER NOT NULL UNIQUE, + name TEXT NOT NULL, + x25519_pubkey BLOB NOT NULL, + x25519_privkey BLOB, + ed25519_pubkey BLOB, + created_at INTEGER DEFAULT (unixepoch()), + updated_at INTEGER DEFAULT (unixepoch()) +); +``` -### `my_nodes` — все мои устройства) -между моими устройствами синхронизируется контент - список чатов и сообщения. -... todo - +### `nodes` — все известные узлы ```sql CREATE TABLE IF NOT EXISTS nodes ( - id INTEGER PRIMARY KEY, - node_id INTEGER NOT NULL UNIQUE, // id из utun - name TEXT NOT NULL, // ник - pubkey (ed+x) - last_seen INTEGER DEFAULT 0, - online_rating INTEGER DEFAULT 0, // рейтинг доступности узла. чем больше число тем чаще узел онлайн - speed_rating INTEGER DEFAULT 0, // рейтинг скорости обмена с узлом - rtt_rating INTEGER DEFAULT 0, // rtt до узла (через промежуточные узлы если есть) - connectivity_json TEXT, // как связаться с узлом - адреса, (пока пусто) - parent_node_id INTEGER DEFAULT NULL, // каждая нода должна иметь родителя (кроме корневой) - main_node_id INTEGER DEFAULT NULL, // если несколько устройств у клиента, id следующей ноды (циклический список) - created_at TEXT DEFAULT (datetime('now')), + node_id INTEGER PRIMARY KEY, + name TEXT, + x25519_pubkey BLOB NOT NULL, + ed25519_pubkey BLOB, + last_seen_at INTEGER, + created_at INTEGER DEFAULT (unixepoch()) ); ``` +| Поле | Описание | +|------|---------| +| `node_id` | ID узла в сети uTun | +| `name` | Никнейм | +| `x25519_pubkey` | Публичный ключ для key exchange (32 байта) | +| `ed25519_pubkey` | Публичный ключ для подписей (32 байта), может быть NULL | +| `last_seen_at` | Unix-время последней активности | + +Статус online определяется через conn_mgr, а не полем в БД. -### `channels` — чаты (группы узлов) +### `node_addresses` — адреса узлов + +```sql +CREATE TABLE IF NOT EXISTS node_addresses ( + id INTEGER PRIMARY KEY AUTOINCREMENT, + node_id INTEGER NOT NULL REFERENCES nodes(node_id) ON DELETE CASCADE, + family INTEGER NOT NULL CHECK(family IN (4, 6)), + protocol INTEGER NOT NULL DEFAULT 1, + address BLOB NOT NULL, + port INTEGER NOT NULL CHECK(port > 0 AND port <= 65535), + rtt INTEGER, + is_nat INTEGER DEFAULT 0, + created_at INTEGER DEFAULT (unixepoch()) +); +CREATE INDEX IF NOT EXISTS idx_na_node ON node_addresses(node_id); +``` + +| Поле | Описание | +|------|---------| +| `family` | 4=IPv4, 6=IPv6 | +| `protocol` | Битовая маска: 1=UDP (NODE_PROTO_UDP), 2=TCP (NODE_PROTO_TCP), 3=UDP+TCP | +| `address` | BLOB: 4 байта для IPv4, 16 байт для IPv6 | +| `port` | Порт | +| `rtt` | RTT в миллисекундах, NULL если не измерен | +| `is_nat` | 1 = узел за NAT | + +### `accounts` — профили узлов + +```sql +CREATE TABLE IF NOT EXISTS accounts ( + node_id INTEGER PRIMARY KEY REFERENCES nodes(node_id), + display_name TEXT NOT NULL, + avatar_color TEXT DEFAULT '#4A90E2', + avatar_letter TEXT NOT NULL, + is_contact INTEGER DEFAULT 1, + created_at INTEGER DEFAULT (unixepoch()) +); +``` + +### `channels` — метаданные каналов ```sql CREATE TABLE IF NOT EXISTS channels ( - id INTEGER PRIMARY KEY AUTOINCREMENT, - name TEXT NOT NULL, - last_msg_at TEXT, + channel_id TEXT PRIMARY KEY, + name TEXT NOT NULL, + owner_node_id INTEGER, + is_dm INTEGER DEFAULT 0, + x25519_pubkey BLOB NOT NULL, + x25519_privkey BLOB, + ed25519_pubkey BLOB NOT NULL, + ed25519_privkey BLOB, + signature BLOB NOT NULL, + last_read_msg_id INTEGER, + last_pos_msg_id INTEGER, + created_at INTEGER DEFAULT (unixepoch()) +); +``` + +| Поле | Описание | +|------|---------| +| `channel_id` | Уникальный ID канала (напр. `"ch:general"`) | +| `name` | Отображаемое имя (напр. `"#general"`) | +| `owner_node_id` | ID создателя канала | +| `is_dm` | 1 = личная переписка (direct message) | +| `x25519_pubkey` | Публичный ключ канала (32 байта) | +| `x25519_privkey` | Приватный ключ канала, NULL если локальный узел не админ | +| `ed25519_pubkey` | Публичный ключ подписи канала (32 байта) | +| `ed25519_privkey` | Приватный ключ подписи, NULL если локальный узел не админ | +| `signature` | Ed25519 подпись всех полей записи приватным ключом канала | +| `last_read_msg_id` | ID последнего прочитанного сообщения (локальное) | +| `last_pos_msg_id` | ID сообщения на позиции скролла (для восстановления при открытии) | + +### `msg_` — per-channel таблица сообщений + +Для каждого канала создаётся отдельная таблица. Имя: `msg_` + sanitized `channel_id` +(не-буквоцифры заменяются на `_`). + +```sql +CREATE TABLE IF NOT EXISTS "msg_" ( + 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 NOT NULL, + chain_hash BLOB NOT NULL, + signature BLOB, + is_outgoing INTEGER DEFAULT 0, + is_read INTEGER DEFAULT 0, + sync_flags INTEGER DEFAULT 0, + UNIQUE(timestamp, datahash) ); +CREATE INDEX IF NOT EXISTS "idx_msg__ts" ON "msg_"(timestamp, datahash); ``` -### `messages` — список пользователей в чате +| Поле | Описание | +|------|---------| +| `node_id` | Автор сообщения | +| `content_type` | MIME-тип: `"text/plain"`, `"image/png"`, ... | +| `data` | Содержимое сообщения (BLOB) | +| `timestamp` | Время отправки в микросекундах (первая часть ключа sync) | +| `datahash` | Первые 8 байт SHA256(data), вторая часть ключа sync | +| `chain_hash` | SHA256 цепочки (32 байта) — как в blockchain, связывает записи по порядку | +| `signature` | Ed25519 подпись автора (64 байта), пока заглушка (64 нуля) | +| `is_outgoing` | 1 = отправлено нами | +| `is_read` | 1 = прочитано | +| `sync_flags` | Битовая маска: 0x01 = WAS_SENT (запись подтверждена пирами) | + +**Дедупликация:** уникальный индекс `(timestamp, datahash)` предотвращает повторную вставку +одного и того же сообщения при получении через разные пути. + +**Chain hash:** `chain_hash = SHA256(prev_chain_hash || timestamp || datahash)`. +Обеспечивает проверку целостности порядка сообщений при синхронизации. + +### `peers_` — per-channel таблица участников ```sql -CREATE TABLE IF NOT EXISTS channel_members ( - channel_id INTEGER, - node_id INTEGER, // если несколько устройств то ноды всех стройств +CREATE TABLE IF NOT EXISTS "peers_" ( + node_id INTEGER NOT NULL, + join_sig BLOB NOT NULL, + creator_sig BLOB, + comment TEXT, + joined_at INTEGER DEFAULT (unixepoch()), + PRIMARY KEY (node_id) ); ``` -### `messages` — история сообщений (для каждой группы создаем новую таблицу) +| Поле | Описание | +|------|---------| +| `join_sig` | Ed25519 подпись присоединяющегося: `sign(sk_user, channel_id || node_id)` | +| `creator_sig` | Ed25519 подпись создателя канала: `sign(sk_creator, channel_id || node_id)`. NULL = гость (ограниченные права) | +| `comment` | Комментарий/заметка об участнике | + +### `ui_state` — локальное состояние UI ```sql -CREATE TABLE IF NOT EXISTS messages_ch ( - id INTEGER PRIMARY KEY AUTOINCREMENT, - channel_id TEXT NOT NULL REFERENCES channels(channel_id) ON DELETE CASCADE, - author_node_id INTEGER NOT NULL, - content TEXT NOT NULL, - content_type TEXT DEFAULT 'text/plain', - timestamp INTEGER NOT NULL, - local_seq INTEGER DEFAULT 0, - signature BLOB, - is_outgoing INTEGER DEFAULT 0, - is_read INTEGER DEFAULT 0, - created_at TEXT DEFAULT (datetime('now')) +CREATE TABLE IF NOT EXISTS ui_state ( + key TEXT PRIMARY KEY, + value TEXT ); -CREATE INDEX IF NOT EXISTS idx_msg_channel_time ON messages(channel_id, timestamp); -CREATE INDEX IF NOT EXISTS idx_msg_author ON messages(author_node_id); -CREATE UNIQUE INDEX IF NOT EXISTS idx_msg_dedup ON messages(channel_id, author_node_id, timestamp); ``` -**Пояснение:** каждое сообщение привязано к каналу. `author_node_id` — кто автор -(из `nodes.node_id`). `timestamp` — unix время в миллисекундах *по часам отправителя* -(в децентрализованной системе нет глобальных часов, но для порядка в UI достаточно). -`local_seq` — монотонно возрастающий номер в пределах канала для детерминированной -сортировки при одинаковых timestamp. - -**Дедупликация:** сообщение идентифицируется тройкой `(channel_id, author_node_id, -timestamp)`. Уникальный индекс `idx_msg_dedup` предотвращает дубликаты при -повторной доставке через разные пути. - -| Поле | Тип | Описание | -|---|---|---| -| `channel_id` | TEXT | Ссылка на канал | -| `author_node_id` | INTEGER | Кто отправил (node_id) | -| `content` | TEXT | Текст сообщения | -| `content_type` | TEXT | MIME-тип: `text/plain`, в будущем `image/png` и т.д. | -| `timestamp` | INTEGER | Unix timestamp (мс) — когда отправлено | -| `local_seq` | INTEGER | Локальный порядковый номер в канале | -| `signature` | BLOB | Ed25519 подпись (64 байта), NULL если без подписи | -| `is_outgoing` | INTEGER | 1 = отправлено нами | -| `is_read` | INTEGER | 1 = прочитано (для галочек в UI) | +Key-value хранилище для произвольных настроек UI (позиции окон, выбранный канал, ...). + +--- + +## Формат binary-записей sync (gui_bridge / chat_sync) + +Записи передаются между потоками и по сети в компактном бинарном формате: + +### Insert record (для GUI_OP_INSERT_RECORD / PUSH wire) +``` +[timestamp:8 LE][datahash:8 LE][node_id:8 LE][ct_len:1][content_type:ct_len][data_len:4 LE][data:data_len][chain_hash:32] +``` + +### NodeInfo (GUI_OP_LOAD_NODEINFO) +``` +[x25519_pubkey:32][ed25519_pubkey:32][addr_count:1]([family:1][proto:1][addr_len:1][addr:4|16][port:2 LE][rtt:2 LE])... +``` + +### SEND_DATA cursor record (GUI_OP_CURSOR_NEXT) +``` +[ts:8 LE][dh:8 LE][node_id:8 LE][ct_len:1][ct:var][data_len:4 LE][data:var] +``` + +--- + +## Порядок синхронизации (chat_sync) + +1. **Peer online** → для каждого общего канала: `INIT_SYNC(my_count, last_chain_hash)` +2. **INIT_RESP** → сравнение chain_hash на позиции `min(counts)-1`. Если match → synced. + Иначе sparse hashes (-1, -2, -4, -8, ...) для поиска точки расхождения. +3. **SEND_DATA** → обмен недостающими записями пакетами до 32 штук +4. **SYNC_DONE** → финальная сверка count + last_chain_hash +5. **PUSH** → реал-тайм доставка новых сообщений (→ ACK_PUSH) +6. **TTL** → раз в час удаление своих неподтверждённых записей старше 24ч +Все операции с БД асинхронные: `chat_sync (uasync)` → `gui_bridge_call()` → +GUI-поток (SQL) → `uasync_post()` → callback в chat_sync. diff --git a/tools/chatgui/libutun/CMakeLists.txt b/tools/chatgui/libutun/CMakeLists.txt index 7ae486c5..f485daa5 100644 --- a/tools/chatgui/libutun/CMakeLists.txt +++ b/tools/chatgui/libutun/CMakeLists.txt @@ -2,6 +2,7 @@ cmake_minimum_required(VERSION 3.16) set(LIB_DIR ${CMAKE_CURRENT_SOURCE_DIR}/../../../lib) set(SRC_DIR ${CMAKE_CURRENT_SOURCE_DIR}/../../../src) +set(TRANSPORT_DIR ${CMAKE_CURRENT_SOURCE_DIR}/../transport) # ============================================================================ # libuasync — low-level runtime library @@ -61,7 +62,8 @@ set(UTUN_COMMON_SOURCES ${SRC_DIR}/route_node_lmdb.c ${SRC_DIR}/route_connectivity.c ${SRC_DIR}/conn_mgr.c - ${SRC_DIR}/db_sync.c + ${TRANSPORT_DIR}/db_sync_stub.c + ${TRANSPORT_DIR}/chat_sync.c ${SRC_DIR}/routing.c ${SRC_DIR}/tun_if.c ${SRC_DIR}/tun_route.c @@ -115,6 +117,7 @@ add_library(utun STATIC ${UTUN_COMMON_SOURCES} ${UTUN_TUN_SOURCES}) target_include_directories(utun PUBLIC ${SRC_DIR} ${LIB_DIR} + ${TRANSPORT_DIR} ${SRC_DIR}/uip ) diff --git a/tools/chatgui/src/channellist.cpp b/tools/chatgui/src/channellist.cpp index 9425c4d0..911a8400 100644 --- a/tools/chatgui/src/channellist.cpp +++ b/tools/chatgui/src/channellist.cpp @@ -77,7 +77,7 @@ void ChannelList::loadChannels() { auto *item = new QStandardItem(); item->setData(ch.name, ChannelNameRole); item->setData(ch.channelId, ChannelChannelIdRole); - item->setData(ch.lastMessage, ChannelLastMessageRole); + item->setData(QString(), ChannelLastMessageRole); item->setData(makeChannelIcon(colors[ci % 6], letter), ChannelIconRole); m_model->appendRow(item); ci++; diff --git a/tools/chatgui/src/mainwindow.cpp b/tools/chatgui/src/mainwindow.cpp index f8dc6483..9a941960 100644 --- a/tools/chatgui/src/mainwindow.cpp +++ b/tools/chatgui/src/mainwindow.cpp @@ -7,7 +7,7 @@ #include "../db/db_manager.h" #include "../transport/utun_node.h" #include "../transport/node_config.h" -#include "../transport/chat_propagator.h" +#include "../transport/gui_bridge.h" #include #include #include @@ -109,11 +109,11 @@ void MainWindow::setupNode() { void MainWindow::setupMessaging() { if (!m_node || !m_node->isRunning()) return; - m_propagator = new ChatPropagator(m_node, m_db, this); - m_messageList->setPropagator(m_propagator); - connect(m_propagator, &ChatPropagator::messageReceived, - m_messageList, &MessageList::addMessage); + /* Initialize gui_bridge */ + gui_bridge_init(NULL, NULL, this); + gui_bridge_set_db(m_db); + m_db->setUasyncInstance(NULL); /* будет обновлён при необходимости */ } void MainWindow::setupTray() { diff --git a/tools/chatgui/src/mainwindow.h b/tools/chatgui/src/mainwindow.h index b2b8e7d1..143253df 100644 --- a/tools/chatgui/src/mainwindow.h +++ b/tools/chatgui/src/mainwindow.h @@ -10,7 +10,6 @@ class AccountList; class DbManager; class UtunNode; class NodeConfig; -class ChatPropagator; class MainWindow : public QMainWindow { Q_OBJECT @@ -39,5 +38,4 @@ private: DbManager *m_db; UtunNode *m_node = nullptr; NodeConfig *m_config = nullptr; - ChatPropagator *m_propagator = nullptr; }; diff --git a/tools/chatgui/src/messagelist.cpp b/tools/chatgui/src/messagelist.cpp index f3c7e81d..ac3634ef 100644 --- a/tools/chatgui/src/messagelist.cpp +++ b/tools/chatgui/src/messagelist.cpp @@ -7,7 +7,6 @@ #include "animtimer.h" #include "lottieicon.h" #include "../db/db_manager.h" -#include "../transport/chat_propagator.h" #include #include #include @@ -121,13 +120,11 @@ MessageList::MessageList(DbManager* db, QWidget *parent) m_model->appendRow(item); m_view->scrollToBottom(); - if (m_db && !m_currentChannelId.isEmpty()) - m_db->insertMessage(m_currentChannelId, m_db->myNodeId(), text, ts, 1); - - if (m_propagator && !m_currentChannelId.isEmpty()) { - static const unsigned char zeroSig[64] = {}; - m_propagator->propagateNewMessage(m_currentChannelId, text.toUtf8(), ts, - QByteArray(reinterpret_cast(zeroSig), 64)); + if (m_db && !m_currentChannelId.isEmpty()) { + /* insert into DB — sync triggered internally via gui_bridge */ + QByteArray textData = text.toUtf8(); + m_db->insertMessage(m_currentChannelId, m_db->myNodeId(), + QStringLiteral("text/plain"), textData, ts, 1, nullptr, nullptr); } }); @@ -184,7 +181,7 @@ void MessageList::loadChannel(const QString& channelId) { QVariantList reacts; auto *item = new QStandardItem(); - setMsg(item, author, m.content, time, avatar, false, + setMsg(item, author, QString::fromUtf8(m.data), time, avatar, false, false, QString(), QString(), reacts); m_model->appendRow(item); } diff --git a/tools/chatgui/src/messagelist.h b/tools/chatgui/src/messagelist.h index 265f2de5..cdf40cf3 100644 --- a/tools/chatgui/src/messagelist.h +++ b/tools/chatgui/src/messagelist.h @@ -8,7 +8,6 @@ class ChatView; class InputBar; class EmojiPanel; class DbManager; -class ChatPropagator; class MessageList : public QWidget { Q_OBJECT @@ -18,7 +17,6 @@ public: void loadChannel(const QString& channelId); void addMessage(const QString& channelId, quint64 authorNodeId, const QByteArray& content, qint64 timestamp); - void setPropagator(ChatPropagator* p) { m_propagator = p; } private: ChatView *m_view; @@ -26,7 +24,6 @@ private: InputBar *m_inputBar; EmojiPanel *m_emojiPanel; DbManager *m_db; - ChatPropagator *m_propagator = nullptr; QString m_currentChannelId; bool m_animHoverActive = false; QHash m_avatarCache; diff --git a/tools/chatgui/transport/chat_sync.c b/tools/chatgui/transport/chat_sync.c new file mode 100644 index 00000000..9e6b2d23 --- /dev/null +++ b/tools/chatgui/transport/chat_sync.c @@ -0,0 +1,660 @@ +#include "chat_sync.h" +#include "gui_bridge.h" + +#include "../../../src/utun_instance.h" +#include "../../../src/etcp_router.h" +#include "../../../src/etcp_api.h" +#include "../../../src/etcp.h" +#include "../../../src/conn_mgr.h" +#include "../../../src/route_bgp.h" +#include "../../../src/route_node.h" +#include "../../../lib/u_async.h" +#include "../../../lib/ll_queue.h" +#include "../../../lib/debug_config.h" +#include "../../../lib/mem.h" +#include "../../../lib/sha256.h" +#include "../../../lib/platform_compat.h" + +#include + +static struct chat_sync* g_cs = NULL; + +/* ── Per-channel cache ── */ + +struct channel_cache { + char channel_id[64]; + uint64_t* peer_ids; + int peer_count; + uint32_t msg_count; + uint8_t last_chain_hash[32]; + uint8_t synced; +}; + +struct chat_sync { + struct UTUN_INSTANCE* inst; + struct channel_cache* channels; + int channel_count; + void* refresh_timer; + void* ttl_timer; + uint8_t initialized; +}; + +#define CS_ID "chat_sync" + +/* ── Send ── */ + +static int cs_send(struct chat_sync* cs, const char* ch_id, uint64_t dst, + const uint8_t* payload, size_t len) { + uint8_t ch_len = (uint8_t)strlen(ch_id); + if (ch_len > 63) return -1; + size_t total = 1 + 1 + ch_len + len; + uint8_t* buf = u_malloc(total); + if (!buf) return -1; + uint8_t* p = buf; + *p++ = ETCP_RT_ID_CHAT_SYNC; + *p++ = ch_len; + memcpy(p, ch_id, ch_len); p += ch_len; + memcpy(p, payload, len); + struct ll_entry* entry = queue_entry_new(0); + if (!entry) { u_free(buf); return -1; } + entry->dgram = buf; entry->len = total; + return etcp_route_send(cs->inst, dst, entry, 0); +} + +/* ── Channel cache ── */ + +static struct channel_cache* cs_find(struct chat_sync* cs, const char* ch_id) { + for (int i = 0; i < cs->channel_count; i++) + if (strcmp(cs->channels[i].channel_id, ch_id) == 0) return &cs->channels[i]; + return NULL; +} + +/* ── Protocol message context (async state machine) ── */ + +struct cs_ctx { + struct chat_sync* sync; + uint64_t peer; + char ch_id[64]; + union { + struct { + uint32_t peer_count; + uint8_t peer_last_ch[32]; + uint32_t my_count; + uint32_t test_pos; + } init_sync; + struct { + uint32_t peer_count; + uint32_t test_pos; + uint8_t peer_ch[37]; /* 1 byte for sparse_count + 36 for first sparse */ + int parse_off; /* offset in original payload for sparse data */ + size_t parse_len; + } init_resp; + struct { + uint32_t cursor_id; + uint32_t from_pos; + uint16_t sent_count; + } send_data; + struct { + uint64_t ts; + uint64_t dh; + } push; + }; +}; + +static struct cs_ctx* cs_ctx_new(struct chat_sync* cs, uint64_t peer, const char* ch_id) { + struct cs_ctx* ctx = u_calloc(1, sizeof(*ctx)); + if (!ctx) return NULL; + ctx->sync = cs; ctx->peer = peer; + size_t n = strlen(ch_id); if (n > 63) n = 63; + memcpy(ctx->ch_id, ch_id, n); + return ctx; +} + +/* ── gui_bridge async callbacks ── */ + +static void on_init_sync_count(void* ud, int op, const uint8_t* r, int rl); +static void on_init_sync_hash(void* ud, int op, const uint8_t* r, int rl); +static void on_init_sync_sparse(void* ud, int op, const uint8_t* r, int rl); +static void on_init_resp_hash(void* ud, int op, const uint8_t* r, int rl); +static void on_send_data_cursor_open(void* ud, int op, const uint8_t* r, int rl); +static void on_send_data_cursor_next(void* ud, int op, const uint8_t* r, int rl); +static void on_push_inserted(void* ud, int op, const uint8_t* r, int rl); +static void on_refresh_list(void* ud, int op, const uint8_t* r, int rl); +static void on_refresh_peers_and_count(void* ud, int op, const uint8_t* r, int rl); + +/* helper: build request [ch_id_len:1][ch_id:var] */ +static int req_ch(const char* ch_id, uint8_t* buf) { + uint8_t n = (uint8_t)strlen(ch_id); + buf[0] = n; memcpy(buf + 1, ch_id, n); + return 1 + n; +} + +/* ── 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) return; + struct cs_ctx* ctx = cs_ctx_new(cs, peer, ch_id); + if (!ctx) return; + memcpy(&ctx->init_sync.peer_count, pl, 4); + memcpy(ctx->init_sync.peer_last_ch, pl + 4, 32); + uint8_t req[65]; int rl = req_ch(ch_id, req); + gui_bridge_call(GUI_OP_COUNT, req, rl, ctx, on_init_sync_count); +} + +static void on_init_sync_count(void* ud, int op, const uint8_t* r, int rl) { + (void)op; + struct cs_ctx* ctx = (struct cs_ctx*)ud; + uint32_t my_count = (rl >= 4) ? *(const uint32_t*)r : 0; + ctx->init_sync.my_count = my_count; + + uint32_t tp = ctx->init_sync.peer_count < my_count ? ctx->init_sync.peer_count : my_count; + if (tp > 0) tp--; + ctx->init_sync.test_pos = tp; + + uint8_t req[69]; int nr = req_ch(ctx->ch_id, req); + memcpy(req + nr, &tp, 4); nr += 4; + gui_bridge_call(GUI_OP_CHAIN_HASH_AT, req, nr, ctx, on_init_sync_hash); +} + +static void on_init_sync_hash(void* ud, int op, const uint8_t* r, int rl) { + (void)op; + struct cs_ctx* ctx = (struct cs_ctx*)ud; + uint32_t tp = ctx->init_sync.test_pos; + uint8_t my_ch[32]; memset(my_ch, 0, 32); + if (rl >= 32) memcpy(my_ch, r, 32); + + if (memcmp(my_ch, ctx->init_sync.peer_last_ch, 32) == 0) { + /* synced */ + uint8_t resp[64]; resp[0] = CS_MSG_INIT_RESP; + memcpy(resp + 1, &ctx->init_sync.my_count, 4); resp[5] = 0; + cs_send(ctx->sync, ctx->ch_id, ctx->peer, resp, 6); + struct channel_cache* ch = cs_find(ctx->sync, ctx->ch_id); + if (ch) ch->synced = CS_SYNC_DONE; + u_free(ctx); return; + } + + /* Build sparse hashes: async chain */ + uint8_t resp[1024]; uint32_t off = 0; + resp[off++] = CS_MSG_INIT_RESP; + memcpy(resp + off, &ctx->init_sync.my_count, 4); off += 4; + memcpy(resp + off, &tp, 4); off += 4; + memcpy(resp + off, my_ch, 32); off += 32; + uint8_t sparse_pos = off; off++; /* placeholder */ + uint8_t sparse_count = 0; + + /* First sparse hash to start chain */ + uint32_t first_step = 1; + if (tp >= first_step) { + uint32_t pos = tp - first_step; + ctx->init_sync.test_pos = tp; /* сохраняем для продолжения */ + /* Запрашиваем первый sparse hash */ + uint8_t req2[69]; int nr2 = req_ch(ctx->ch_id, req2); + memcpy(req2 + nr2, &pos, 4); nr2 += 4; + gui_bridge_call(GUI_OP_CHAIN_HASH_AT, req2, nr2, ctx, on_init_sync_sparse); + return; + } + resp[sparse_pos] = 0; + cs_send(ctx->sync, ctx->ch_id, ctx->peer, resp, off); + u_free(ctx); +} + +static void on_init_sync_sparse(void* ud, int op, const uint8_t* r, int rl) { + (void)op; + struct cs_ctx* ctx = (struct cs_ctx*)ud; + /* For initial version: just send INIT_RESP with 0 sparse. + Peer will handle divergence. */ + uint32_t tp = ctx->init_sync.test_pos; + uint8_t resp[64]; resp[0] = CS_MSG_INIT_RESP; + memcpy(resp + 1, &ctx->init_sync.my_count, 4); + resp[5] = 0; /* sparse_count=0 */ + cs_send(ctx->sync, ctx->ch_id, ctx->peer, resp, 6); + u_free(ctx); +} + +/* ── 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 < 37) return; + struct cs_ctx* ctx = cs_ctx_new(cs, peer, ch_id); + if (!ctx) return; + ctx->init_resp.peer_count = *(const uint32_t*)pl; + ctx->init_resp.test_pos = *(const uint32_t*)(pl + 4); + memcpy(ctx->init_resp.peer_ch, pl + 8, 37); /* hash:32 + sparse_count:1 + maybe more */ + + uint8_t req[69]; int nr = req_ch(ch_id, req); + memcpy(req + nr, &ctx->init_resp.test_pos, 4); nr += 4; + gui_bridge_call(GUI_OP_CHAIN_HASH_AT, req, nr, ctx, on_init_resp_hash); +} + +static void on_init_resp_hash(void* ud, int op, const uint8_t* r, int rl) { + (void)op; + struct cs_ctx* ctx = (struct cs_ctx*)ud; + uint8_t my_ch[32]; memset(my_ch, 0, 32); + if (rl >= 32) memcpy(my_ch, r, 32); + + uint8_t sparse_count = ctx->init_resp.peer_ch[32]; + + /* Check match */ + if (memcmp(my_ch, ctx->init_resp.peer_ch, 32) == 0) { + /* Divergence is in our extra records (or peer's) */ + /* Request data from peer starting at test_pos+1 if peer has more */ + struct channel_cache* ch = cs_find(ctx->sync, ctx->ch_id); + if (ch && ctx->init_resp.peer_count > ctx->init_resp.test_pos + 1) { + /* We need peer's extra records */ + uint32_t from = ctx->init_resp.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(ctx->sync, ctx->ch_id, ctx->peer, snd, 7); + } else { + ch->synced = CS_SYNC_DONE; + } + u_free(ctx); return; + } + + /* Mismatch — request data from start */ + 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(ctx->sync, ctx->ch_id, ctx->peer, snd, 7); + u_free(ctx); +} + +/* ── 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) 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' */ + struct cs_ctx* ctx = cs_ctx_new(cs, peer, ch_id); + if (!ctx) return; + ctx->send_data.from_pos = from; + ctx->send_data.sent_count = 0; + + uint8_t req[65]; int nr = req_ch(ch_id, req); + gui_bridge_call(GUI_OP_CURSOR_OPEN, req, nr, ctx, on_send_data_cursor_open); + return; + } + + /* Peer sent US data — insert each record */ + const uint8_t* ptr = pl + 6; + size_t remain = len - 6; + + for (uint16_t i = 0; i < count && remain > 0; i++) { + /* record: [ts:8][dh:8][node_id:8][ct_len:1][ct:var][data_len:4][data:var] */ + /* For this, we call GUI_OP_INSERT_RECORD with the full record + chain_hash */ + /* But SEND_DATA doesn't include chain_hash — only PUSH does. + For SEND_DATA: skip chain_hash, just insert the content fields */ + if (remain < 29) 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) break; + ptr += dlen; remain -= dlen; + + size_t reclen = (size_t)(ptr - rec_start); + + /* Insert via GUI */ + uint8_t req2[2048]; int nr2 = req_ch(ch_id, req2); + if (nr2 + (int)reclen <= (int)sizeof(req2)) { + memcpy(req2 + nr2, rec_start, reclen); nr2 += (int)reclen; + /* Append zero chain_hash (not verified for SEND_DATA) */ + uint8_t zero_ch[32]; memset(zero_ch, 0, 32); + if (nr2 + 32 <= (int)sizeof(req2)) { + memcpy(req2 + nr2, zero_ch, 32); nr2 += 32; + } + gui_bridge_call(GUI_OP_INSERT_RECORD, req2, nr2, NULL, NULL); + } + } + + /* Request more if needed */ + 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; +} + +static void on_send_data_cursor_open(void* ud, int op, const uint8_t* r, int rl) { + (void)op; + struct cs_ctx* ctx = (struct cs_ctx*)ud; + if (rl < 4) { u_free(ctx); return; } + memcpy(&ctx->send_data.cursor_id, r, 4); + gui_bridge_call(GUI_OP_CURSOR_NEXT, (const uint8_t*)&ctx->send_data.cursor_id, 4, ctx, on_send_data_cursor_next); +} + +static void on_send_data_cursor_next(void* ud, int op, const uint8_t* r, int rl) { + (void)op; + struct cs_ctx* ctx = (struct cs_ctx*)ud; + if (rl == 0) { + gui_bridge_call(GUI_OP_CURSOR_CLOSE, (const uint8_t*)&ctx->send_data.cursor_id, 4, NULL, NULL); + u_free(ctx); return; + } + + /* Send record to peer */ + uint8_t snd[8192]; uint32_t off = 0; + snd[off++] = CS_MSG_SEND_DATA; + uint32_t pos = ctx->send_data.from_pos + ctx->send_data.sent_count; + memcpy(snd + off, &pos, 4); off += 4; + uint16_t c = 1; memcpy(snd + off, &c, 2); off += 2; + if (off + rl <= (uint32_t)sizeof(snd)) { memcpy(snd + off, r, rl); off += rl; } + cs_send(ctx->sync, ctx->ch_id, ctx->peer, snd, off); + ctx->send_data.sent_count++; + + gui_bridge_call(GUI_OP_CURSOR_NEXT, (const uint8_t*)&ctx->send_data.cursor_id, 4, ctx, on_send_data_cursor_next); +} + +/* ── 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) return; + uint64_t ts, dh, nid; + memcpy(&ts, pl, 8); memcpy(&dh, pl + 8, 8); memcpy(&nid, pl + 16, 8); + + struct cs_ctx* ctx = cs_ctx_new(cs, peer, ch_id); + if (!ctx) return; + ctx->push.ts = ts; ctx->push.dh = dh; + + /* Request: [ch_len][ch_id] + [rest of PUSH payload starting from ct_len] */ + /* PUSH payload: [ts:8][dh:8][node_id:8][ct_len:1][ct:var][data_len:4][data:var][chain_hash:32] */ + uint8_t req[2048]; int nr = req_ch(ch_id, req); + const uint8_t* rec = pl + 24; /* skip ts,dh,nid */ + int reclen = (int)(len - 24); + if (nr + reclen <= (int)sizeof(req)) { + memcpy(req + nr, rec, reclen); nr += reclen; + } + gui_bridge_call(GUI_OP_INSERT_RECORD, req, nr, ctx, on_push_inserted); +} + +static void on_push_inserted(void* ud, int op, const uint8_t* r, int rl) { + (void)op; + struct cs_ctx* ctx = (struct cs_ctx*)ud; + int rc = (rl >= 1) ? (int8_t)r[0] : -1; + if (rc == 0) { + uint8_t ack[17]; ack[0] = CS_MSG_ACK_PUSH; + memcpy(ack + 1, &ctx->push.dh, 8); + memcpy(ack + 9, &ctx->push.ts, 8); + cs_send(ctx->sync, ctx->ch_id, ctx->peer, ack, 17); + + uint8_t evt[65]; evt[0] = (uint8_t)strlen(ctx->ch_id); + memcpy(evt + 1, ctx->ch_id, evt[0]); + gui_bridge_post(GUI_EVT_MSG_RECEIVED, evt, 1 + evt[0]); + + struct channel_cache* ch = cs_find(ctx->sync, ctx->ch_id); + if (ch) ch->msg_count++; + } + u_free(ctx); +} + +/* ── 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) { + if (len < 16) return; + uint64_t ts, dh; memcpy(&dh, pl, 8); memcpy(&ts, pl + 8, 8); + uint8_t req[1 + 64 + 24]; int nr = req_ch(ch_id, req); + memcpy(req + nr, &ts, 8); nr += 8; + memcpy(req + nr, &dh, 8); nr += 8; + uint64_t myid = cs->inst->node_id; + memcpy(req + nr, &myid, 8); nr += 8; + gui_bridge_call(GUI_OP_MARK_SENT, req, nr, NULL, NULL); +} + +/* ── 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) { + if (len < 36) 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; } +} + +/* ── Recv dispatcher ── */ + +static void chat_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { + if (!entry || entry->len < 4) { if (entry) { if (entry->dgram) u_free(entry->dgram); queue_entry_free(entry); } return; } + if (!g_cs || !g_cs->initialized) { u_free(entry->dgram); queue_entry_free(entry); return; } + + uint64_t peer = conn ? conn->peer_node_id : 0; + const uint8_t* d = entry->dgram; + size_t dlen = entry->len; + + if (dlen < 4) { u_free(entry->dgram); queue_entry_free(entry); return; } + uint8_t ch_len = d[1]; + if (dlen < (size_t)(2 + ch_len + 1)) { u_free(entry->dgram); queue_entry_free(entry); return; } + + char ch_id[64]; memcpy(ch_id, d + 2, ch_len); ch_id[ch_len] = '\0'; + uint8_t type = d[2 + ch_len]; + const uint8_t* pl = d + 3 + ch_len; + size_t plen = dlen - 3 - ch_len; + + 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; + default: break; + } + u_free(entry->dgram); queue_entry_free(entry); +} + +/* ── Connection callbacks ── */ + +static void cs_on_conn_up(struct ETCP_CONN* conn, void* arg) { + (void)arg; + if (!conn || !g_cs) return; + uint64_t peer = conn->peer_node_id; + if (peer == 0 || peer == g_cs->inst->node_id) return; + 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; + uint8_t msg[37]; msg[0] = CS_MSG_INIT_SYNC; + 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); + } +} + +static void cs_on_conn_down(struct ETCP_CONN* conn, void* arg) { + (void)arg; + if (!conn || !g_cs) return; + uint64_t peer = conn->peer_node_id; + for (int i = 0; i < g_cs->channel_count; i++) { + int found = 0; + for (int j = 0; j < g_cs->channels[i].peer_count; j++) + if (g_cs->channels[i].peer_ids[j] == peer) { found = 1; break; } + if (found) g_cs->channels[i].synced = CS_SYNC_NONE; + } +} + +static void cs_on_new_conn(struct ETCP_CONN* conn, void* arg) { + (void)arg; + if (!conn) return; + etcp_conn_add_up_cbk(conn, cs_on_conn_up, NULL); + etcp_conn_add_down_cbk(conn, cs_on_conn_down, NULL); +} + +/* ── Periodic refresh from DB ── */ + +static void on_refresh_list(void* ud, int op, const uint8_t* r, int rl) { + (void)op; + if (!g_cs || !r || rl < 2) return; + uint16_t cnt; memcpy(&cnt, r, 2); + const uint8_t* p = r + 2; int rem = rl - 2; + + for (int i = 0; i < g_cs->channel_count; i++) + if (g_cs->channels[i].peer_ids) u_free(g_cs->channels[i].peer_ids); + if (g_cs->channels) u_free(g_cs->channels); + g_cs->channels = u_calloc(cnt, sizeof(struct channel_cache)); + g_cs->channel_count = cnt; + if (!g_cs->channels) { g_cs->channel_count = 0; return; } + + for (int i = 0; i < (int)cnt && rem >= 1; i++) { + uint8_t id_len = *p++; + rem--; + if (rem < id_len) break; + memcpy(g_cs->channels[i].channel_id, p, id_len); + g_cs->channels[i].channel_id[id_len] = '\0'; + p += id_len; rem -= id_len; + + uint8_t req[65]; req[0] = id_len; + memcpy(req + 1, g_cs->channels[i].channel_id, id_len); + gui_bridge_call(GUI_OP_LIST_PEERS, req, 1 + id_len, g_cs, on_refresh_peers_and_count); + } +} + +static void on_refresh_peers_and_count(void* ud, int op, const uint8_t* r, int rl) { + (void)op; (void)ud; + if (!g_cs || !r || rl < 2) return; + uint16_t cnt; memcpy(&cnt, r, 2); + /* Find first channel without peers (simplified: assumes sequential refresh) */ + for (int i = 0; i < g_cs->channel_count; i++) { + if (g_cs->channels[i].peer_ids) continue; + g_cs->channels[i].peer_ids = u_calloc(cnt, sizeof(uint64_t)); + g_cs->channels[i].peer_count = cnt; + for (uint16_t j = 0; j < cnt && 2 + j * 8 < (uint32_t)rl; j++) + memcpy(&g_cs->channels[i].peer_ids[j], r + 2 + j * 8, 8); + + /* Also get count */ + uint8_t req[65]; req[0] = (uint8_t)strlen(g_cs->channels[i].channel_id); + memcpy(req + 1, g_cs->channels[i].channel_id, req[0]); + gui_bridge_call(GUI_OP_COUNT, req, 1 + req[0], NULL, NULL); + break; + } +} + +static void refresh_timer_cb(void* arg) { + (void)arg; + if (!g_cs || !g_cs->initialized) return; + gui_bridge_call(GUI_OP_LIST_CHANNELS, NULL, 0, g_cs, on_refresh_list); + g_cs->refresh_timer = uasync_set_timeout(g_cs->inst->ua, 30u * 10000u, g_cs, refresh_timer_cb, "cs_refresh"); +} + +/* ── 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; + uint64_t myid = g_cs->inst->node_id; + for (int i = 0; i < g_cs->channel_count; i++) { + uint8_t req[1 + 64 + 16]; int nr = req_ch(g_cs->channels[i].channel_id, req); + memcpy(req + nr, &myid, 8); nr += 8; + memcpy(req + nr, &cutoff, 8); nr += 8; + gui_bridge_call(GUI_OP_TTL_DELETE, req, nr, NULL, NULL); + } + g_cs->ttl_timer = uasync_set_timeout(g_cs->inst->ua, 3600u * 10000u, g_cs, ttl_timer_cb, "cs_ttl"); +} + +/* ── Public API ── */ + +int chat_sync_init(struct UTUN_INSTANCE* inst, + void (*gui_cb)(void*, int, const uint8_t*, int)) { + (void)gui_cb; + if (!inst) return -1; + struct chat_sync* cs = u_calloc(1, sizeof(*cs)); + if (!cs) return -1; + cs->inst = inst; + cs->initialized = 1; + g_cs = cs; + + etcp_router_bind(inst, ETCP_RT_ID_CHAT_SYNC, chat_sync_recv_cb); + etcp_add_new_conn_cbk(inst, cs_on_new_conn, NULL); + + struct ETCP_CONN* c = inst->connections; + while (c) { + etcp_conn_add_up_cbk(c, cs_on_conn_up, NULL); + etcp_conn_add_down_cbk(c, cs_on_conn_down, NULL); + c = c->next; + } + + 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"); + + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: initialized", CS_ID); + return 0; +} + +void chat_sync_destroy(struct UTUN_INSTANCE* inst) { + struct chat_sync* cs = g_cs; + if (!cs || !inst) return; + cs->initialized = 0; g_cs = NULL; + + etcp_router_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; } + + struct ETCP_CONN* c = inst->connections; + while (c) { + etcp_conn_remove_up_cbk(c, cs_on_conn_up, NULL); + etcp_conn_remove_down_cbk(c, cs_on_conn_down, NULL); + c = c->next; + } + + for (int i = 0; i < cs->channel_count; i++) + if (cs->channels[i].peer_ids) u_free(cs->channels[i].peer_ids); + if (cs->channels) u_free(cs->channels); + 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) { + if (!g_cs || !g_cs->initialized) return -1; + struct channel_cache* ch = cs_find(g_cs, ch_id); + + uint8_t ct_len = content_type ? (uint8_t)strlen(content_type) : 0; + if (ct_len > 63) ct_len = 63; + size_t total = 1 + 8 + 8 + 8 + 1 + ct_len + 4 + data_len + 32; + 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; + 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; } + memset(p, 0, 32); /* chain_hash placeholder */ + + 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 + 32 - buf)); + sent++; + } + } + u_free(buf); + return sent; +} + +void chat_sync_connect_node(struct UTUN_INSTANCE* inst, uint64_t node_id) { + /* Placeholder — full impl requires parsing nodeinfo and calling conn_mgr_connect_node */ + uint8_t req[8]; memcpy(req, &node_id, 8); + gui_bridge_call(GUI_OP_LOAD_NODEINFO, req, 8, NULL, NULL); +} diff --git a/tools/chatgui/transport/chat_sync.h b/tools/chatgui/transport/chat_sync.h new file mode 100644 index 00000000..a20d9c8e --- /dev/null +++ b/tools/chatgui/transport/chat_sync.h @@ -0,0 +1,53 @@ +#ifndef CHAT_SYNC_H +#define CHAT_SYNC_H + +#ifdef __cplusplus +extern "C" { +#endif + +#include +#include + +struct UTUN_INSTANCE; +struct UASYNC; + +/* etcp_router service ID */ +#define ETCP_RT_ID_CHAT_SYNC 0x30 + +/* Message types */ +#define CS_MSG_INIT_SYNC 0x01 +#define CS_MSG_INIT_RESP 0x02 +#define CS_MSG_SEND_DATA 0x03 +#define CS_MSG_PUSH 0x04 +#define CS_MSG_ACK_PUSH 0x05 +#define CS_MSG_SYNC_DONE 0x06 + +/* Protocol constants */ +#define CS_SEND_DATA_MAX 32 +#define CS_MAX_SPARSE 16 +#define CS_PEER_TIMEOUT_MS 5000 + +/* Sync states per channel */ +#define CS_SYNC_NONE 0 +#define CS_SYNC_IN_PROGRESS 1 +#define CS_SYNC_DONE 2 + +/* ── Public API ── */ + +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); + +/* Прослойка conn_mgr: загрузить nodeinfo из БД и подключиться */ +void chat_sync_connect_node(struct UTUN_INSTANCE* inst, uint64_t node_id); + +#ifdef __cplusplus +} +#endif +#endif /* CHAT_SYNC_H */ diff --git a/tools/chatgui/transport/db_sync_stub.c b/tools/chatgui/transport/db_sync_stub.c new file mode 100644 index 00000000..06ad165c --- /dev/null +++ b/tools/chatgui/transport/db_sync_stub.c @@ -0,0 +1,31 @@ +#include "db_sync.h" +#include "utun_instance.h" +#include "../lib/u_async.h" +#include "../lib/debug_config.h" +#include "../lib/mem.h" + +int db_sync_init(struct UTUN_INSTANCE* inst) { + (void)inst; + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "db_sync: stub (disabled for chatgui)"); + return 0; +} + +void db_sync_destroy(struct UTUN_INSTANCE* inst) { + (void)inst; + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "db_sync: stub destroyed"); +} + +int db_sync_insert_len(struct UTUN_INSTANCE* inst, const char* json_data, size_t len) { + (void)inst; (void)json_data; (void)len; + return 0; +} + +int db_sync_insert(struct UTUN_INSTANCE* inst, const char* json_data) { + (void)inst; (void)json_data; + return 0; +} + +uint32_t db_sync_count(struct UTUN_INSTANCE* inst) { + (void)inst; + return 0; +} diff --git a/tools/chatgui/transport/gui_bridge.h b/tools/chatgui/transport/gui_bridge.h new file mode 100644 index 00000000..03e8b825 --- /dev/null +++ b/tools/chatgui/transport/gui_bridge.h @@ -0,0 +1,60 @@ +#ifndef GUI_BRIDGE_H +#define GUI_BRIDGE_H + +#ifdef __cplusplus +extern "C" { +#endif + +#include +#include + +struct UASYNC; + +/* ── Операции uasync→GUI (chat_sync запрашивает работу с БД) ── */ + +#define GUI_OP_LIST_CHANNELS 0 /* → массив: [count:2][ch_id_len:1][ch_id:var][name_len:2][name:var]... */ +#define GUI_OP_COUNT 1 /* req: [ch_id_len:1][ch_id:var] → resp: uint32 count (4 байта) */ +#define GUI_OP_CHAIN_HASH_AT 2 /* req: [ch_id_len:1][ch_id:var][pos:4] → resp: 32 байта chain_hash */ +#define GUI_OP_INSERT_RECORD 3 /* req: [ch_id_len:1][ch_id:var][record:var] → resp: 0=ok, 1=dup, -1=err (1 байт) */ +#define GUI_OP_CURSOR_OPEN 4 /* req: [ch_id_len:1][ch_id:var] → resp: uint32 cursor_id (4 байта) */ +#define GUI_OP_CURSOR_NEXT 5 /* req: cursor_id (4 байта) → resp: record или 0 байт = EOF */ +#define GUI_OP_CURSOR_CLOSE 6 /* req: cursor_id (4 байта) → resp: 0 байт */ +#define GUI_OP_MARK_SENT 7 /* req: [ch_id_len:1][ch_id:var][ts:8][dh:8][my_node_id:8] → resp: 0 байт */ +#define GUI_OP_TTL_DELETE 8 /* req: [ch_id_len:1][ch_id:var][my_node_id:8][cutoff:8] → resp: 0 байт */ +#define GUI_OP_LIST_PEERS 9 /* req: [ch_id_len:1][ch_id:var] → resp: [count:2][node_id:8]... */ +#define GUI_OP_LOAD_NODEINFO 10 /* req: node_id (8 байт) → resp: nodeinfo или 0 байт если нет */ + +/* ── Типы уведомлений uasync→GUI (fire-and-forget) ── */ + +#define GUI_EVT_MSG_RECEIVED 1 /* data: [ch_id_len:1][ch_id:var][record:var] */ +#define GUI_EVT_CONNECT_RESULT 2 /* data: [node_id:8][result:4] */ +#define GUI_EVT_NEW_PEER 3 /* data: [node_id:8] */ +#define GUI_EVT_CHANNEL_UPDATED 4 /* data: [ch_id_len:1][ch_id:var] */ + +/* ── Callback-типы ── */ + +typedef void (*gui_result_fn)(void* userdata, int op, + const uint8_t* result, int result_len); + +/* ── API ── */ + +/* Регистрация (вызывается из GUI при старте) */ +void gui_bridge_init(struct UASYNC* ua, gui_result_fn callback, void* gui_target); + +/* Установить указатель на DbManager (до первого вызова gui_bridge_call) */ +void gui_bridge_set_db(void* db_ptr); + +/* uasync → GUI: запросить SQL-операцию */ +void gui_bridge_call(int op, const uint8_t* request, int req_len, + void* userdata, gui_result_fn callback); + +/* uasync → GUI: уведомление (fire-and-forget, callback не вызывается) */ +void gui_bridge_post(int event_type, const uint8_t* data, int data_len); + +/* GUI → uasync: выполнить функцию в uasync-потоке */ +void gui_bridge_post_uasync(struct UASYNC* ua, void (*fn)(void*), void* arg); + +#ifdef __cplusplus +} +#endif +#endif /* GUI_BRIDGE_H */ diff --git a/tools/chatgui/transport/gui_bridge_impl.cpp b/tools/chatgui/transport/gui_bridge_impl.cpp new file mode 100644 index 00000000..719a9922 --- /dev/null +++ b/tools/chatgui/transport/gui_bridge_impl.cpp @@ -0,0 +1,252 @@ +#include "gui_bridge.h" + +#include +#include +#include +#include +#include + +#include "../db/db_manager.h" + +extern "C" { +#include "../../../lib/u_async.h" +#include "../../../lib/mem.h" +} + +/* ── Внутренний объект-приёмник в GUI-потоке ── */ + +class GuiBridgeReceiver : public QObject { + Q_OBJECT +public: + void* ua = nullptr; + gui_result_fn callback = nullptr; + DbManager* db = nullptr; + + GuiBridgeReceiver(QObject* parent = nullptr) : QObject(parent) {} + + Q_INVOKABLE void processCall(int op, QByteArray packedReq); + Q_INVOKABLE void processPost(int eventType, QByteArray data); +}; + +static GuiBridgeReceiver* g_receiver = nullptr; + +/* ── Трамплин для uasync_post: распаковывает результат и вызывает оригинальный callback ── */ + +struct TrampolineCtx { + gui_result_fn cb; + void* userdata; + int op; + int result_len; + /* result data follows after this struct in memory */ +}; + +static void bridge_trampoline(void* arg) { + uint8_t* packed = (uint8_t*)arg; + TrampolineCtx ctx; + memcpy(&ctx, packed, sizeof(ctx)); + const uint8_t* result = packed + sizeof(ctx); + if (ctx.cb) ctx.cb(ctx.userdata, ctx.op, result, ctx.result_len); + u_free(packed); +} + +static void post_result_to_uasync(int op, const QByteArray& result, + void* userdata, gui_result_fn cb) { + int res_len = result.size(); + int total = (int)sizeof(TrampolineCtx) + res_len; + uint8_t* buf = (uint8_t*)u_malloc(total); + if (!buf) return; + + TrampolineCtx ctx; + ctx.cb = cb; + ctx.userdata = userdata; + ctx.op = op; + ctx.result_len = res_len; + memcpy(buf, &ctx, sizeof(ctx)); + if (res_len > 0) memcpy(buf + sizeof(ctx), result.constData(), res_len); + + uasync_post((struct UASYNC*)g_receiver->ua, bridge_trampoline, buf); +} + +/* ── GuiBridgeReceiver implementation ── */ + +static QString readChId(const uint8_t*& r, int& rem) { + if (rem < 1) return {}; + uint8_t len = r[0]; r++; rem--; + if (rem < len) return {}; + QString s = QString::fromUtf8((const char*)r, len); + r += len; rem -= len; + return s; +} + +void GuiBridgeReceiver::processCall(int op, QByteArray packedReq) { + if (!db || packedReq.size() < (int)(sizeof(gui_result_fn) + sizeof(void*))) return; + + gui_result_fn cb; + void* userdata; + memcpy(&cb, packedReq.constData(), sizeof(cb)); + memcpy(&userdata, packedReq.constData() + sizeof(cb), sizeof(userdata)); + + const uint8_t* r = (const uint8_t*)packedReq.constData() + sizeof(cb) + sizeof(userdata); + int rem = packedReq.size() - (int)(sizeof(cb) + sizeof(userdata)); + + QByteArray result; + + switch (op) { + case GUI_OP_LIST_CHANNELS: + result = db->syncListChannels(); + break; + + case GUI_OP_COUNT: { + QString chId = readChId(r, rem); + if (!chId.isEmpty()) result = db->syncCount(chId); + break; + } + case GUI_OP_CHAIN_HASH_AT: { + QString chId = readChId(r, rem); + if (!chId.isEmpty() && rem >= 4) { + uint32_t pos; memcpy(&pos, r, 4); + result = db->syncChainHashAt(chId, pos); + } + break; + } + case GUI_OP_INSERT_RECORD: { + QString chId = readChId(r, rem); + if (!chId.isEmpty() && rem > 0) { + QByteArray rec((const char*)r, rem); + result = db->syncInsert(chId, rec); + } + break; + } + case GUI_OP_CURSOR_OPEN: { + QString chId = readChId(r, rem); + if (!chId.isEmpty()) result = db->syncCursorOpen(chId); + break; + } + case GUI_OP_CURSOR_NEXT: + if (rem >= 4) { + uint32_t curId; memcpy(&curId, r, 4); + result = db->syncCursorNext(curId); + } + break; + case GUI_OP_CURSOR_CLOSE: + if (rem >= 4) { + uint32_t curId; memcpy(&curId, r, 4); + db->syncCursorClose(curId); + } + break; + case GUI_OP_MARK_SENT: { + QString chId = readChId(r, rem); + if (!chId.isEmpty() && rem >= 24) { + uint64_t ts, dh, myId; + memcpy(&ts, r, 8); memcpy(&dh, r + 8, 8); memcpy(&myId, r + 16, 8); + result = db->syncMarkSent(chId, ts, dh, myId); + } + break; + } + case GUI_OP_TTL_DELETE: { + QString chId = readChId(r, rem); + if (!chId.isEmpty() && rem >= 16) { + uint64_t myId, cutoff; + memcpy(&myId, r, 8); memcpy(&cutoff, r + 8, 8); + result = db->syncTtlDelete(chId, myId, cutoff); + } + break; + } + case GUI_OP_LIST_PEERS: { + QString chId = readChId(r, rem); + if (!chId.isEmpty()) result = db->syncListPeers(chId); + break; + } + case GUI_OP_LOAD_NODEINFO: + if (rem >= 8) { + uint64_t nodeId; memcpy(&nodeId, r, 8); + result = db->loadNodeInfo(nodeId); + } + break; + default: + break; + } + + post_result_to_uasync(op, result, userdata, cb); +} + +void GuiBridgeReceiver::processPost(int eventType, QByteArray data) { + if (!db) return; + const uint8_t* d = (const uint8_t*)data.constData(); + int dlen = data.size(); + + switch (eventType) { + case GUI_EVT_MSG_RECEIVED: + if (dlen >= 2) { + uint8_t chLen = d[0]; + if (dlen >= 1 + chLen) { + QString chId = QString::fromUtf8((const char*)d + 1, chLen); + db->onSyncMessageReceived(chId); + } + } + break; + case GUI_EVT_CONNECT_RESULT: + if (dlen >= 12) { + uint64_t nodeId; int result; + memcpy(&nodeId, d, 8); memcpy(&result, d + 8, 4); + db->setNodeOnline(nodeId, result == 0); + } + break; + case GUI_EVT_NEW_PEER: + case GUI_EVT_CHANNEL_UPDATED: + default: + break; + } +} + +void* gui_bridge_get_db() { return g_receiver ? (void*)g_receiver->db : nullptr; } + +/* ── C API ── */ + +extern "C" { + +void gui_bridge_init(struct UASYNC* ua, gui_result_fn callback, void* gui_target) { + g_receiver = new GuiBridgeReceiver((QObject*)gui_target); + g_receiver->ua = ua; + g_receiver->callback = callback; +} + +void gui_bridge_set_db(void* db_ptr) { + if (g_receiver) g_receiver->db = (DbManager*)db_ptr; +} + +void gui_bridge_call(int op, const uint8_t* request, int req_len, + void* userdata, gui_result_fn callback) { + if (!g_receiver) return; + + /* Пакуем cb и userdata в начало, потом request */ + QByteArray packed; + packed.append((const char*)&callback, sizeof(callback)); + packed.append((const char*)&userdata, sizeof(userdata)); + if (req_len > 0) packed.append((const char*)request, req_len); + + QMetaObject::invokeMethod(g_receiver, "processCall", + Qt::QueuedConnection, + Q_ARG(int, op), + Q_ARG(QByteArray, packed)); +} + +void gui_bridge_post(int event_type, const uint8_t* data, int data_len) { + if (!g_receiver) return; + + QByteArray d; + if (data_len > 0) d = QByteArray((const char*)data, data_len); + + QMetaObject::invokeMethod(g_receiver, "processPost", + Qt::QueuedConnection, + Q_ARG(int, event_type), + Q_ARG(QByteArray, d)); +} + +void gui_bridge_post_uasync(struct UASYNC* ua, void (*fn)(void*), void* arg) { + if (ua && fn) uasync_post(ua, fn, arg); +} + +} /* extern "C" */ + +#include "gui_bridge_impl.moc" diff --git a/tools/chatgui/transport/utun_node.cpp b/tools/chatgui/transport/utun_node.cpp index f7043a15..345499f6 100644 --- a/tools/chatgui/transport/utun_node.cpp +++ b/tools/chatgui/transport/utun_node.cpp @@ -17,6 +17,8 @@ extern "C" { #include "../lib/memory_pool.h" #include "../lib/debug_config.h" #include "../lib/mem.h" +#include "chat_sync.h" +#include "gui_bridge.h" } #define ETCP_RT_ID_CHAT 0x11 @@ -122,11 +124,16 @@ void UtunNode::runLoop() { etcp_router_bind(m_instance, ETCP_RT_ID_CHAT, recvCallback); + /* Initialize chat_sync */ + chat_sync_init(m_instance, NULL); + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "chat_sync initialized"); + while (!m_stop) { uasync_poll(ua, 100); } etcp_router_bind(m_instance, ETCP_RT_ID_CHAT, nullptr); + chat_sync_destroy(m_instance); utun_instance_destroy(m_instance); m_instance = nullptr; uasync_destroy(ua, 0);