diff --git a/AGENTS.md b/AGENTS.md index 7ac3a0c8..6f4f1a70 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -428,6 +428,8 @@ Chatgui — десктопный GUI-чат на Qt 6 (Qt 5 fallback), отде В chat gui интегрированы библиотеки utun. Чат и библиотеки работают в разных потоках. Поэтому нужно использовать семафоры, сокеты или другие механизмы синхронизации (uasync_post, uasync_memsync, uasync_get_wakeup_fd) +База данных sqlite в chatgui: можно писать (update) из потока utun / instance. можно только читать из gui потока. + ### Технологии - **Язык:** C++20 - **Фреймворк:** Qt 6 (предпочитаемый) или Qt 5.15+ (Widgets + Network) diff --git a/src/topo_group.c b/src/topo_group.c index bc2f5c36..c5ae7d13 100644 --- a/src/topo_group.c +++ b/src/topo_group.c @@ -347,6 +347,10 @@ void topo_groups_set_sqlite_db(struct TOPO_GROUPS* g, sqlite3* db) { void topo_group_set_nat_check_local(struct TOPO_GROUP* group, int allow) { if (!group) return; group->allow_nat_check_local = allow ? 1 : 0; } +void topo_groups_set_node_updated_cb(struct TOPO_GROUPS* groups, topo_node_updated_fn fn) { + if (groups) groups->node_updated_cb = fn; +} + void topo_group_new_conn(struct TOPO_GROUP* group, struct ETCP_CONN* conn) { if (!group || !conn) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "invalid args"); return; } if (!conn->instance) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "conn->instance is NULL"); return; } @@ -636,6 +640,9 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from } if (group->instance->control_srv) control_server_notify_node_change(group->instance->control_srv, nodeinfo1); + if (group->instance->topo_groups->node_updated_cb) + group->instance->topo_groups->node_updated_cb(group->instance, node_id, ni->public_key, ni->ed25519_public_key); + int hop_count = nodeinfo1->hop_count; struct ll_entry* se = group->senders_list ? group->senders_list->head : NULL; while (se) { diff --git a/src/topo_group.h b/src/topo_group.h index 53178179..d6a66de4 100644 --- a/src/topo_group.h +++ b/src/topo_group.h @@ -18,6 +18,13 @@ extern "C" { #include #endif +struct UTUN_INSTANCE; + +/** Callback when a node's info (pubkeys, addresses) is persisted in the DB */ +typedef void (*topo_node_updated_fn)(struct UTUN_INSTANCE* inst, uint64_t node_id, + const uint8_t* x25519_pubkey, + const uint8_t* ed25519_pubkey); + // ETCP ID для пакетов топологии #define ETCP_ID_TOPO_ENTRY 0x01 @@ -137,6 +144,7 @@ struct TOPO_GROUPS { #ifdef USE_SQLITE sqlite3* topo_sqlite_db; #endif + topo_node_updated_fn node_updated_cb; /* optional — set by chatgui's member_sync */ }; /** @@ -267,6 +275,8 @@ void topo_group_send_nat_check_req(struct ETCP_CONN* conn, uint8_t socket_id); */ void topo_group_set_nat_check_local(struct TOPO_GROUP* group, int allow); +void topo_groups_set_node_updated_cb(struct TOPO_GROUPS* groups, topo_node_updated_fn fn); + /** * @brief Запускает NAT check для всех линков всех соединений. */ diff --git a/tests/test_member_sync.c b/tests/test_member_sync.c new file mode 100644 index 00000000..f7280a89 --- /dev/null +++ b/tests/test_member_sync.c @@ -0,0 +1,169 @@ +/** test_member_sync.c — standalone hash tree protocol test with 1000 members */ +#include +#include +#include +#include +#include +#include + +#define HASH_SZ 32 +#define BUCKETS 32 +#define MAX_LVL 5 + +static void _sanitize(const char* ch_id, char* out, size_t sz) { + size_t i = 0; + while (*ch_id && i < sz - 1) { + char c = *ch_id++; + if ((c>='a'&&c<='z')||(c>='A'&&c<='Z')||(c>='0'&&c<='9')||c=='_') out[i++]=c; else out[i++]='_'; + } + out[i]=0; +} + +static void _init_db(sqlite3* db) { + sqlite3_exec(db, + "CREATE TABLE IF NOT EXISTS nodes(node_id INTEGER PRIMARY KEY," + " x25519_pubkey BLOB NOT NULL, ed25519_pubkey BLOB, online INTEGER DEFAULT 0);" + "CREATE TABLE IF NOT EXISTS member_tree_hash(channel_id TEXT NOT NULL," + " level INTEGER NOT NULL, prefix64 INTEGER NOT NULL, hash BLOB NOT NULL," + " member_count INTEGER NOT NULL, PRIMARY KEY(channel_id, level, prefix64));", + NULL, NULL, NULL); +} + +static void _create_peers(sqlite3* db, const char* ch_id) { + char san[64]; _sanitize(ch_id, san, sizeof(san)); + char sql[256]; snprintf(sql, sizeof(sql), + "CREATE TABLE IF NOT EXISTS peers_%s(node_id INTEGER NOT NULL," + " join_sig BLOB NOT NULL, PRIMARY KEY(node_id))", san); + sqlite3_exec(db, sql, NULL, NULL, NULL); +} + +static void _member_hash(uint64_t nid, const uint8_t* x25, const uint8_t* ed, + const uint8_t* sig, uint8_t out[HASH_SZ]) { + SHA256_CTX ctx; SHA256_Init(&ctx); + SHA256_Update(&ctx, &nid, 8); + SHA256_Update(&ctx, x25, 32); + SHA256_Update(&ctx, ed, 32); + SHA256_Update(&ctx, sig, 64); + SHA256_Final(out, &ctx); +} + +static uint64_t _mask(int lvl) { int s=63-lvl*5; return s>=0 ? ~0ULL<MAX_LVL)nl=MAX_LVL; + for(int i=0;i0?63-nl*5:0)); + _tree_hash(a,ch,nl,cp,ha[i]);_tree_hash(b,ch,nl,cp,hb[i]); + int ae=1,be=1; for(int j=0;j>i)&1,bh=(bm_b>>i)&1; if(!ah&&!bh)continue; + if(ah&&bh&&memcmp(ha[i],hb[i],HASH_SZ)==0)continue; + uint64_t cp=pref|((uint64_t)i<<(63-nl*5>0?63-nl*5:0)); + char san[64];_sanitize(ch,san,sizeof(san)); int ca=0,cb=0; + sqlite3_stmt* st=NULL; char sql[512]; + snprintf(sql,sizeof(sql),"SELECT COUNT(*) FROM peers_%s WHERE (node_id & %lld)==%lld", + san,(long long)_mask(nl),(long long)cp); + sqlite3_prepare_v2(a,sql,-1,&st,NULL);if(st&&sqlite3_step(st)==SQLITE_ROW)ca=sqlite3_column_int(st,0);if(st)sqlite3_finalize(st); + st=NULL;sqlite3_prepare_v2(b,sql,-1,&st,NULL);if(st&&sqlite3_step(st)==SQLITE_ROW)cb=sqlite3_column_int(st,0);if(st)sqlite3_finalize(st); + if(nl>=MAX_LVL||ca<8||cb<8) { + for(int side=0;side<2;side++) { sqlite3*src=side?a:b,*dst=side?b:a; sqlite3_stmt*r=NULL; + sqlite3_prepare_v2(src,sql,-1,&r,NULL); /* count query already has our WHERE — rebuild for SELECT */ + char sq2[512]; snprintf(sq2,sizeof(sq2), + "SELECT p.node_id,n.x25519_pubkey,n.ed25519_pubkey,p.join_sig" + " FROM peers_%s p JOIN nodes n ON p.node_id=n.node_id" + " WHERE (p.node_id & %lld)==%lld ORDER BY p.node_id", + san,(long long)_mask(nl),(long long)cp); + r=NULL; sqlite3_prepare_v2(src,sq2,-1,&r,NULL); + if(r){while(sqlite3_step(r)==SQLITE_ROW){uint64_t nid=(uint64_t)sqlite3_column_int64(r,0); + const uint8_t* x=sqlite3_column_blob(r,1),*e=sqlite3_column_blob(r,2),*s=sqlite3_column_blob(r,3); + if(x&&e&&s)_insert(dst,ch,nid,x,e,s);} sqlite3_finalize(r);} + } + } else _resolve(a,b,ch,nl,cp); + } +} + +/* ── Tests ── */ +static void t_empty(void) { + printf("t_empty... "); sqlite3 *a,*b; + sqlite3_open(":memory:",&a);sqlite3_open(":memory:",&b); _init_db(a);_init_db(b); _create_peers(a,"e");_create_peers(b,"e"); + uint8_t x[32],e[32],s[64]; _gen(1,x,e,s); _insert(a,"e",1,x,e,s); _insert(b,"e",1,x,e,s); + assert(_count(a,"e")==1&&_count(b,"e")==1); + uint8_t ha[HASH_SZ],hb[HASH_SZ]; _tree_hash(a,"e",1,_mask(1)&1,ha);_tree_hash(b,"e",1,_mask(1)&1,hb); + assert(memcmp(ha,hb,HASH_SZ)==0); sqlite3_close(a);sqlite3_close(b); printf("OK\n"); +} +static void t_1000(void) { + printf("t_1000 (500+500)... "); sqlite3 *a,*b; + sqlite3_open(":memory:",&a);sqlite3_open(":memory:",&b); _init_db(a);_init_db(b); _create_peers(a,"B");_create_peers(b,"B"); + for(int i=0;i<1000;i++){uint64_t id=((uint64_t)(i+1))&0x7FFFFFFFFFFFFFFFULL; uint8_t x[32],e[32],s[64]; _gen(id,x,e,s); + if(i<500)_insert(a,"B",id,x,e,s); else _insert(b,"B",id,x,e,s);} + assert(_count(a,"B")==500&&_count(b,"B")==500); _resolve(a,b,"B",1,0); + printf("A=%d B=%d ",_count(a,"B"),_count(b,"B")); assert(_count(a,"B")==1000&&_count(b,"B")==1000); + for(int i=0;i=200)_insert(b,"o",id,x,e,s);} + assert(_count(a,"o")==500&&_count(b,"o")==500); _resolve(a,b,"o",1,0); + assert(_count(a,"o")==700&&_count(b,"o")==700); + int all_ok=1; for(int i=0;i DbManager::getChannelMembers(const QString& chId, int offset, int limit) const { + QList list; + QString sql = QStringLiteral( + "SELECT p.node_id, COALESCE(n.online,0), COALESCE(n.name,'')" + " FROM \"%1\" p LEFT JOIN nodes n ON p.node_id=n.node_id" + " ORDER BY n.online DESC, p.node_id ASC LIMIT ? OFFSET ?" + ).arg(peersTableName(chId)); + sqlite3_stmt* stmt = prepareOrNull(sql.toUtf8().constData()); + if (!stmt) return list; + sqlite3_bind_int(stmt, 1, limit); + sqlite3_bind_int(stmt, 2, offset); + while (sqlite3_step(stmt) == SQLITE_ROW) { + ChannelMember m; + m.nodeId = (quint64)sqlite3_column_int64(stmt, 0); + m.online = sqlite3_column_int(stmt, 1) != 0; + m.name = colText(stmt, 2); + list.append(m); + } + sqlite3_finalize(stmt); + return list; +} + +QString DbManager::getDisplayName(quint64 nodeId) const { + sqlite3_stmt* st = prepareOrNull( + "SELECT display_name FROM accounts WHERE node_id=?"); + if (st) { + sqlite3_bind_int64(st, 1, (sqlite3_int64)nodeId); + if (sqlite3_step(st) == SQLITE_ROW) { + QString s = colText(st, 0); + sqlite3_finalize(st); + if (!s.isEmpty()) return s; + } else sqlite3_finalize(st); + } + st = prepareOrNull("SELECT name FROM nodes WHERE node_id=?"); + if (st) { + sqlite3_bind_int64(st, 1, (sqlite3_int64)nodeId); + if (sqlite3_step(st) == SQLITE_ROW) { + QString s = colText(st, 0); + sqlite3_finalize(st); + if (!s.isEmpty()) return s; + } else sqlite3_finalize(st); + } + return QString("node_%1").arg(nodeId); } /* ── nodes (запись теперь через topo_node_sqlite) ── */ diff --git a/tools/chatgui/db/db_manager.h b/tools/chatgui/db/db_manager.h index c1b3a1cb..4fdd02ba 100644 --- a/tools/chatgui/db/db_manager.h +++ b/tools/chatgui/db/db_manager.h @@ -46,6 +46,12 @@ struct AccountRow { bool isContact = true; }; +struct ChannelMember { + quint64 nodeId = 0; + bool online = false; + QString name; // from nodes.name, empty if not set +}; + struct NodeAddr { quint64 nodeId = 0; int family = 0; @@ -64,6 +70,7 @@ public: bool isOpen() const { return m_db != nullptr; } quint64 myNodeId() const { return m_myNodeId; } + void setMyNodeId(quint64 nodeId) { m_myNodeId = nodeId; } /* ── Каналы ── */ QList getChannels() const; @@ -80,11 +87,10 @@ public: 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); + /* ── Участники канала (UI) ── */ + int getChannelMemberCount(const QString& chId) const; + QList getChannelMembers(const QString& chId, int offset, int limit) const; + QString getDisplayName(quint64 nodeId) const; /* ── Узлы ── */ bool upsertNode(quint64 nodeId, const QString& name, diff --git a/tools/chatgui/libutun/CMakeLists.txt b/tools/chatgui/libutun/CMakeLists.txt index 24057018..45a55494 100644 --- a/tools/chatgui/libutun/CMakeLists.txt +++ b/tools/chatgui/libutun/CMakeLists.txt @@ -64,6 +64,7 @@ set(UTUN_COMMON_SOURCES ${SRC_DIR}/conn_mgr.c ${TRANSPORT_DIR}/db_sync_stub.c ${TRANSPORT_DIR}/chat_sync.c + ${TRANSPORT_DIR}/member_sync.c ${SRC_DIR}/routing.c ${SRC_DIR}/tun_if.c ${SRC_DIR}/tun_route.c diff --git a/tools/chatgui/src/accountlist.cpp b/tools/chatgui/src/accountlist.cpp index 0569aaa4..17644e5d 100644 --- a/tools/chatgui/src/accountlist.cpp +++ b/tools/chatgui/src/accountlist.cpp @@ -1,39 +1,20 @@ #include "accountlist.h" +#include "memberlistmodel.h" #include "../db/db_manager.h" #include #include -#include -#include - -static QPixmap makeAvatar(const QColor &color, const QString &letter) { - QPixmap pix(24, 24); - pix.fill(Qt::transparent); - QPainter p(&pix); - p.setRenderHint(QPainter::Antialiasing); - p.setBrush(color); - p.setPen(Qt::NoPen); - p.drawEllipse(0, 0, 24, 24); - p.setPen(Qt::white); - QFont f = p.font(); - f.setBold(true); - f.setPointSize(10); - p.setFont(f); - p.drawText(QRect(0, 0, 24, 24), Qt::AlignCenter, letter); - p.end(); - return pix; -} AccountList::AccountList(DbManager* db, QWidget *parent) : QWidget(parent) , m_listView(new QListView(this)) - , m_model(new QStandardItemModel(this)) + , m_model(new MemberListModel(db, this)) , m_db(db) { auto *layout = new QVBoxLayout(this); layout->setContentsMargins(0, 0, 0, 0); layout->setSpacing(0); - auto *header = new QLabel("Accounts", this); + auto *header = new QLabel("Members", this); header->setAlignment(Qt::AlignCenter); header->setStyleSheet("font-weight: bold; padding: 8px; background: palette(window);"); layout->addWidget(header); @@ -47,15 +28,5 @@ AccountList::AccountList(DbManager* db, QWidget *parent) layout->addWidget(m_listView); } -void AccountList::loadAccounts() { - m_model->clear(); - if (!m_db || !m_db->isOpen()) return; - - auto accounts = m_db->getAccounts(true); - for (const auto& a : accounts) { - auto *item = new QStandardItem(); - item->setText(a.displayName); - item->setIcon(QIcon(makeAvatar(QColor(a.avatarColor), a.avatarLetter))); - m_model->appendRow(item); - } -} +void AccountList::setChannel(const QString& channelId) { m_model->setChannel(channelId); } +void AccountList::refresh() { m_model->refresh(); } diff --git a/tools/chatgui/src/accountlist.h b/tools/chatgui/src/accountlist.h index 9b5600e3..5856c0a2 100644 --- a/tools/chatgui/src/accountlist.h +++ b/tools/chatgui/src/accountlist.h @@ -2,8 +2,8 @@ #include #include -#include +class MemberListModel; class DbManager; class AccountList : public QWidget { @@ -11,10 +11,11 @@ class AccountList : public QWidget { public: explicit AccountList(DbManager* db = nullptr, QWidget *parent = nullptr); - void loadAccounts(); + void setChannel(const QString& channelId); + void refresh(); private: - QListView *m_listView; - QStandardItemModel *m_model; - DbManager *m_db; + QListView *m_listView; + MemberListModel *m_model; + DbManager *m_db; }; diff --git a/tools/chatgui/src/channellist.cpp b/tools/chatgui/src/channellist.cpp index 0ddc9f14..c4f976d0 100644 --- a/tools/chatgui/src/channellist.cpp +++ b/tools/chatgui/src/channellist.cpp @@ -99,6 +99,14 @@ ChannelList::ChannelList(DbManager* db, QWidget *parent) } void ChannelList::loadChannels() { + QModelIndex curIdx = m_listView->currentIndex(); + QString selectedCid; + int selectedRow = -1; + if (curIdx.isValid()) { + selectedCid = curIdx.data(ChannelChannelIdRole).toString(); + selectedRow = curIdx.row(); + } + m_model->clear(); if (!m_db || !m_db->isOpen()) return; @@ -121,8 +129,15 @@ void ChannelList::loadChannels() { ci++; } - if (m_model->rowCount() > 0) - m_listView->setCurrentIndex(m_model->index(0, 0)); + if (!selectedCid.isEmpty()) + selectChannel(selectedCid); + if (!m_listView->currentIndex().isValid()) { + int fallbackRow = (selectedRow >= 0) ? selectedRow : 0; + if (fallbackRow >= m_model->rowCount()) + fallbackRow = m_model->rowCount() - 1; + if (fallbackRow >= 0) + m_listView->setCurrentIndex(m_model->index(fallbackRow, 0)); + } } void ChannelList::selectChannel(const QString& channelId) { diff --git a/tools/chatgui/src/mainwindow.cpp b/tools/chatgui/src/mainwindow.cpp index 521ec157..5f190148 100644 --- a/tools/chatgui/src/mainwindow.cpp +++ b/tools/chatgui/src/mainwindow.cpp @@ -40,6 +40,12 @@ static void onChannelUpdatedCallback(const char* ch_id, int ch_id_len) { s_mainWindow->reloadChannels(); } +static void onMembersChangedCallback(const char* ch_id, int ch_id_len) { + Q_UNUSED(ch_id_len); + if (s_mainWindow) + s_mainWindow->onMembersChanged(ch_id); +} + MainWindow::MainWindow(QWidget *parent, DbManager* db, const QString& cfgPath, const QString& dbPath, const QString& debugFile, const QString& debugLevel, const QString& debugCategories) @@ -65,6 +71,11 @@ MainWindow::MainWindow(QWidget *parent, DbManager* db, const QString& cfgPath, s_mainWindow = this; gui_bridge_set_msg_received_cb(onMsgReceivedCallback); gui_bridge_set_channel_updated_cb(onChannelUpdatedCallback); + gui_bridge_set_members_changed_cb(onMembersChangedCallback); + gui_bridge_set_my_node_id_cb([](uint64_t nid) { + qDebug("MainWindow: MY_NODE_ID setMyNodeId(%016llx)", (unsigned long long)nid); + if (s_mainWindow && s_mainWindow->m_db) s_mainWindow->m_db->setMyNodeId((quint64)nid); + }); setupNode(); setupMessaging(); @@ -73,6 +84,7 @@ MainWindow::MainWindow(QWidget *parent, DbManager* db, const QString& cfgPath, this, [this](const QString& channelId) { m_db->setUiState("last_channel_id", channelId); m_messageList->loadChannel(channelId); + m_accountList->setChannel(channelId); }); connect(m_channelList, &ChannelList::inviteRequested, @@ -85,11 +97,12 @@ MainWindow::MainWindow(QWidget *parent, DbManager* db, const QString& cfgPath, this, &MainWindow::showSettings); m_channelList->loadChannels(); - m_accountList->loadAccounts(); QString lastCid = m_db->getUiState("last_channel_id"); - if (!lastCid.isEmpty()) + if (!lastCid.isEmpty()) { m_channelList->selectChannel(lastCid); + m_accountList->setChannel(lastCid); + } } } @@ -239,9 +252,23 @@ void MainWindow::onMessageReceived(const char* ch_id) { m_messageList->refresh(); } +void MainWindow::onMembersChanged(const char* ch_id) { + QString cur = m_db->getUiState("last_channel_id"); + if (cur == QString::fromUtf8(ch_id)) + m_accountList->refresh(); +} + void MainWindow::reloadChannels() { - if (m_channelList) + if (m_channelList) { m_channelList->loadChannels(); + if (!m_pendingChannelSelect.isEmpty()) { + m_channelList->selectChannel(m_pendingChannelSelect); + m_accountList->setChannel(m_pendingChannelSelect); + m_messageList->loadChannel(m_pendingChannelSelect); + m_db->setUiState("last_channel_id", m_pendingChannelSelect); + m_pendingChannelSelect.clear(); + } + } } void MainWindow::onCreateGroupRequested() { @@ -344,6 +371,7 @@ void MainWindow::onCreateGroupRequested() { memcpy(req->signature, sig, 64); gui_bridge_post_uasync_fn(chat_core_create_channel_trampoline, req); + m_pendingChannelSelect = channelId; DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "MainWindow: create group request posted ch=%s name=%s", qPrintable(channelId), qPrintable(name)); diff --git a/tools/chatgui/src/mainwindow.h b/tools/chatgui/src/mainwindow.h index 7886a749..b730d8c6 100644 --- a/tools/chatgui/src/mainwindow.h +++ b/tools/chatgui/src/mainwindow.h @@ -32,6 +32,7 @@ private slots: public: void onMessageReceived(const char* ch_id); + void onMembersChanged(const char* ch_id); void reloadChannels(); private: @@ -52,4 +53,5 @@ private: QString m_cfgDebugFile; QString m_cfgDebugLevel; QString m_cfgDebugCategories; + QString m_pendingChannelSelect; }; diff --git a/tools/chatgui/src/memberlistmodel.cpp b/tools/chatgui/src/memberlistmodel.cpp new file mode 100644 index 00000000..1b054778 --- /dev/null +++ b/tools/chatgui/src/memberlistmodel.cpp @@ -0,0 +1,112 @@ +#include "memberlistmodel.h" +#include "../db/db_manager.h" +#include +#include +#include + +MemberListModel::MemberListModel(DbManager* db, QObject* parent) + : QAbstractListModel(parent), m_db(db) {} + +int MemberListModel::rowCount(const QModelIndex& parent) const { + (void)parent; + return m_totalCount; +} + +QVariant MemberListModel::data(const QModelIndex& index, int role) const { + if (!index.isValid() || index.row() < 0 || index.row() >= m_totalCount) + return {}; + + auto it = m_cache.find(index.row()); + if (it != m_cache.end()) { + if (role == Qt::DisplayRole) return it->name; + if (role == Qt::DecorationRole) return QVariant::fromValue(QIcon(it->avatar)); + if (role == Qt::ToolTipRole) + return QString(it->online ? "online" : "offline") + " | 0x" + + QString("%1").arg(it->nodeId, 16, 16, QChar('0')); + return {}; + } + + if (role == Qt::DisplayRole) { + scheduleBatch(index.row()); + return QStringLiteral("\u2026"); + } + return {}; +} + +void MemberListModel::setChannel(const QString& channelId) { + beginResetModel(); + m_channelId = channelId; + m_totalCount = m_db ? m_db->getChannelMemberCount(channelId) : 0; + m_cache.clear(); + m_pendingFrom = -1; + endResetModel(); + if (m_totalCount > 0) scheduleBatch(0); +} + +void MemberListModel::refresh() { if (!m_channelId.isEmpty()) setChannel(m_channelId); } + +void MemberListModel::scheduleBatch(int centerRow) const { + if (m_pendingFrom >= 0) return; + int from = centerRow - BATCH / 2; + if (from < 0) from = 0; + int maxFrom = m_totalCount - BATCH; + if (from > maxFrom) from = maxFrom > 0 ? maxFrom : 0; + const_cast(this)->m_pendingFrom = from; + const_cast(this)->m_pendingCount = BATCH; + QTimer::singleShot(0, const_cast(this), &MemberListModel::loadPendingBatch); +} + +void MemberListModel::loadPendingBatch() { + if (m_pendingFrom < 0) return; + int from = m_pendingFrom, count = m_pendingCount; + m_pendingFrom = -1; + loadBatch(from, count); +} + +void MemberListModel::loadBatch(int from, int count) { + if (!m_db || m_channelId.isEmpty()) return; + int maxFrom = m_totalCount - count; + if (from > maxFrom) from = maxFrom > 0 ? maxFrom : 0; + + auto members = m_db->getChannelMembers(m_channelId, from, count); + for (int i = 0; i < members.size(); i++) { + int row = from + i; + MemberListItem& item = m_cache[row]; + item.nodeId = members[i].nodeId; + item.online = members[i].online; + QString shortId = QString("%1").arg(members[i].nodeId, 8, 16, QChar('0')).right(4).toUpper(); + item.name = members[i].name.isEmpty() ? shortId : members[i].name; + QColor c = members[i].online ? QColor("#4CAF50") : QColor("#9E9E9E"); + item.avatar = makeAvatar(c, item.name[0]); + } + evictDistant(from + BATCH / 2); + + int last = from + (int)members.size() - 1; + if (last >= from) emit dataChanged(index(from), index(last)); +} + +void MemberListModel::evictDistant(int centerRow) { + if (m_cache.size() <= MAX_CACHE) return; + int half = MAX_CACHE / 2; + QList toRemove; + for (auto it = m_cache.begin(); it != m_cache.end(); ++it) + if (qAbs(it.key() - centerRow) > half) + toRemove.append(it.key()); + for (int k : toRemove) m_cache.remove(k); +} + +QPixmap MemberListModel::makeAvatar(QColor color, QChar letter) const { + QPixmap pm(24, 24); + pm.fill(Qt::transparent); + QPainter p(&pm); + p.setRenderHint(QPainter::Antialiasing); + p.setBrush(color); + p.setPen(Qt::NoPen); + p.drawEllipse(2, 2, 20, 20); + p.setPen(Qt::white); + QFont f; f.setBold(true); f.setPixelSize(12); + p.setFont(f); + p.drawText(QRect(0, 0, 24, 24), Qt::AlignCenter, QString(letter)); + p.end(); + return pm; +} diff --git a/tools/chatgui/src/memberlistmodel.h b/tools/chatgui/src/memberlistmodel.h new file mode 100644 index 00000000..67920d5b --- /dev/null +++ b/tools/chatgui/src/memberlistmodel.h @@ -0,0 +1,51 @@ +#ifndef MEMBERLISTMODEL_H +#define MEMBERLISTMODEL_H + +#include +#include +#include +#include +#include +#include + +class DbManager; + +struct MemberListItem { + quint64 nodeId = 0; + bool online = false; + QPixmap avatar; + QString name; +}; + +class MemberListModel : public QAbstractListModel { + Q_OBJECT +public: + static const int BATCH = 40; + static const int MAX_CACHE = 120; + + explicit MemberListModel(DbManager* db, QObject* parent = nullptr); + + int rowCount(const QModelIndex& parent = QModelIndex()) const override; + QVariant data(const QModelIndex& index, int role = Qt::DisplayRole) const override; + + void setChannel(const QString& channelId); + void refresh(); + +private slots: + void loadPendingBatch(); + +private: + void scheduleBatch(int centerRow) const; + void loadBatch(int from, int count); + void evictDistant(int centerRow); + QPixmap makeAvatar(QColor color, QChar letter) const; + + DbManager* m_db; + QString m_channelId; + int m_totalCount = 0; + mutable QMap m_cache; + mutable int m_pendingFrom = -1; + mutable int m_pendingCount = 0; +}; + +#endif diff --git a/tools/chatgui/src/messagelist.cpp b/tools/chatgui/src/messagelist.cpp index 1e1bbace..abc2c3f5 100644 --- a/tools/chatgui/src/messagelist.cpp +++ b/tools/chatgui/src/messagelist.cpp @@ -125,21 +125,6 @@ MessageList::MessageList(DbManager* db, QWidget *parent) QByteArray textData = text.toUtf8(); auto ts = QDateTime::currentMSecsSinceEpoch(); - uint64_t datahash = 0; - QByteArray chainHash; - bool ok = m_db->insertMessage(m_currentChannelId, m_db->myNodeId(), - "text/plain", textData, ts, true, - &datahash, &chainHash); - if (!ok) { - qWarning("MessageList: insertMessage FAILED ch=%s ts=%lld", - qPrintable(m_currentChannelId), ts); - return; - } - qDebug("MessageList: msg stored ch=%s ts=%lld dh=0x%016llx", - qPrintable(m_currentChannelId), ts, (unsigned long long)datahash); - - loadChannel(m_currentChannelId); - struct chat_msg_submit* req = (struct chat_msg_submit*) u_malloc(sizeof(struct chat_msg_submit) + textData.size()); memset(req, 0, sizeof(*req)); @@ -150,7 +135,7 @@ MessageList::MessageList(DbManager* db, QWidget *parent) req->data_len = (uint32_t)textData.size(); memcpy(req->data, textData.constData(), textData.size()); req->timestamp = (uint64_t)ts; - gui_bridge_post_uasync_fn(chat_core_push_trampoline, req); + gui_bridge_post_uasync_fn(chat_core_submit_trampoline, req); }); connect(AnimTimer::instance(), &AnimTimer::ticked, this, [this]() { @@ -192,8 +177,7 @@ void MessageList::loadChannel(const QString& channelId) { for (const auto& m : msgs) { AccountRow acc = m_db->getAccount(m.authorNodeId); - QString author = acc.displayName.isEmpty() - ? QString("node_%1").arg(m.authorNodeId) : acc.displayName; + QString author = m_db->getDisplayName(m.authorNodeId); QString letter = acc.avatarLetter.isEmpty() ? author.mid(0, 1) : acc.avatarLetter; QColor color(acc.avatarColor); @@ -226,8 +210,7 @@ void MessageList::addMessage(const QString& channelId, quint64 authorNodeId, if (channelId != m_currentChannelId) return; AccountRow acc = m_db ? m_db->getAccount(authorNodeId) : AccountRow(); - QString author = acc.displayName.isEmpty() - ? QString("node_%1").arg(authorNodeId) : acc.displayName; + QString author = m_db ? m_db->getDisplayName(authorNodeId) : QString("node_%1").arg(authorNodeId); QString letter = acc.avatarLetter.isEmpty() ? author.mid(0, 1) : acc.avatarLetter; QColor color(acc.avatarColor); auto avatar = makeAvatar(color, letter); diff --git a/tools/chatgui/src/settingsdialog.cpp b/tools/chatgui/src/settingsdialog.cpp index 031b3319..086896dc 100644 --- a/tools/chatgui/src/settingsdialog.cpp +++ b/tools/chatgui/src/settingsdialog.cpp @@ -1,10 +1,19 @@ // settingsdialog.cpp — Settings dialog with category list #include "settingsdialog.h" #include "networksettingspage.h" +#include "../transport/gui_bridge.h" +#include "../../lib/mem.h" +extern "C" { +#include "../transport/chat_core.h" +} #include #include #include #include +#include +#include +#include +#include SettingsDialog::SettingsDialog(const QString& configPath, QWidget* parent) : QDialog(parent) @@ -26,6 +35,7 @@ SettingsDialog::SettingsDialog(const QString& configPath, QWidget* parent) " background: palette(window); font-size: 13px; }" "QListWidget::item { padding: 10px 14px; }" "QListWidget::item:selected { background: #3390EC; color: white; }"); + m_categoryList->addItem("Profile"); m_categoryList->addItem("Network"); mainLayout->addWidget(m_categoryList); @@ -36,6 +46,22 @@ SettingsDialog::SettingsDialog(const QString& configPath, QWidget* parent) rightLayout->setSpacing(0); m_pages = new QStackedWidget(rightPanel); + + // Profile page + auto* profilePage = new QWidget(m_pages); + auto* profLayout = new QFormLayout(profilePage); + profLayout->setContentsMargins(0, 12, 0, 0); + auto* nameLabel = new QLabel("Display name:", profilePage); + m_nicknameEdit = new QLineEdit(profilePage); + m_nicknameEdit->setPlaceholderText("Your name (shown to others)"); + m_nicknameEdit->setMaxLength(32); + m_nicknameEdit->setStyleSheet( + "QLineEdit { border: 1px solid palette(mid); border-radius: 4px;" + " padding: 6px 10px; font-size: 13px; background: palette(base); }"); + profLayout->addRow(nameLabel, m_nicknameEdit); + m_pages->addWidget(profilePage); + + // Network page m_networkPage = new NetworkSettingsPage(m_pages); m_pages->addWidget(m_networkPage); rightLayout->addWidget(m_pages, 1); @@ -71,15 +97,68 @@ SettingsDialog::SettingsDialog(const QString& configPath, QWidget* parent) connect(m_categoryList, &QListWidget::currentRowChanged, this, &SettingsDialog::onCategoryChanged); + // Load current values + loadProfileFromConfig(m_configPath); m_networkPage->loadFromConfig(m_configPath); m_categoryList->setCurrentRow(0); } +void SettingsDialog::loadProfileFromConfig(const QString& path) { + QFile f(path); + if (!f.open(QIODevice::ReadOnly | QIODevice::Text)) return; + QTextStream in(&f); + while (!in.atEnd()) { + QString line = in.readLine().trimmed(); + if (line.startsWith("my_node_name=")) { + QString nm = line.mid(QString("my_node_name=").size()).trimmed(); + m_nicknameEdit->setText(nm); + m_oldNickname = nm; + break; + } + } +} + void SettingsDialog::onCategoryChanged(int row) { m_pages->setCurrentIndex(row); } void SettingsDialog::onSave() { + QString newName = m_nicknameEdit->text().trimmed(); + + /* update name in uTun if changed */ + if (newName != m_oldNickname) { + char* name_copy = u_strdup(newName.toUtf8().constData()); + if (name_copy) gui_bridge_post_uasync_fn(chat_core_update_my_name_trampoline, name_copy); + m_oldNickname = newName; + } + + // Save nickname to config + QFile f(m_configPath); + if (!f.open(QIODevice::ReadOnly | QIODevice::Text)) { accept(); return; } + QString content = f.readAll(); + f.close(); + + QTextStream ts(&content); + QStringList lines; + while (!ts.atEnd()) lines << ts.readLine(); + + bool found = false; + for (int i = 0; i < lines.size(); i++) { + if (lines[i].trimmed().startsWith("my_node_name=")) { + lines[i] = "my_node_name=" + newName; + found = true; + } else if (!found && lines[i].trimmed().startsWith("[gui]")) { + // insert before [gui] if not found yet + lines.insert(i, "my_node_name=" + newName); + found = true; + } + } + if (!found) lines << "my_node_name=" + newName; + + f.open(QIODevice::WriteOnly | QIODevice::Truncate | QIODevice::Text); + f.write(lines.join('\n').toUtf8()); + f.close(); + m_networkPage->saveToConfig(m_configPath); accept(); } diff --git a/tools/chatgui/src/settingsdialog.h b/tools/chatgui/src/settingsdialog.h index ec8eadb0..f59e3d2c 100644 --- a/tools/chatgui/src/settingsdialog.h +++ b/tools/chatgui/src/settingsdialog.h @@ -5,6 +5,7 @@ #include #include #include +#include #include class NetworkSettingsPage; @@ -19,9 +20,13 @@ private slots: void onSave(); private: + void loadProfileFromConfig(const QString& path); + QString m_configPath; QListWidget* m_categoryList; QStackedWidget* m_pages; QPushButton* m_saveBtn; + QLineEdit* m_nicknameEdit; + QString m_oldNickname; NetworkSettingsPage* m_networkPage; }; diff --git a/tools/chatgui/transport/chat_core.c b/tools/chatgui/transport/chat_core.c index 6b09089f..eccdb425 100644 --- a/tools/chatgui/transport/chat_core.c +++ b/tools/chatgui/transport/chat_core.c @@ -9,6 +9,7 @@ #include "chat_sync.h" #include "gui_bridge.h" #include "topo_node_sqlite.h" +#include "member_sync.h" #include "../../../src/utun_instance.h" #include "../../../src/etcp_router.h" @@ -26,6 +27,7 @@ #include #include +#include #include #include @@ -134,6 +136,33 @@ int chat_core_init(struct UTUN_INSTANCE* inst, const char* db_path) { topo_node_sqlite_init(g_cc.db); topo_groups_set_sqlite_db(inst->topo_groups, g_cc.db); + /* записать себя в nodes + local_identity с реальным node_id и именем */ + { + const char* my_name = inst->name[0] ? inst->name : "Me"; + sqlite3_stmt* st = NULL; + sqlite3_prepare_v2(g_cc.db, + "INSERT OR REPLACE INTO nodes(node_id, name, x25519_pubkey, ed25519_pubkey, online)" + " VALUES(?,?,?,?,1)", -1, &st, NULL); + if (st) { + sqlite3_bind_int64(st, 1, (sqlite3_int64)g_cc.my_node_id); + sqlite3_bind_text(st, 2, my_name, -1, SQLITE_STATIC); + sqlite3_bind_blob(st, 3, inst->my_keys.public_key, 32, SQLITE_STATIC); + sqlite3_bind_blob(st, 4, inst->my_ed25519_pubkey, 32, SQLITE_STATIC); + sqlite3_step(st); sqlite3_finalize(st); + } + st = NULL; + sqlite3_prepare_v2(g_cc.db, + "INSERT OR REPLACE INTO local_identity(id,node_id,name,x25519_pubkey,ed25519_pubkey)" + " VALUES(1,?,?,?,?)", -1, &st, NULL); + if (st) { + sqlite3_bind_int64(st, 1, (sqlite3_int64)g_cc.my_node_id); + sqlite3_bind_text(st, 2, my_name, -1, SQLITE_STATIC); + sqlite3_bind_blob(st, 3, inst->my_keys.public_key, 32, SQLITE_STATIC); + sqlite3_bind_blob(st, 4, inst->my_ed25519_pubkey, 32, SQLITE_STATIC); + sqlite3_step(st); sqlite3_finalize(st); + } + } + db_exec( "CREATE TABLE IF NOT EXISTS local_identity (" " id INTEGER PRIMARY KEY CHECK (id = 1)," @@ -164,6 +193,9 @@ int chat_core_init(struct UTUN_INSTANCE* inst, const char* db_path) { g_cc.initialized = 1; DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: initialized, db=%s node_id=0x%016llx", CC_ID, db_path, (unsigned long long)g_cc.my_node_id); + + { uint8_t nid[8]; memcpy(nid, &g_cc.my_node_id, 8); gui_bridge_post(GUI_EVT_MY_NODE_ID, nid, 8); } + return 0; } @@ -184,6 +216,73 @@ void chat_core_set_my_node_id(uint64_t node_id) { g_cc.my_node_id = node_id; } +void chat_core_update_my_name(const char* name) { + if (!g_cc.initialized || !name) return; + sqlite3_stmt* st = NULL; + sqlite3_prepare_v2(g_cc.db, "UPDATE nodes SET name=? WHERE node_id=?", -1, &st, NULL); + if (st) { + sqlite3_bind_text(st, 1, name, -1, SQLITE_STATIC); + sqlite3_bind_int64(st, 2, (sqlite3_int64)g_cc.my_node_id); + sqlite3_step(st); sqlite3_finalize(st); + } + /* update instance name for NODEINFO */ + snprintf(g_cc.inst->name, sizeof(g_cc.inst->name), "%s", name); + + /* recompute join_sigs for all channels where I'm a member */ + uint64_t myid = g_cc.my_node_id; + sqlite3_stmt* cs = NULL; + sqlite3_prepare_v2(g_cc.db, "SELECT channel_id FROM channels", -1, &cs, NULL); + if (cs) { + while (sqlite3_step(cs) == SQLITE_ROW) { + const char* ch = (const char*)sqlite3_column_text(cs, 0); + if (!ch) continue; + char tbl[80]; peers_table_name(ch, tbl, sizeof(tbl)); + char buf[256]; snprintf(buf, sizeof(buf), "SELECT 1 FROM \"%s\" WHERE node_id=?", tbl); + sqlite3_stmt* ps = NULL; + if (sqlite3_prepare_v2(g_cc.db, buf, -1, &ps, NULL) == SQLITE_OK) { + sqlite3_bind_int64(ps, 1, (sqlite3_int64)myid); + if (sqlite3_step(ps) == SQLITE_ROW) { + uint8_t join_msg[256]; size_t mlen = 0; + mlen += snprintf((char*)join_msg + mlen, sizeof(join_msg) - mlen, "%s", ch) + 1; + memcpy(join_msg + mlen, &myid, 8); mlen += 8; + memcpy(join_msg + mlen, g_cc.inst->my_keys.public_key, 32); mlen += 32; + { const char* nm = name[0] ? name : ""; + size_t nl = strlen(nm); memcpy(join_msg + mlen, nm, nl); mlen += nl; join_msg[mlen++] = '\0'; } + uint8_t new_sig[64]; memset(new_sig, 0, 64); + EVP_PKEY* pkey = EVP_PKEY_new_raw_private_key(EVP_PKEY_ED25519, NULL, + g_cc.inst->my_ed25519_privkey, 32); + if (pkey) { + EVP_MD_CTX* mdctx = EVP_MD_CTX_new(); + if (mdctx) { + if (EVP_DigestSignInit(mdctx, NULL, NULL, NULL, pkey) == 1) + EVP_DigestSign(mdctx, new_sig, &(size_t){64}, join_msg, mlen); + EVP_MD_CTX_free(mdctx); + } + EVP_PKEY_free(pkey); + } + topo_node_sqlite_member_put(g_cc.db, ch, myid, new_sig, NULL); + member_sync_put(g_cc.inst, ch, myid, g_cc.inst->my_keys.public_key, + g_cc.inst->my_ed25519_pubkey, new_sig, NULL, 0); + uint8_t evt[65]; uint8_t cl = (uint8_t)strlen(ch); + evt[0] = cl; memcpy(evt + 1, ch, cl); + gui_bridge_post(GUI_EVT_MEMBERS_CHANGED, evt, 1 + cl); + } + sqlite3_finalize(ps); + } + } + sqlite3_finalize(cs); + } + + /* broadcast updated NODEINFO to connected peers */ + struct TOPO_GROUP* grp = topo_groups_get_default(g_cc.inst->topo_groups); + if (grp) topo_group_update_my_nodeinfo(g_cc.inst, grp); +} + +void chat_core_update_my_name_trampoline(void* arg) { + chat_core_update_my_name((const char*)arg); + u_free(arg); +} + /* ─── отправка сообщения (GUI → uasync) ─── */ void chat_core_submit_message(struct chat_msg_submit* req) { @@ -234,6 +333,10 @@ void chat_core_submit_message(struct chat_msg_submit* req) { DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: msg inserted ch=%s ts=%llu dh=0x%016llx", CC_ID, ch_id, (unsigned long long)ts, (unsigned long long)dh); + /* notify GUI to reload messages */ + { uint8_t evt[65]; uint8_t cl = (uint8_t)strlen(ch_id); evt[0] = cl; memcpy(evt + 1, ch_id, cl); + gui_bridge_post(GUI_EVT_MSG_RECEIVED, evt, 1 + cl); } + /* push пирам */ chat_sync_push(g_cc.inst, ch_id, g_cc.my_node_id, req->content_type, data, data_len, ts, dh); @@ -866,14 +969,43 @@ void chat_core_create_channel(struct chat_channel_create* req) { DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: channel_put OK ch=%s name=%s owner=0x%016llx", CC_ID, req->channel_id, req->name, (unsigned long long)req->owner_node_id); - /* уведомляем GUI */ uint8_t ch_id_len = (uint8_t)strlen(req->channel_id); uint8_t data[65]; data[0] = ch_id_len; memcpy(data + 1, req->channel_id, ch_id_len); + + /* add self as first member BEFORE posting channel to GUI */ + { + uint64_t myid = g_cc.inst->node_id; + uint8_t join_msg[256]; size_t mlen = 0; + mlen += snprintf((char*)join_msg + mlen, sizeof(join_msg) - mlen, "%s", req->channel_id) + 1; + memcpy(join_msg + mlen, &myid, 8); mlen += 8; + memcpy(join_msg + mlen, g_cc.inst->my_keys.public_key, 32); mlen += 32; + { const char* nm = g_cc.inst->name[0] ? g_cc.inst->name : ""; + size_t nl = strlen(nm); memcpy(join_msg + mlen, nm, nl); mlen += nl; join_msg[mlen++] = '\0'; } + uint8_t join_sig[64]; + EVP_PKEY* pkey = EVP_PKEY_new_raw_private_key(EVP_PKEY_ED25519, NULL, + g_cc.inst->my_ed25519_privkey, 32); + if (pkey) { + EVP_MD_CTX* mdctx = EVP_MD_CTX_new(); + if (mdctx) { + if (EVP_DigestSignInit(mdctx, NULL, NULL, NULL, pkey) == 1) + EVP_DigestSign(mdctx, join_sig, &(size_t){64}, join_msg, mlen); + else + memset(join_sig, 0, 64); + EVP_MD_CTX_free(mdctx); + } + EVP_PKEY_free(pkey); + } else { + memset(join_sig, 0, 64); + } + member_sync_put(g_cc.inst, req->channel_id, myid, + g_cc.inst->my_keys.public_key, g_cc.inst->my_ed25519_pubkey, + join_sig, NULL, 0); + } + + /* notify GUI — members already in DB, channel will show with self as participant */ gui_bridge_post(GUI_EVT_CHANNEL_UPDATED, data, 1 + ch_id_len); - DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: gui_bridge_post sent GUI_EVT_CHANNEL_UPDATED ch=%s", - CC_ID, req->channel_id); } } diff --git a/tools/chatgui/transport/chat_core.h b/tools/chatgui/transport/chat_core.h index afe08749..c47a18b8 100644 --- a/tools/chatgui/transport/chat_core.h +++ b/tools/chatgui/transport/chat_core.h @@ -81,6 +81,8 @@ void chat_core_create_channel_trampoline(void* arg); /* ── Утилиты ── */ void chat_core_set_my_node_id(uint64_t node_id); +void chat_core_update_my_name(const char* name); +void chat_core_update_my_name_trampoline(void* arg); /* Трамплин для gui_bridge_post_uasync (GUI → uasync) */ void chat_core_submit_trampoline(void* arg); diff --git a/tools/chatgui/transport/chat_sync.c b/tools/chatgui/transport/chat_sync.c index 0609f6ce..e16491dd 100644 --- a/tools/chatgui/transport/chat_sync.c +++ b/tools/chatgui/transport/chat_sync.c @@ -2,6 +2,7 @@ #include "chat_core.h" #include "gui_bridge.h" #include "topo_node_sqlite.h" +#include "member_sync.h" #include "../../../src/utun_instance.h" #include "../../../src/etcp_router.h" @@ -404,7 +405,9 @@ static void cs_on_conn_up(struct ETCP_CONN* conn, void* arg) { memcpy(msg + 1, &ch->msg_count, 4); memcpy(msg + 5, ch->last_chain_hash, 32); cs_send(g_cs, ch->channel_id, peer, msg, 37); + member_sync_start(g_cs->inst, peer, ch->channel_id); } + member_sync_set_online(g_cs->inst, peer, 1); } static void cs_on_conn_down(struct ETCP_CONN* conn, void* arg) { @@ -532,6 +535,8 @@ int chat_sync_init(struct UTUN_INSTANCE* inst, cs->info_req_timer = NULL; cs->join_timer = NULL; + member_sync_init(inst); + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: initialized", CS_ID); return 0; } @@ -539,6 +544,7 @@ int chat_sync_init(struct UTUN_INSTANCE* inst, void chat_sync_destroy(struct UTUN_INSTANCE* inst) { struct chat_sync* cs = g_cs; if (!cs || !inst) return; + member_sync_destroy(inst); cs->initialized = 0; g_cs = NULL; etcp_router_unbind(inst, ETCP_RT_ID_CHAT_SYNC); @@ -679,6 +685,22 @@ static void cs_propagate(struct chat_sync* cs, const char* ch_id, uint64_t exclu } } +static int _get_node_name(sqlite3* db, uint64_t node_id, char* out, size_t sz) { + if (!db || !out || !sz) return -1; + sqlite3_stmt* st = NULL; + sqlite3_prepare_v2(db, "SELECT name FROM nodes WHERE node_id=?", -1, &st, NULL); + if (!st) return -1; + sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); + int rc = -1; + if (sqlite3_step(st) == SQLITE_ROW) { + const unsigned char* nm = sqlite3_column_text(st, 0); + if (nm) { snprintf(out, sz, "%s", nm); rc = 0; } + else { out[0] = '\0'; rc = 0; } + } + sqlite3_finalize(st); + return rc; +} + /* ─── CHANNEL_INFO_REQ (0x09): joiner → inviter ─── */ static void cs_handle_channel_info_req(struct chat_sync* cs, uint64_t peer, @@ -697,6 +719,8 @@ static void cs_handle_channel_info_req(struct chat_sync* cs, uint64_t peer, mlen += snprintf((char*)join_msg + mlen, sizeof(join_msg) - mlen, "%s", ch_id) + 1; memcpy(join_msg + mlen, &myid, 8); mlen += 8; memcpy(join_msg + mlen, cs->inst->my_keys.public_key, 32); mlen += 32; + { const char* nm = cs->inst->name[0] ? cs->inst->name : ""; + size_t nl = strlen(nm); memcpy(join_msg + mlen, nm, nl); mlen += nl; join_msg[mlen++] = '\0'; } cs_ed25519_sign(cs->inst->my_ed25519_privkey, join_msg, mlen, my_join_sig); uint8_t buf[1024]; size_t boff = 0; @@ -744,6 +768,28 @@ static void cs_handle_channel_info_resp(struct chat_sync* cs, uint64_t peer, topo_node_sqlite_channel_put(cs->inst->topo_groups->topo_sqlite_db, ch_id, name, (int)is_dm, owner, x25519, NULL, ed_pub, NULL, ch_sig); + /* verify inviter's join_sig */ + { sqlite3* vdb = cs->inst->topo_groups->topo_sqlite_db; + /* read inviter's Ed25519 pubkey */ + uint8_t inv_ed[32] = {0}; + sqlite3_stmt* es = NULL; + sqlite3_prepare_v2(vdb, "SELECT ed25519_pubkey FROM nodes WHERE node_id=?", -1, &es, NULL); + if (es) { sqlite3_bind_int64(es, 1, (sqlite3_int64)peer); + if (sqlite3_step(es) == SQLITE_ROW) memcpy(inv_ed, sqlite3_column_blob(es, 0), 32); + sqlite3_finalize(es); } + /* verify */ + uint8_t ivmsg[256]; size_t ilen = 0; + ilen += snprintf((char*)ivmsg + ilen, sizeof(ivmsg) - ilen, "%s", ch_id) + 1; + memcpy(ivmsg + ilen, &peer, 8); ilen += 8; + memcpy(ivmsg + ilen, x25519, 32); ilen += 32; + { char nnm[64] = ""; _get_node_name(vdb, peer, nnm, sizeof(nnm)); + size_t nl = strlen(nnm); memcpy(ivmsg + ilen, nnm, nl); ilen += nl; ivmsg[ilen++] = '\0'; } + if (cs_ed25519_verify(inv_ed, ivmsg, ilen, inviter_join_sig) != 0) { + DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: CHANNEL_INFO_RESP invalid inviter_join_sig peer=%016llx", + CS_ID, (unsigned long long)peer); + } + } + /* save inviter as node and member */ topo_node_sqlite_member_put(cs->inst->topo_groups->topo_sqlite_db, ch_id, peer, inviter_join_sig, inviter_join_sig); @@ -758,6 +804,8 @@ static void cs_handle_channel_info_resp(struct chat_sync* cs, uint64_t peer, mlen += snprintf((char*)msg + mlen, sizeof(msg) - mlen, "%s", ch_id) + 1; memcpy(msg + mlen, &myid, 8); mlen += 8; memcpy(msg + mlen, my_x25519, 32); mlen += 32; + { const char* nm = cs->inst->name[0] ? cs->inst->name : ""; + size_t nl = strlen(nm); memcpy(msg + mlen, nm, nl); mlen += nl; msg[mlen++] = '\0'; } cs_ed25519_sign(cs->inst->my_ed25519_privkey, msg, mlen, join_sig); } @@ -823,6 +871,8 @@ static void cs_handle_channel_join(struct chat_sync* cs, uint64_t peer, vlen += snprintf((char*)vmsg + vlen, sizeof(vmsg) - vlen, "%s", ch_id) + 1; memcpy(vmsg + vlen, &node_id, 8); vlen += 8; memcpy(vmsg + vlen, x25519, 32); vlen += 32; + { sqlite3* tdb = cs->inst->topo_groups->topo_sqlite_db; char nnm[64] = ""; _get_node_name(tdb, node_id, nnm, sizeof(nnm)); + size_t nl = strlen(nnm); memcpy(vmsg + vlen, nnm, nl); vlen += nl; vmsg[vlen++] = '\0'; } if (cs_ed25519_verify(ed_pub, vmsg, vlen, join_sig) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: JOIN invalid sig node=0x%016llx ch=%s", CS_ID, (unsigned long long)node_id, ch_id); @@ -934,8 +984,10 @@ static void cs_handle_welcome(struct chat_sync* cs, uint64_t peer, sqlite3_bind_blob(stmt, 3, ip, ip_len, SQLITE_STATIC); sqlite3_bind_int(stmt, 4, (int)port); sqlite3_step(stmt); sqlite3_finalize(stmt); - } - } + } + member_sync_cancel_peer(g_cs->inst, peer); + member_sync_set_online(g_cs->inst, peer, 0); +} } cs_refresh_channels(cs); @@ -944,6 +996,8 @@ static void cs_handle_welcome(struct chat_sync* cs, uint64_t peer, evt[0] = ch_id_len; memcpy(evt + 1, ch_id, ch_id_len); gui_bridge_post(GUI_EVT_CHANNEL_UPDATED, evt, 1 + ch_id_len); + gui_bridge_post(GUI_EVT_MEMBERS_CHANGED, evt, 1 + ch_id_len); + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: WELCOME processed ch=%s peers=%d", CS_ID, ch_id, pc); } @@ -965,6 +1019,8 @@ static void cs_handle_peer_upsert(struct chat_sync* cs, uint64_t peer, vlen += snprintf((char*)vmsg + vlen, sizeof(vmsg) - vlen, "%s", ch_id) + 1; memcpy(vmsg + vlen, &node_id, 8); vlen += 8; memcpy(vmsg + vlen, x25519, 32); vlen += 32; + { sqlite3* tdb = cs->inst->topo_groups->topo_sqlite_db; char nnm[64] = ""; _get_node_name(tdb, node_id, nnm, sizeof(nnm)); + size_t nl = strlen(nnm); memcpy(vmsg + vlen, nnm, nl); vlen += nl; vmsg[vlen++] = '\0'; } if (cs_ed25519_verify(ed_pub, vmsg, vlen, join_sig) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: PEER_UPSERT invalid sig node=0x%016llx", CS_ID, (unsigned long long)node_id); @@ -972,6 +1028,7 @@ static void cs_handle_peer_upsert(struct chat_sync* cs, uint64_t peer, } sqlite3* db = cs->inst->topo_groups->topo_sqlite_db; + topo_node_sqlite_member_put(db, ch_id, node_id, join_sig, NULL); for (uint8_t i = 0; i < ac; i++) { @@ -997,6 +1054,8 @@ static void cs_handle_peer_upsert(struct chat_sync* cs, uint64_t peer, cs_refresh_channels(cs); + { uint8_t mevt[65]; uint8_t ml = (uint8_t)strlen(ch_id); mevt[0] = ml; memcpy(mevt + 1, ch_id, ml); gui_bridge_post(GUI_EVT_MEMBERS_CHANGED, mevt, 1 + ml); } + /* propagate to others (except sender and the subject node) */ cs_propagate(cs, ch_id, peer, pl, len); @@ -1015,6 +1074,8 @@ static void cs_handle_peer_remove(struct chat_sync* cs, uint64_t peer, topo_node_sqlite_member_del(db, ch_id, node_id); cs_refresh_channels(cs); + { uint8_t mevt[65]; uint8_t ml = (uint8_t)strlen(ch_id); mevt[0] = ml; memcpy(mevt + 1, ch_id, ml); gui_bridge_post(GUI_EVT_MEMBERS_CHANGED, mevt, 1 + ml); } + cs_propagate(cs, ch_id, peer, pl, len); DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: PEER_REMOVE node=0x%016llx ch=%s", diff --git a/tools/chatgui/transport/chat_sync.h b/tools/chatgui/transport/chat_sync.h index 21b8ed43..077e520f 100644 --- a/tools/chatgui/transport/chat_sync.h +++ b/tools/chatgui/transport/chat_sync.h @@ -12,7 +12,8 @@ struct UTUN_INSTANCE; struct UASYNC; /* etcp_router service ID */ -#define ETCP_RT_ID_CHAT_SYNC 0x30 +#define ETCP_RT_ID_CHAT_SYNC 0x30 +#define ETCP_RT_ID_MEMBER_SYNC 0x31 /* Message types */ #define CS_MSG_INIT_SYNC 0x01 diff --git a/tools/chatgui/transport/gui_bridge.h b/tools/chatgui/transport/gui_bridge.h index 692ca4da..c9c9e122 100644 --- a/tools/chatgui/transport/gui_bridge.h +++ b/tools/chatgui/transport/gui_bridge.h @@ -16,6 +16,8 @@ struct UASYNC; #define GUI_EVT_CONNECT_RESULT 2 /* data: [node_id:8][result:4][channel_id:8] */ #define GUI_EVT_NEW_PEER 3 /* data: [node_id:8] */ #define GUI_EVT_CHANNEL_UPDATED 4 /* data: [ch_id_len:1][ch_id:var] */ +#define GUI_EVT_MEMBERS_CHANGED 5 /* data: [ch_id_len:1][ch_id:var] */ +#define GUI_EVT_MY_NODE_ID 6 /* data: [node_id:8] */ /* ── API ── */ @@ -46,6 +48,14 @@ void gui_bridge_set_msg_received_cb(gui_msg_received_fn cb); typedef void (*gui_channel_updated_fn)(const char* ch_id, int ch_id_len); void gui_bridge_set_channel_updated_cb(gui_channel_updated_fn cb); +/* Callback для изменения списка участников (вызывается из GUI-потока) */ +typedef void (*gui_members_changed_fn)(const char* ch_id, int ch_id_len); +void gui_bridge_set_members_changed_cb(gui_members_changed_fn cb); + +/* Callback для получения своего node_id после инициализации uTun */ +typedef void (*gui_my_node_id_fn)(uint64_t node_id); +void gui_bridge_set_my_node_id_cb(gui_my_node_id_fn cb); + #ifdef __cplusplus } #endif diff --git a/tools/chatgui/transport/gui_bridge_impl.cpp b/tools/chatgui/transport/gui_bridge_impl.cpp index 98272695..27c9a2e9 100644 --- a/tools/chatgui/transport/gui_bridge_impl.cpp +++ b/tools/chatgui/transport/gui_bridge_impl.cpp @@ -26,7 +26,9 @@ public: static GuiBridgeReceiver* g_receiver = nullptr; static gui_connect_result_fn g_connect_result_cb = nullptr; static gui_msg_received_fn g_msg_received_cb = nullptr; -static gui_channel_updated_fn g_channel_updated_cb = nullptr; +static gui_channel_updated_fn g_channel_updated_cb = nullptr; +static gui_members_changed_fn g_members_changed_cb = nullptr; +static gui_my_node_id_fn g_my_node_id_cb = nullptr; static struct UASYNC* g_ua = nullptr; /* ── GuiBridgeReceiver implementation ── */ @@ -81,6 +83,21 @@ void GuiBridgeReceiver::processPost(int eventType, QByteArray data) { DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "gui_bridge: CHANNEL_UPDATED data too short %d", dlen); } break; + case GUI_EVT_MEMBERS_CHANGED: + if (dlen >= 2) { + uint8_t chLen = d[0]; + if (dlen >= 1 + chLen) { + QString chId = QString::fromUtf8((const char*)d + 1, chLen); + if (g_members_changed_cb) g_members_changed_cb((const char*)d + 1, chLen); + } + } + break; + case GUI_EVT_MY_NODE_ID: + if (dlen >= 8) { + uint64_t nid; memcpy(&nid, d, 8); + if (g_my_node_id_cb) g_my_node_id_cb(nid); + } + break; default: DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "gui_bridge: unknown event type %d", eventType); break; @@ -139,6 +156,14 @@ void gui_bridge_set_channel_updated_cb(gui_channel_updated_fn cb) { g_channel_updated_cb = cb; } +void gui_bridge_set_members_changed_cb(gui_members_changed_fn cb) { + g_members_changed_cb = cb; +} + +void gui_bridge_set_my_node_id_cb(gui_my_node_id_fn cb) { + g_my_node_id_cb = cb; +} + } /* extern "C" */ #include "gui_bridge_impl.moc" diff --git a/tools/chatgui/transport/member_sync.c b/tools/chatgui/transport/member_sync.c new file mode 100644 index 00000000..3461e8f7 --- /dev/null +++ b/tools/chatgui/transport/member_sync.c @@ -0,0 +1,940 @@ +#include "member_sync.h" +#include "topo_node_sqlite.h" + +#include "../../../src/utun_instance.h" +#include "../../../src/etcp_router.h" +#include "../../../src/etcp_api.h" +#include "../../../src/etcp.h" +#include "../../../src/topo_group.h" +#include "../../../lib/debug_config.h" +#include "../../../lib/mem.h" +#include "../../../lib/u_async.h" + +#include +#include +#include +#include + +#define MS_ID "member_sync" +#define MS_SYNC_TIMEOUT_MS 10000 +#define MS_BG_INTERVAL_MS 100 + +static struct member_sync* g_ms = NULL; + +struct bucket_entry { + uint8_t level; + uint8_t prefix_bytes; + uint64_t prefix; +}; + +struct addr_item { + uint8_t family; + uint8_t addr[16]; + uint16_t port; +}; + +struct ms_session { + struct ms_session* next; + char ch_id[64]; + uint64_t peer; + uint8_t active; + void* timer; + uint32_t started_tb; + uint8_t retries; +}; + +struct member_sync { + struct UTUN_INSTANCE* inst; + struct ms_session* sessions; + uint8_t initialized; + void* bg_timer; +}; + +static sqlite3* _db(struct UTUN_INSTANCE* inst) { + return inst && inst->topo_groups ? inst->topo_groups->topo_sqlite_db : NULL; +} + +static uint64_t _level_prefix(uint64_t member_id, uint8_t level) { + int shift = 63 - (int)level * 5; + if (shift < 0) shift = 0; + return (member_id >> shift) << shift; +} + +static uint8_t _prefix_bytes(uint8_t level) { + int bits = level * 5; + return (uint8_t)((bits + 7) / 8); +} + +static void _prefix_write(uint8_t* out, uint64_t prefix, uint8_t pb) { + for (int i = (int)pb - 1; i >= 0; i--) + out[pb - 1 - i] = (uint8_t)(prefix >> (i * 8)); +} + +static uint64_t _prefix_read(const uint8_t* in, uint8_t pb) { + uint64_t v = 0; + for (uint8_t i = 0; i < pb && i < 8; i++) v = (v << 8) | in[i]; + return v; +} + +static int _addr_cmp(const void* a, const void* b) { + const struct addr_item* ia = (const struct addr_item*)a; + const struct addr_item* ib = (const struct addr_item*)b; + if (ia->family != ib->family) return (int)ia->family - (int)ib->family; + int r = memcmp(ia->addr, ib->addr, (size_t)(ia->family == 4 ? 4 : 16)); + if (r) return r; + return (int)ia->port - (int)ib->port; +} + +static void _compute_member_hash(uint64_t node_id, const uint8_t* x25519, + const uint8_t* ed25519, const uint8_t* join_sig, + const uint8_t* addrs_data, int addr_count, int online, + uint8_t hash_out[MS_HASH_SIZE]) { + SHA256_CTX ctx; + SHA256_Init(&ctx); + SHA256_Update(&ctx, &node_id, 8); + SHA256_Update(&ctx, x25519, 32); + SHA256_Update(&ctx, ed25519, 32); + SHA256_Update(&ctx, join_sig, 64); + uint8_t ac = (uint8_t)addr_count; + SHA256_Update(&ctx, &ac, 1); + + if (addr_count > 0 && addrs_data) { + struct addr_item items[256]; int n = 0; + const uint8_t* p = addrs_data; + for (int i = 0; i < addr_count && n < 256; i++) { + uint8_t fam = *p++; + items[n].family = fam; + int ip_len = fam == 4 ? 4 : 16; + memcpy(items[n].addr, p, (size_t)ip_len); p += ip_len; + items[n].port = ((uint16_t)p[0] << 8) | p[1]; p += 2; + n++; + } + qsort(items, (size_t)n, sizeof(struct addr_item), _addr_cmp); + for (int i = 0; i < n; i++) { + SHA256_Update(&ctx, &items[i].family, 1); + SHA256_Update(&ctx, items[i].addr, (size_t)(items[i].family == 4 ? 4 : 16)); + uint8_t port_be[2] = { (uint8_t)(items[i].port >> 8), (uint8_t)(items[i].port & 0xFF) }; + SHA256_Update(&ctx, port_be, 2); + } + } + + uint8_t ol = (uint8_t)(online ? 1 : 0); + SHA256_Update(&ctx, &ol, 1); + SHA256_Final(hash_out, &ctx); +} + +static void _peers_table(const char* ch_id, char* buf, size_t sz) { + size_t i = 0; + while (*ch_id && i < sz - 1) { + char c = *ch_id++; + if ((c >= 'a' && c <= 'z') || (c >= 'A' && c <= 'Z') || (c >= '0' && c <= '9') || c == '_') + buf[i++] = c; + else + buf[i++] = '_'; + } + buf[i] = '\0'; + char tbl[128]; snprintf(tbl, sizeof(tbl), "peers_%s", buf); + snprintf(buf, sz, "%s", tbl); +} + +static int _recompute_bucket(struct UTUN_INSTANCE* inst, const char* ch_id, + uint8_t level, uint64_t prefix64) { + sqlite3* db = _db(inst); if (!db) return -1; + char peers_tbl[128]; _peers_table(ch_id, peers_tbl, sizeof(peers_tbl)); + char sql[512]; + int mask_shift = 63 - (int)level * 5; + uint64_t mask = (mask_shift >= 0) ? (~0ULL << mask_shift) : UINT64_MAX; + snprintf(sql, sizeof(sql), + "SELECT p.node_id, n.x25519_pubkey, n.ed25519_pubkey, p.join_sig, n.online" + " FROM \"%s\" p JOIN nodes n ON p.node_id=n.node_id" + " WHERE (p.node_id & %lld) == %lld ORDER BY p.node_id ASC", + peers_tbl, (long long)mask, (long long)prefix64); + + sqlite3_stmt* stmt = NULL; + if (sqlite3_prepare_v2(db, sql, -1, &stmt, NULL) != SQLITE_OK) return -1; + + SHA256_CTX ctx; + SHA256_Init(&ctx); + int count = 0; + + while (sqlite3_step(stmt) == SQLITE_ROW) { + uint64_t nid = (uint64_t)sqlite3_column_int64(stmt, 0); + const uint8_t* x25 = (const uint8_t*)sqlite3_column_blob(stmt, 1); + const uint8_t* ed = (const uint8_t*)sqlite3_column_blob(stmt, 2); + const uint8_t* sig = (const uint8_t*)sqlite3_column_blob(stmt, 3); + int online = sqlite3_column_int(stmt, 4); + if (!x25 || !ed || !sig) continue; + + uint8_t mh[MS_HASH_SIZE]; + _compute_member_hash(nid, x25, ed, sig, NULL, 0, online, mh); + SHA256_Update(&ctx, mh, MS_HASH_SIZE); + count++; + } + sqlite3_finalize(stmt); + + if (count == 0) { + snprintf(sql, sizeof(sql), + "DELETE FROM member_tree_hash WHERE channel_id=? AND level=? AND prefix64=?"); + sqlite3_stmt* ds = NULL; + if (sqlite3_prepare_v2(db, sql, -1, &ds, NULL) == SQLITE_OK) { + sqlite3_bind_text(ds, 1, ch_id, -1, SQLITE_STATIC); + sqlite3_bind_int(ds, 2, level); + sqlite3_bind_int64(ds, 3, (sqlite3_int64)prefix64); + sqlite3_step(ds); sqlite3_finalize(ds); + } + return 0; + } + + uint8_t hash[MS_HASH_SIZE]; + SHA256_Final(hash, &ctx); + + snprintf(sql, sizeof(sql), + "INSERT OR REPLACE INTO member_tree_hash(channel_id, level, prefix64, hash, member_count)" + " VALUES(?,?,?,?,?)"); + sqlite3_stmt* is = NULL; + if (sqlite3_prepare_v2(db, sql, -1, &is, NULL) == SQLITE_OK) { + sqlite3_bind_text(is, 1, ch_id, -1, SQLITE_STATIC); + sqlite3_bind_int(is, 2, level); + sqlite3_bind_int64(is, 3, (sqlite3_int64)prefix64); + sqlite3_bind_blob(is, 4, hash, MS_HASH_SIZE, SQLITE_STATIC); + sqlite3_bind_int(is, 5, count); + sqlite3_step(is); sqlite3_finalize(is); + } + return 0; +} + +static void _recompute_path(struct UTUN_INSTANCE* inst, const char* ch_id, uint64_t member_id) { + for (uint8_t level = 1; level <= MS_MAX_LEVEL; level++) + _recompute_bucket(inst, ch_id, level, _level_prefix(member_id, level)); +} + +static void _ensure_table(void) { + sqlite3* db = g_ms ? _db(g_ms->inst) : NULL; if (!db) return; + sqlite3_exec(db, + "CREATE TABLE IF NOT EXISTS member_tree_hash (" + " channel_id TEXT NOT NULL," + " level INTEGER NOT NULL CHECK(level BETWEEN 1 AND 5)," + " prefix64 INTEGER NOT NULL," + " hash BLOB NOT NULL," + " member_count INTEGER NOT NULL," + " PRIMARY KEY (channel_id, level, prefix64))", + NULL, NULL, NULL); +} + +const uint8_t* member_sync_get_hash(struct UTUN_INSTANCE* inst, const char* ch_id, + uint8_t level, uint64_t prefix64) { + static uint8_t zero[MS_HASH_SIZE]; /* returns static — not thread-safe but single-threaded uasync */ + sqlite3* db = _db(inst); if (!db) { memset(zero, 0, MS_HASH_SIZE); return zero; } + sqlite3_stmt* stmt = NULL; + if (sqlite3_prepare_v2(db, + "SELECT hash FROM member_tree_hash WHERE channel_id=? AND level=? AND prefix64=?", + -1, &stmt, NULL) != SQLITE_OK) { memset(zero, 0, MS_HASH_SIZE); return zero; } + sqlite3_bind_text(stmt, 1, ch_id, -1, SQLITE_STATIC); + sqlite3_bind_int(stmt, 2, level); + sqlite3_bind_int64(stmt, 3, (sqlite3_int64)prefix64); + memset(zero, 0, MS_HASH_SIZE); + if (sqlite3_step(stmt) == SQLITE_ROW) { + const void* h = sqlite3_column_blob(stmt, 0); + if (h) memcpy(zero, h, MS_HASH_SIZE); + } + sqlite3_finalize(stmt); + return zero; +} + +int member_sync_put(struct UTUN_INSTANCE* inst, const char* ch_id, + uint64_t member_id, const uint8_t* x25519, + const uint8_t* ed25519, const uint8_t* join_sig, + const uint8_t* addrs_data, int addr_count) { + if (!inst || !ch_id) return -1; + sqlite3* db = _db(inst); if (!db) return -1; + _ensure_table(); + + sqlite3_stmt* ns = NULL; + sqlite3_prepare_v2(db, + "INSERT OR REPLACE INTO nodes(node_id, name, x25519_pubkey, ed25519_pubkey, online)" + " VALUES(?,COALESCE((SELECT name FROM nodes WHERE node_id=?),''),?,?," + " COALESCE((SELECT online FROM nodes WHERE node_id=?),0))", + -1, &ns, NULL); + if (ns) { + sqlite3_bind_int64(ns, 1, (sqlite3_int64)member_id); + sqlite3_bind_int64(ns, 2, (sqlite3_int64)member_id); + sqlite3_bind_blob(ns, 3, x25519, 32, SQLITE_STATIC); + sqlite3_bind_blob(ns, 4, ed25519, 32, SQLITE_STATIC); + sqlite3_bind_int64(ns, 5, (sqlite3_int64)member_id); + sqlite3_step(ns); sqlite3_finalize(ns); + } + + topo_node_sqlite_member_put(db, ch_id, member_id, join_sig, NULL); + + if (addrs_data && addr_count > 0) { + sqlite3_exec(db, "BEGIN", NULL, NULL, NULL); + sqlite3_stmt* ds = NULL; + sqlite3_prepare_v2(db, "DELETE FROM node_addresses WHERE node_id=? AND is_nat=0", -1, &ds, NULL); + if (ds) { sqlite3_bind_int64(ds, 1, (sqlite3_int64)member_id); sqlite3_step(ds); sqlite3_finalize(ds); } + + sqlite3_stmt* as = NULL; + sqlite3_prepare_v2(db, + "INSERT INTO node_addresses(node_id,family,protocol,address,port,is_nat)" + " VALUES(?,?,1,?,?,0)", -1, &as, NULL); + if (as) { + const uint8_t* p = addrs_data; + for (int i = 0; i < addr_count; i++) { + uint8_t fam = *p++; int ip_sz = fam == 4 ? 4 : 16; + sqlite3_bind_int64(as, 1, (sqlite3_int64)member_id); + sqlite3_bind_int(as, 2, fam); + sqlite3_bind_blob(as, 3, p, ip_sz, SQLITE_STATIC); + p += ip_sz; + uint16_t port = ((uint16_t)p[0] << 8) | p[1]; p += 2; + sqlite3_bind_int(as, 4, port); + sqlite3_step(as); sqlite3_reset(as); + } + sqlite3_finalize(as); + } + sqlite3_exec(db, "COMMIT", NULL, NULL, NULL); + } + + _recompute_path(inst, ch_id, member_id); + return 0; +} + +int member_sync_del(struct UTUN_INSTANCE* inst, const char* ch_id, uint64_t member_id) { + if (!inst || !ch_id) return -1; + sqlite3* db = _db(inst); if (!db) return -1; + topo_node_sqlite_member_del(db, ch_id, member_id); + _recompute_path(inst, ch_id, member_id); + return 0; +} + +void member_sync_set_online(struct UTUN_INSTANCE* inst, uint64_t node_id, int online) { + if (!inst) return; + sqlite3* db = _db(inst); if (!db) return; + topo_node_sqlite_node_set_online(db, node_id, online); + if (!g_ms || !g_ms->initialized) return; + sqlite3_stmt* stmt = NULL; + if (sqlite3_prepare_v2(db, "SELECT channel_id FROM channels", -1, &stmt, NULL) != SQLITE_OK) return; + while (sqlite3_step(stmt) == SQLITE_ROW) { + const char* ch_id = (const char*)sqlite3_column_text(stmt, 0); + if (!ch_id) continue; + char peers_tbl[128]; _peers_table(ch_id, peers_tbl, sizeof(peers_tbl)); + char buf[256]; snprintf(buf, sizeof(buf), + "SELECT 1 FROM \"%s\" WHERE node_id=?", peers_tbl); + sqlite3_stmt* cs = NULL; + if (sqlite3_prepare_v2(db, buf, -1, &cs, NULL) == SQLITE_OK) { + sqlite3_bind_int64(cs, 1, (sqlite3_int64)node_id); + if (sqlite3_step(cs) == SQLITE_ROW) _recompute_path(inst, ch_id, node_id); + sqlite3_finalize(cs); + } + } + sqlite3_finalize(stmt); +} + +int member_sync_count(struct UTUN_INSTANCE* inst, const char* ch_id) { + sqlite3* db = _db(inst); + if (!db || !ch_id) return 0; + char peers_tbl[128]; _peers_table(ch_id, peers_tbl, sizeof(peers_tbl)); + char sql[256]; snprintf(sql, sizeof(sql), "SELECT COUNT(*) FROM \"%s\"", peers_tbl); + sqlite3_stmt* stmt = NULL; + if (sqlite3_prepare_v2(db, sql, -1, &stmt, NULL) != SQLITE_OK) return 0; + int c = 0; + if (sqlite3_step(stmt) == SQLITE_ROW) c = sqlite3_column_int(stmt, 0); + sqlite3_finalize(stmt); + return c; +} + +int member_sync_get_level_hashes(struct UTUN_INSTANCE* inst, const char* ch_id, + uint8_t level, uint64_t prefix, uint8_t prefix_bytes, + uint32_t* bitmap, uint8_t hashes[MS_BUCKETS][MS_HASH_SIZE]) { + (void)prefix_bytes; + memset(hashes, 0, sizeof(uint8_t) * MS_BUCKETS * MS_HASH_SIZE); + *bitmap = 0; + if (level >= MS_MAX_LEVEL) return 0; + int next_shift = 63 - ((int)level + 1) * 5; + + for (int i = 0; i < MS_BUCKETS; i++) { + uint64_t child_prefix = prefix | ((uint64_t)i << (next_shift > 0 ? next_shift : 0)); + const uint8_t* h = member_sync_get_hash(inst, ch_id, (uint8_t)(level + 1), child_prefix); + int empty = 1; + for (int j = 0; j < MS_HASH_SIZE; j++) { if (h[j] != 0) { empty = 0; break; } } + if (!empty) { *bitmap |= (1u << i); memcpy(hashes[i], h, MS_HASH_SIZE); } + } + return 0; +} + +int member_sync_get_bucket_members(struct UTUN_INSTANCE* inst, const char* ch_id, + uint8_t level, uint64_t prefix, uint8_t prefix_bytes, + uint8_t* buf, size_t* len) { + (void)prefix_bytes; + sqlite3* db = _db(inst); if (!db || !buf || !len) return -1; + char peers_tbl[128]; _peers_table(ch_id, peers_tbl, sizeof(peers_tbl)); + char sql[512]; + int mask_shift = 63 - (int)level * 5; + uint64_t mask = (mask_shift >= 0) ? (~0ULL << mask_shift) : UINT64_MAX; + snprintf(sql, sizeof(sql), + "SELECT p.node_id, n.x25519_pubkey, n.ed25519_pubkey, p.join_sig, n.online" + " FROM \"%s\" p JOIN nodes n ON p.node_id=n.node_id" + " WHERE (p.node_id & %lld) == %lld ORDER BY p.node_id ASC", + peers_tbl, (long long)mask, (long long)prefix); + + sqlite3_stmt* stmt = NULL; + if (sqlite3_prepare_v2(db, sql, -1, &stmt, NULL) != SQLITE_OK) return -1; + + size_t off = 0; + if (off + 2 > *len) { sqlite3_finalize(stmt); return -2; } + uint16_t* cnt = (uint16_t*)(buf + off); off += 2; *cnt = 0; + + while (sqlite3_step(stmt) == SQLITE_ROW) { + uint64_t nid = (uint64_t)sqlite3_column_int64(stmt, 0); + const uint8_t* x25 = (const uint8_t*)sqlite3_column_blob(stmt, 1); + const uint8_t* ed = (const uint8_t*)sqlite3_column_blob(stmt, 2); + const uint8_t* sig = (const uint8_t*)sqlite3_column_blob(stmt, 3); + int online = sqlite3_column_int(stmt, 4); + if (!x25 || !ed || !sig) continue; + + sqlite3_stmt* as = NULL; + sqlite3_prepare_v2(db, + "SELECT family, address, port FROM node_addresses WHERE node_id=? AND is_nat=0" + " ORDER BY family, address, port", -1, &as, NULL); + uint8_t addrs[2048]; int addr_off = 0; int addr_count = 0; + if (as) { + sqlite3_bind_int64(as, 1, (sqlite3_int64)nid); + while (sqlite3_step(as) == SQLITE_ROW && addr_off < (int)sizeof(addrs) - 7) { + int fam = sqlite3_column_int(as, 0); + addrs[addr_off++] = (uint8_t)fam; + int ip_sz = fam == 4 ? 4 : 16; + memcpy(addrs + addr_off, sqlite3_column_blob(as, 1), (size_t)ip_sz); + addr_off += ip_sz; + uint16_t p = (uint16_t)sqlite3_column_int(as, 2); + addrs[addr_off++] = (uint8_t)(p >> 8); + addrs[addr_off++] = (uint8_t)(p & 0xFF); + addr_count++; + } + sqlite3_finalize(as); + } + + size_t need = 8 + 32 + 32 + 64 + 1 + 1 + (size_t)addr_off; + if (off + need > *len) { sqlite3_finalize(stmt); return -2; } + memcpy(buf + off, &nid, 8); off += 8; + memcpy(buf + off, x25, 32); off += 32; + memcpy(buf + off, ed, 32); off += 32; + memcpy(buf + off, sig, 64); off += 64; + buf[off++] = (uint8_t)(online ? 1 : 0); + buf[off++] = (uint8_t)addr_count; + memcpy(buf + off, addrs, (size_t)addr_off); off += (size_t)addr_off; + (*cnt)++; + } + sqlite3_finalize(stmt); + *len = off; + return 0; +} + +static int _send_msg(struct UTUN_INSTANCE* inst, uint64_t peer, + const uint8_t* payload, size_t len) { + if (!inst || len < 1) return -1; + uint8_t* buf = u_malloc(1 + len); + if (!buf) return -1; + buf[0] = ETCP_RT_ID_MEMBER_SYNC; + memcpy(buf + 1, payload, len); + struct ll_entry* entry = queue_entry_new(0); + if (!entry) { u_free(buf); return -1; } + entry->dgram = buf; entry->len = 1 + len; + int r = etcp_route_send(inst, peer, entry, 0); + if (r != 0) { u_free(buf); queue_entry_free(entry); } + return r; +} + +static int _popcount_u32(uint32_t v) { + return __builtin_popcount(v); +} +#ifndef __has_builtin +#define __has_builtin(x) 0 +#endif +#if !defined(__GNUC__) && !defined(__clang__) +static int _popcount_u32(uint32_t v) { + v = v - ((v >> 1) & 0x55555555); + v = (v & 0x33333333) + ((v >> 2) & 0x33333333); + return ((v + (v >> 4) & 0x0F0F0F0F) * 0x01010101) >> 24; +} +#endif + +static int _send_hashes(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ch_id, + uint8_t level, uint64_t prefix, uint8_t prefix_bytes, int is_data) { + uint8_t ch_len = (uint8_t)strlen(ch_id); + size_t max_sz = 1 + 1 + ch_len + 1 + 1 + prefix_bytes + 1 + 4 + MS_BUCKETS * MS_HASH_SIZE + 2 + 65536; + uint8_t* buf = u_malloc(max_sz); + if (!buf) return -1; + uint8_t* p = buf; + *p++ = MS_MSG_HASHES; + *p++ = ch_len; memcpy(p, ch_id, ch_len); p += ch_len; + *p++ = level; + *p++ = prefix_bytes; _prefix_write(p, prefix, prefix_bytes); p += prefix_bytes; + *p++ = (uint8_t)(is_data ? 1 : 0); + + if (is_data) { + size_t mlen = 65536; + uint8_t* mbuf = u_malloc(mlen); + if (mbuf) { + if (member_sync_get_bucket_members(inst, ch_id, level, prefix, prefix_bytes, mbuf, &mlen) == 0) + { memcpy(p, mbuf, mlen); p += mlen; } + else { uint16_t zero = 0; memcpy(p, &zero, 2); p += 2; } + u_free(mbuf); + } + } else { + uint32_t bitmap; uint8_t hashes[MS_BUCKETS][MS_HASH_SIZE]; + member_sync_get_level_hashes(inst, ch_id, level, prefix, prefix_bytes, &bitmap, hashes); + memcpy(p, &bitmap, 4); p += 4; + for (int i = 0; i < MS_BUCKETS; i++) + if (bitmap & (1u << i)) { memcpy(p, hashes[i], MS_HASH_SIZE); p += MS_HASH_SIZE; } + } + + int r = _send_msg(inst, peer, buf, (size_t)(p - buf)); + u_free(buf); + return r; +} + +static int _send_batch(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ch_id, + struct bucket_entry* buckets, int count) { + uint8_t ch_len = (uint8_t)strlen(ch_id); + size_t max_sz = 1 + 1 + ch_len + 1 + (size_t)count * (1 + 1 + 8 + 1 + 4 + MS_BUCKETS * MS_HASH_SIZE + 65536); + uint8_t* buf = u_malloc(max_sz); + if (!buf) return -1; + uint8_t* p = buf; + *p++ = MS_MSG_BATCH; + *p++ = ch_len; memcpy(p, ch_id, ch_len); p += ch_len; + *p++ = (uint8_t)count; + + for (int i = 0; i < count; i++) { + *p++ = buckets[i].level; + *p++ = buckets[i].prefix_bytes; + _prefix_write(p, buckets[i].prefix, buckets[i].prefix_bytes); p += buckets[i].prefix_bytes; + + int bcount = 0; + uint64_t mask = 0; + int mask_shift = 63 - (int)buckets[i].level * 5; + if (mask_shift >= 0) mask = ~0ULL << mask_shift; else mask = UINT64_MAX; + char peers_tbl[128]; _peers_table(ch_id, peers_tbl, sizeof(peers_tbl)); + sqlite3* db = _db(inst); + if (db) { + char sql[256]; snprintf(sql, sizeof(sql), + "SELECT COUNT(*) FROM \"%s\" p WHERE (p.node_id & %lld)==%lld", + peers_tbl, (long long)mask, (long long)buckets[i].prefix); + sqlite3_stmt* cs = NULL; + if (sqlite3_prepare_v2(db, sql, -1, &cs, NULL) == SQLITE_OK) { + if (sqlite3_step(cs) == SQLITE_ROW) bcount = sqlite3_column_int(cs, 0); + sqlite3_finalize(cs); + } + } + + int is_terminal = (buckets[i].level >= MS_MAX_LEVEL || bcount < 8); + *p++ = (uint8_t)(is_terminal ? 1 : 0); + + if (is_terminal) { + size_t mlen = 65536; + uint8_t* mbuf = u_malloc(mlen); + if (mbuf) { + if (member_sync_get_bucket_members(inst, ch_id, buckets[i].level, + buckets[i].prefix, buckets[i].prefix_bytes, mbuf, &mlen) == 0) + { memcpy(p, mbuf, mlen); p += mlen; } + else { uint16_t z = 0; memcpy(p, &z, 2); p += 2; } + u_free(mbuf); + } + } else { + uint32_t bm; uint8_t hs[MS_BUCKETS][MS_HASH_SIZE]; + member_sync_get_level_hashes(inst, ch_id, buckets[i].level, + buckets[i].prefix, buckets[i].prefix_bytes, &bm, hs); + memcpy(p, &bm, 4); p += 4; + for (int j = 0; j < MS_BUCKETS; j++) + if (bm & (1u << j)) { memcpy(p, hs[j], MS_HASH_SIZE); p += MS_HASH_SIZE; } + } + } + + int r = _send_msg(inst, peer, buf, (size_t)(p - buf)); + u_free(buf); + return r; +} + +/* ── sessions ── */ + +static void _session_start_timer(struct ms_session* s); + +static void _session_timeout_cb(void* arg) { + struct ms_session* s = (struct ms_session*)arg; + if (!s || !s->active || !g_ms || !g_ms->initialized) return; + s->timer = NULL; + s->retries++; + if (s->retries > 3) { + DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: sync timeout peer=%016llx ch=%s", MS_ID, + (unsigned long long)s->peer, s->ch_id); + s->active = 0; + return; + } + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: retry %d peer=%016llx ch=%s", MS_ID, + s->retries, (unsigned long long)s->peer, s->ch_id); + _send_hashes(g_ms->inst, s->peer, s->ch_id, 1, 0, 1, 0); + _session_start_timer(s); +} + +static void _session_start_timer(struct ms_session* s) { + if (!g_ms || !g_ms->inst) return; + s->timer = uasync_set_timeout(g_ms->inst->ua, (uint32_t)(MS_SYNC_TIMEOUT_MS * 10), + s, _session_timeout_cb, "ms_sync"); +} + +static struct ms_session* _session_find(uint64_t peer, const char* ch_id) { + for (struct ms_session* s = g_ms ? g_ms->sessions : NULL; s; s = s->next) + if (s->peer == peer && strcmp(s->ch_id, ch_id) == 0) return s; + return NULL; +} + +void member_sync_start(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ch_id) { + if (!inst || !ch_id || !g_ms || !g_ms->initialized) return; + _ensure_table(); + + struct ms_session* s = _session_find(peer, ch_id); + if (!s) { + s = u_calloc(1, sizeof(*s)); + if (!s) return; + snprintf(s->ch_id, sizeof(s->ch_id), "%s", ch_id); + s->peer = peer; + s->next = g_ms->sessions; + g_ms->sessions = s; + } + s->active = 1; s->retries = 0; + + _send_hashes(inst, peer, ch_id, 1, 0, 1, 0); + _session_start_timer(s); +} + +void member_sync_cancel_peer(struct UTUN_INSTANCE* inst, uint64_t peer) { + (void)inst; + struct ms_session** p = g_ms ? &g_ms->sessions : NULL; + while (p && *p) { + struct ms_session* s = *p; + if (s->peer == peer) { + if (s->timer) { uasync_cancel_timeout(g_ms->inst->ua, s->timer); s->timer = NULL; } + *p = s->next; u_free(s); + } else { + p = &(*p)->next; + } + } +} + +/* ── recv handlers ── */ + +static void _handle_hashes(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ch_id, + const uint8_t* pl, size_t plen) { + if (plen < 3) return; + uint8_t level = pl[0]; uint8_t pb = pl[1]; + uint64_t prefix = _prefix_read(pl + 2, pb); + uint8_t is_data = pl[2 + pb]; + const uint8_t* payload = pl + 3 + pb; + size_t paylen = plen - 3 - pb; + + struct ms_session* s = _session_find(peer, ch_id); + if (s && s->timer) { uasync_cancel_timeout(g_ms->inst->ua, s->timer); s->timer = NULL; } + + if (is_data) { + if (paylen < 2) return; + uint16_t count; memcpy(&count, payload, 2); + const uint8_t* mp = payload + 2; size_t mrem = paylen - 2; + for (uint16_t i = 0; i < count && mrem >= 137; i++) { + uint64_t nid; memcpy(&nid, mp, 8); mp += 8; mrem -= 8; + const uint8_t* x25 = mp; mp += 32; mrem -= 32; + const uint8_t* ed = mp; mp += 32; mrem -= 32; + const uint8_t* sig = mp; mp += 64; mrem -= 64; + uint8_t online = *mp++; mrem--; + uint8_t ac = *mp++; mrem--; + const uint8_t* addrs = mp; + int consumed = 0; + for (int a = 0; a < (int)ac && mrem >= (size_t)(1 + consumed); a++) { + uint8_t fam = mp[consumed]; consumed++; + int sz = fam == 4 ? 4 : 16; + consumed += sz + 2; + } + member_sync_put(inst, ch_id, nid, x25, ed, sig, addrs, (int)ac); + mp += consumed; mrem -= (size_t)consumed; + } + } else { + if (!s) { s = u_calloc(1, sizeof(*s)); if (!s) return; + snprintf(s->ch_id, sizeof(s->ch_id), "%s", ch_id); s->peer = peer; + s->next = g_ms->sessions; g_ms->sessions = s; } + s->active = 1; + + if (paylen < 4) return; + uint32_t remote_bm; memcpy(&remote_bm, payload, 4); + const uint8_t* rp = payload + 4; + + uint32_t local_bm; uint8_t lh[MS_BUCKETS][MS_HASH_SIZE]; + member_sync_get_level_hashes(inst, ch_id, level, prefix, pb, &local_bm, lh); + + uint32_t differs = remote_bm ^ local_bm; + for (int i = 0; i < MS_BUCKETS; i++) { + if (!(local_bm & (1u << i)) && !(remote_bm & (1u << i))) continue; + if ((local_bm & (1u << i)) && (remote_bm & (1u << i))) { + const uint8_t* rh_ptr = rp; + int rh_idx = 0; + for (int j = 0; j < i; j++) if (remote_bm & (1u << j)) rh_idx++; + if (memcmp(rh_ptr + (size_t)rh_idx * MS_HASH_SIZE, lh[i], MS_HASH_SIZE) == 0) + differs &= ~(1u << i); + } + } + + if (differs == 0) { s->active = 0; return; } + + int next_shift = 63 - ((int)level + 1) * 5; + struct bucket_entry requests[MS_BUCKETS]; int rcount = 0; + for (int i = 0; i < MS_BUCKETS && rcount < MS_MAX_BATCH; i++) { + if (!(differs & (1u << i))) continue; + uint8_t nl = (uint8_t)(level + 1); + if (nl > MS_MAX_LEVEL) nl = MS_MAX_LEVEL; + requests[rcount].level = nl; + requests[rcount].prefix_bytes = _prefix_bytes(nl); + requests[rcount].prefix = prefix | ((uint64_t)i << (next_shift > 0 ? next_shift : 0)); + rcount++; + } + + if (rcount > 0) { + uint8_t ch_len = (uint8_t)strlen(ch_id); + size_t rs = 1 + 1 + ch_len + 1; + for (int i = 0; i < rcount; i++) rs += 1 + 1 + requests[i].prefix_bytes + 1; + uint8_t* rbuf = u_malloc(rs); + if (rbuf) { + uint8_t* wr = rbuf; + *wr++ = MS_MSG_REQUEST; *wr++ = ch_len; + memcpy(wr, ch_id, ch_len); wr += ch_len; + *wr++ = (uint8_t)rcount; + for (int i = 0; i < rcount; i++) { + *wr++ = requests[i].level; + *wr++ = requests[i].prefix_bytes; + _prefix_write(wr, requests[i].prefix, requests[i].prefix_bytes); + wr += requests[i].prefix_bytes; + int cnt = 0; sqlite3* db = _db(inst); + if (db) { + char peers_tbl[128]; _peers_table(ch_id, peers_tbl, sizeof(peers_tbl)); + char sql[256]; snprintf(sql, sizeof(sql), + "SELECT COUNT(*) FROM \"%s\" WHERE (node_id & " + "(CASE WHEN %d>=0 THEN ~0<<%d ELSE ~0 END)) == %lld", + peers_tbl, 63-(int)requests[i].level*5, 63-(int)requests[i].level*5, + (long long)requests[i].prefix); + sqlite3_stmt* cst = NULL; + if (sqlite3_prepare_v2(db, sql, -1, &cst, NULL) == SQLITE_OK) { + if (sqlite3_step(cst) == SQLITE_ROW) cnt = sqlite3_column_int(cst, 0); + sqlite3_finalize(cst); + } + } + *wr++ = (uint8_t)(cnt < 8 || requests[i].level >= MS_MAX_LEVEL ? 1 : 0); + } + _send_msg(inst, peer, rbuf, (size_t)(wr - rbuf)); + u_free(rbuf); + } + } + } +} + +static void _handle_request(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ch_id, + const uint8_t* pl, size_t plen) { + if (plen < 1) return; + uint8_t count = pl[0]; const uint8_t* bp = pl + 1; size_t off = 0; + + struct bucket_entry buckets[MS_MAX_BATCH]; int bc = 0; + for (uint8_t i = 0; i < count && bc < MS_MAX_BATCH; i++) { + if (off + 2 > plen - 1) break; + uint8_t lvl = bp[off++]; + uint8_t pb = bp[off++]; + if (off + pb > plen - 1) break; + uint64_t pr = _prefix_read(bp + off, pb); off += pb; + buckets[bc].level = lvl; + buckets[bc].prefix_bytes = pb; + buckets[bc].prefix = pr; + bc++; + off++; /* skip full_data */ + } + + if (bc > 0) _send_batch(inst, peer, ch_id, buckets, bc); +} + +static void _handle_batch(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ch_id, + const uint8_t* pl, size_t plen) { + if (plen < 1) return; + uint8_t count = pl[0]; const uint8_t* bp = pl + 1; size_t rem = plen - 1; + + for (uint8_t i = 0; i < count && rem >= 3; i++) { + uint8_t lvl = bp[0]; uint8_t pb = bp[1]; rem -= 2; bp += 2; + if (rem < pb + 1) break; + uint64_t pr = _prefix_read(bp, pb); bp += pb; rem -= pb; + uint8_t is_data = *bp++; rem--; + + if (is_data && rem >= 2) { + uint16_t mc; memcpy(&mc, bp, 2); bp += 2; rem -= 2; + for (uint16_t j = 0; j < mc && rem >= 137; j++) { + uint64_t nid; memcpy(&nid, bp, 8); bp += 8; rem -= 8; + const uint8_t* x25 = bp; bp += 32; rem -= 32; + const uint8_t* ed = bp; bp += 32; rem -= 32; + const uint8_t* sig = bp; bp += 64; rem -= 64; + uint8_t online = *bp++; rem--; + uint8_t ac = *bp++; rem--; + const uint8_t* addrs = bp; + int consumed = 0; + for (int a = 0; a < (int)ac && rem >= (size_t)(1 + consumed); a++) { + uint8_t fam = bp[consumed]; consumed++; + int sz = fam == 4 ? 4 : 16; + consumed += sz + 2; + } + member_sync_put(inst, ch_id, nid, x25, ed, sig, addrs, (int)ac); + bp += consumed; rem -= (size_t)consumed; + } + } else if (!is_data && rem >= 4) { + uint8_t sub_pl[4096]; size_t sub_len = 0; + uint8_t next_lvl = (uint8_t)(lvl < MS_MAX_LEVEL ? lvl + 1 : lvl); + sub_pl[sub_len++] = next_lvl; + sub_pl[sub_len++] = pb; + _prefix_write(sub_pl + sub_len, pr, pb); sub_len += pb; + sub_pl[sub_len++] = 0; + uint32_t bm; memcpy(&bm, bp, 4); bp += 4; rem -= 4; + memcpy(sub_pl + sub_len, &bm, 4); sub_len += 4; + int nh = _popcount_u32(bm); + if (rem >= (size_t)nh * MS_HASH_SIZE) { + memcpy(sub_pl + sub_len, bp, (size_t)nh * MS_HASH_SIZE); + sub_len += (size_t)nh * MS_HASH_SIZE; + bp += (size_t)nh * MS_HASH_SIZE; rem -= (size_t)nh * MS_HASH_SIZE; + } + _handle_hashes(inst, peer, ch_id, sub_pl, sub_len); + } + } +} + +static void _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_ms || !g_ms->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; + + 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 MS_MSG_HASHES: _handle_hashes(g_ms->inst, peer, ch_id, pl, plen); break; + case MS_MSG_REQUEST: _handle_request(g_ms->inst, peer, ch_id, pl, plen); break; + case MS_MSG_BATCH: _handle_batch(g_ms->inst, peer, ch_id, pl, plen); break; + } + + u_free(entry->dgram); queue_entry_free(entry); +} + +/* ── bg_check ── */ + +static void _bg_timer_cb(void* arg) { + struct member_sync* ms = (struct member_sync*)arg; + if (!ms || !ms->initialized || !ms->inst) return; + sqlite3* db = _db(ms->inst); if (!db) return; + + sqlite3_stmt* cs = NULL; + sqlite3_prepare_v2(db, "SELECT channel_id FROM channels", -1, &cs, NULL); + if (!cs) return; + while (sqlite3_step(cs) == SQLITE_ROW) { + const char* ch_id = (const char*)sqlite3_column_text(cs, 0); + if (!ch_id) continue; + int rc = member_sync_bg_check(ms->inst, ch_id); + if (rc < 0) continue; + break; + } + sqlite3_finalize(cs); + + ms->bg_timer = uasync_set_timeout(ms->inst->ua, (uint32_t)(MS_BG_INTERVAL_MS * 10), + ms, _bg_timer_cb, "ms_bg"); +} + +int member_sync_bg_check(struct UTUN_INSTANCE* inst, const char* ch_id) { + sqlite3* db = _db(inst); if (!db || !ch_id) return -1; + + sqlite3_stmt* stmt = NULL; + if (sqlite3_prepare_v2(db, + "SELECT prefix64 FROM member_tree_hash WHERE channel_id=? AND level=? LIMIT 1", + -1, &stmt, NULL) != SQLITE_OK) return -1; + sqlite3_bind_text(stmt, 1, ch_id, -1, SQLITE_STATIC); + sqlite3_bind_int(stmt, 2, MS_MAX_LEVEL); + + if (sqlite3_step(stmt) != SQLITE_ROW) { sqlite3_finalize(stmt); return -1; } + uint64_t pref = (uint64_t)sqlite3_column_int64(stmt, 0); + sqlite3_finalize(stmt); + + uint8_t old_hash[MS_HASH_SIZE]; + memcpy(old_hash, member_sync_get_hash(inst, ch_id, MS_MAX_LEVEL, pref), MS_HASH_SIZE); + + _recompute_bucket(inst, ch_id, MS_MAX_LEVEL, pref); + + const uint8_t* new_hash = member_sync_get_hash(inst, ch_id, MS_MAX_LEVEL, pref); + if (memcmp(old_hash, new_hash, MS_HASH_SIZE) != 0) { + for (uint8_t lv = (uint8_t)(MS_MAX_LEVEL - 1); lv >= 1; lv--) + _recompute_bucket(inst, ch_id, lv, _level_prefix(pref, lv)); + return 1; + } + return 0; +} + +/* called when topo_group stores updated node info (pubkeys/addresses) to DB */ +static void _on_node_updated(struct UTUN_INSTANCE* inst, uint64_t node_id, + const uint8_t* x25519, const uint8_t* ed25519) { + (void)x25519; (void)ed25519; + if (!inst || !g_ms || !g_ms->initialized) return; + sqlite3* db = _db(inst); if (!db) return; + sqlite3_stmt* cs = NULL; + if (sqlite3_prepare_v2(db, "SELECT channel_id FROM channels", -1, &cs, NULL) != SQLITE_OK) return; + while (sqlite3_step(cs) == SQLITE_ROW) { + const char* ch = (const char*)sqlite3_column_text(cs, 0); + if (!ch) continue; + char peers_tbl[128]; _peers_table(ch, peers_tbl, sizeof(peers_tbl)); + char buf[256]; snprintf(buf, sizeof(buf), "SELECT 1 FROM \"%s\" WHERE node_id=?", peers_tbl); + sqlite3_stmt* ps = NULL; + if (sqlite3_prepare_v2(db, buf, -1, &ps, NULL) == SQLITE_OK) { + sqlite3_bind_int64(ps, 1, (sqlite3_int64)node_id); + if (sqlite3_step(ps) == SQLITE_ROW) _recompute_path(inst, ch, node_id); + sqlite3_finalize(ps); + } + } + sqlite3_finalize(cs); +} + +int member_sync_init(struct UTUN_INSTANCE* inst) { + if (!inst) return -1; + struct member_sync* ms = u_calloc(1, sizeof(*ms)); + if (!ms) return -1; + ms->inst = inst; + ms->initialized = 1; + g_ms = ms; + + _ensure_table(); + + etcp_router_bind(inst, ETCP_RT_ID_MEMBER_SYNC, _recv_cb); + + topo_groups_set_node_updated_cb(inst->topo_groups, _on_node_updated); + + ms->bg_timer = uasync_set_timeout(inst->ua, (uint32_t)(MS_BG_INTERVAL_MS * 10), + ms, _bg_timer_cb, "ms_bg"); + + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: initialized", MS_ID); + return 0; +} + +void member_sync_destroy(struct UTUN_INSTANCE* inst) { + if (!g_ms || !inst) return; + g_ms->initialized = 0; + + etcp_router_bind(inst, ETCP_RT_ID_MEMBER_SYNC, NULL); + + if (g_ms->bg_timer) { uasync_cancel_timeout(inst->ua, g_ms->bg_timer); g_ms->bg_timer = NULL; } + + struct ms_session* s = g_ms->sessions; + while (s) { struct ms_session* next = s->next; + if (s->timer) uasync_cancel_timeout(inst->ua, s->timer); + u_free(s); s = next; } + + u_free(g_ms); g_ms = NULL; +} diff --git a/tools/chatgui/transport/member_sync.h b/tools/chatgui/transport/member_sync.h new file mode 100644 index 00000000..6d489eaf --- /dev/null +++ b/tools/chatgui/transport/member_sync.h @@ -0,0 +1,63 @@ +#ifndef MEMBER_SYNC_H +#define MEMBER_SYNC_H + +#ifdef __cplusplus +extern "C" { +#endif + +#include +#include + +struct UTUN_INSTANCE; + +#define ETCP_RT_ID_MEMBER_SYNC 0x31 +#define MS_MAX_BATCH 32 +#define MS_MAX_LEVEL 5 +#define MS_BUCKETS 32 +#define MS_HASH_SIZE 32 + +/* Wire message types (etcp_router service 0x31) */ +#define MS_MSG_HASHES 0x01 +#define MS_MSG_REQUEST 0x02 +#define MS_MSG_BATCH 0x03 + +/* === Initialization === */ +int member_sync_init(struct UTUN_INSTANCE* inst); +void member_sync_destroy(struct UTUN_INSTANCE* inst); + +/* === CRUD (DB write + tree recompute) === */ +int member_sync_put(struct UTUN_INSTANCE* inst, const char* ch_id, + uint64_t member_id, const uint8_t* x25519, + const uint8_t* ed25519, const uint8_t* join_sig, + const uint8_t* addrs_data, int addr_count); +int member_sync_del(struct UTUN_INSTANCE* inst, const char* ch_id, uint64_t member_id); +int member_sync_count(struct UTUN_INSTANCE* inst, const char* ch_id); + +/* === Online status (updates nodes.online + recomputes tree path) === */ +void member_sync_set_online(struct UTUN_INSTANCE* inst, uint64_t node_id, int online); + +/* === Protocol primitives === */ +int member_sync_get_level_hashes(struct UTUN_INSTANCE* inst, + const char* ch_id, uint8_t level, + uint64_t prefix, uint8_t prefix_bytes, + uint32_t* bitmap, uint8_t hashes[MS_BUCKETS][MS_HASH_SIZE]); +int member_sync_get_bucket_members(struct UTUN_INSTANCE* inst, + const char* ch_id, uint8_t level, + uint64_t prefix, uint8_t prefix_bytes, + uint8_t* buf, size_t* len); + +/* === Background consistency check === */ +int member_sync_bg_check(struct UTUN_INSTANCE* inst, const char* ch_id); + +/* === Sync sessions === */ +void member_sync_start(struct UTUN_INSTANCE* inst, uint64_t peer, const char* ch_id); +void member_sync_cancel_peer(struct UTUN_INSTANCE* inst, uint64_t peer); + +/* === Test helper === */ +const uint8_t* member_sync_get_hash(struct UTUN_INSTANCE* inst, + const char* ch_id, uint8_t level, uint64_t prefix64); + +#ifdef __cplusplus +} +#endif +#endif /* MEMBER_SYNC_H */ diff --git a/tools/chatgui/transport/topo_node_sqlite.c b/tools/chatgui/transport/topo_node_sqlite.c index 771a8bb4..b119c29b 100644 --- a/tools/chatgui/transport/topo_node_sqlite.c +++ b/tools/chatgui/transport/topo_node_sqlite.c @@ -34,6 +34,7 @@ int topo_node_sqlite_init(sqlite3* db) { " x25519_pubkey BLOB NOT NULL," " ed25519_pubkey BLOB," " last_seen_at INTEGER," + " online INTEGER DEFAULT 0," " created_at INTEGER DEFAULT (unixepoch())" ");" @@ -72,6 +73,7 @@ int topo_node_sqlite_init(sqlite3* db) { sqlite3_free(err); return -1; } + DEBUG_INFO(DEBUG_CATEGORY_BGP, "topo_node_sqlite tables initialized"); return 0; } @@ -295,3 +297,25 @@ int topo_node_sqlite_channel_peers_all(sqlite3* db, const char* channel_id, *out_len = off; return 0; } + +int topo_node_sqlite_node_set_online(sqlite3* db, uint64_t node_id, int online) { + if (!db) return -1; + sqlite3_stmt* stmt = NULL; + if (sqlite3_prepare_v2(db, "UPDATE nodes SET online=? WHERE node_id=?", -1, &stmt, NULL) != SQLITE_OK) return -1; + sqlite3_bind_int(stmt, 1, online ? 1 : 0); + sqlite3_bind_int64(stmt, 2, (sqlite3_int64)node_id); + sqlite3_step(stmt); + sqlite3_finalize(stmt); + return 0; +} + +int topo_node_sqlite_node_get_online(sqlite3* db, uint64_t node_id) { + if (!db) return 0; + sqlite3_stmt* stmt = NULL; + if (sqlite3_prepare_v2(db, "SELECT online FROM nodes WHERE node_id=?", -1, &stmt, NULL) != SQLITE_OK) return 0; + sqlite3_bind_int64(stmt, 1, (sqlite3_int64)node_id); + int online = 0; + if (sqlite3_step(stmt) == SQLITE_ROW) online = sqlite3_column_int(stmt, 0); + sqlite3_finalize(stmt); + return online; +} diff --git a/tools/chatgui/transport/topo_node_sqlite.h b/tools/chatgui/transport/topo_node_sqlite.h index 8691995b..ac6b886e 100644 --- a/tools/chatgui/transport/topo_node_sqlite.h +++ b/tools/chatgui/transport/topo_node_sqlite.h @@ -31,4 +31,7 @@ int topo_node_sqlite_channel_get(sqlite3* db, const char* channel_id, int topo_node_sqlite_channel_peers_all(sqlite3* db, const char* channel_id, uint8_t* buf, size_t buf_sz, size_t* out_len); +int topo_node_sqlite_node_set_online(sqlite3* db, uint64_t node_id, int online); +int topo_node_sqlite_node_get_online(sqlite3* db, uint64_t node_id); + #endif