diff --git a/src/Makefile.am b/src/Makefile.am index 862424ab..e4a09e34 100644 --- a/src/Makefile.am +++ b/src/Makefile.am @@ -45,6 +45,7 @@ utun_CORE_SOURCES = \ nat_transport.c \ dummynet.c \ ntp_time.c \ + ntp_node_time.c \ proxy/tcp_proxy_client.c \ etcp_router.c \ proxy/tcp_proxy_server.c \ @@ -99,6 +100,7 @@ libutun_a_SOURCES = \ nat_transport.c \ dummynet.c \ ntp_time.c \ + ntp_node_time.c \ proxy/tcp_proxy_client.c \ etcp_router.c \ proxy/tcp_proxy_server.c \ diff --git a/src/etcp_api.h b/src/etcp_api.h index 8aae831f..16b65fdf 100644 --- a/src/etcp_api.h +++ b/src/etcp_api.h @@ -41,6 +41,7 @@ extern "C" { #define ETCP_RT_ID_ICMP_PROXY 0x06 // ICMP echo прокси (ping через exit) #define ETCP_RT_ID_MSG_TRANSPORT 0x10 // msg_transport — локальный IPC транспорт сообщений #define ETCP_RT_ID_CONN_MGR 0x11 // Connection Manager — management connections +#define ETCP_RT_ID_NTP_TIME 0x12 // NTP time sync между узлами // Forward declarations struct ETCP_CONN; diff --git a/src/ntp_node_time.c b/src/ntp_node_time.c new file mode 100644 index 00000000..d1945cdc --- /dev/null +++ b/src/ntp_node_time.c @@ -0,0 +1,184 @@ +#include "ntp_node_time.h" +#include "ntp_time.h" +#include "utun_instance.h" +#include "../lib/debug_config.h" +#include "../lib/mem.h" +#include "etcp_api.h" +#include "etcp_router.h" +#include "etcp.h" +#include +#include +#include + +#define TIME_SYNC_FLAG_SYNCED 0x01 + +#pragma pack(push, 1) +struct time_sync_msg { + uint8_t flags; // bit0 = sender_synced + int64_t t_sender_us; // corrected time if synced, local gettimeofday otherwise +}; +#pragma pack(pop) + +_Static_assert(sizeof(struct time_sync_msg) == 9, "time_sync_msg size mismatch"); + +static void ntp_node_on_conn_ready(struct ETCP_CONN* conn, void* arg); +static void ntp_node_on_new_conn(struct ETCP_CONN* conn, void* arg); +static void ntp_node_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry); + +static int send_time_sync(struct UTUN_INSTANCE* inst, uint64_t dst_node_id) { + struct ETCP_ROUTER_CONN* rconn = etcp_router_conn_get(inst, dst_node_id, ETCP_RT_ID_NTP_TIME); + if (!rconn) return -1; + + struct time_sync_msg msg; + msg.flags = ntp_time_is_synced(inst) ? TIME_SYNC_FLAG_SYNCED : 0; + + struct timeval tv; +#ifdef _WIN32 + utun_gettimeofday(&tv, NULL); +#else + gettimeofday(&tv, NULL); +#endif + msg.t_sender_us = ntp_time_is_synced(inst) + ? ntp_time_get_us(inst) + : (int64_t)tv.tv_sec * 1000000LL + tv.tv_usec; + + return etcp_router_conn_send(rconn, (const uint8_t*)&msg, sizeof(msg)); +} + +static int64_t gettimeofday_us(void) { + struct timeval tv; +#ifdef _WIN32 + utun_gettimeofday(&tv, NULL); +#else + gettimeofday(&tv, NULL); +#endif + return (int64_t)tv.tv_sec * 1000000LL + tv.tv_usec; +} + +static struct NTP_NODE_PEER* find_peer(struct NTP_NODE_TIME* np, uint64_t node_id) { + for (int i = 0; i < np->peer_count; i++) { + if (np->peers[i].node_id == node_id) return &np->peers[i]; + } + return NULL; +} + +static struct NTP_NODE_PEER* add_peer(struct NTP_NODE_TIME* np, uint64_t node_id) { + if (np->peer_count >= NTP_NODE_MAX_PEERS) return NULL; + struct NTP_NODE_PEER* p = &np->peers[np->peer_count]; + p->node_id = node_id; + p->offset_us = 0; + np->peer_count++; + return p; +} + +static void check_drift(uint64_t node_id, int64_t offset_us) { + int64_t abs_us = offset_us < 0 ? -offset_us : offset_us; + if (abs_us > NTP_NODE_DRIFT_ERROR_US) { + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NTP_node: clock drift >10s with node %012llx: %lldus", + (unsigned long long)node_id, (long long)offset_us); + } else if (abs_us > NTP_NODE_DRIFT_WARN_US) { + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "NTP_node: clock drift >2s with node %012llx: %lldus", + (unsigned long long)node_id, (long long)offset_us); + } +} + +static void ntp_node_on_conn_ready(struct ETCP_CONN* conn, void* arg) { + struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)arg; + if (!inst || !conn) return; + + if (send_time_sync(inst, conn->peer_node_id) == 0) { + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "NTP_node: sent TIME_SYNC to node %012llx (synced=%d)", + (unsigned long long)conn->peer_node_id, ntp_time_is_synced(inst)); + } else { + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "NTP_node: failed to send TIME_SYNC to node %012llx", + (unsigned long long)conn->peer_node_id); + } +} + +static void ntp_node_on_new_conn(struct ETCP_CONN* conn, void* arg) { + if (!conn) return; + struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)arg; + if (!inst) return; + etcp_conn_add_ready_cbk(conn, ntp_node_on_conn_ready, inst); +} + +static void ntp_node_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) { + if (!conn || !entry) return; + struct UTUN_INSTANCE* inst = conn->instance; + if (!inst) { queue_dgram_free(entry); return; } + + uint8_t svc_id = entry->dgram[offsetof(struct SVC_ROUTE_HDR, svc_id)]; + if (svc_id != ETCP_RT_ID_NTP_TIME) { queue_dgram_free(entry); return; } + + size_t payload_len = entry->len < SVC_ROUTE_HDR_SIZE ? 0 : entry->len - SVC_ROUTE_HDR_SIZE; + if (payload_len < sizeof(struct time_sync_msg)) { queue_dgram_free(entry); return; } + + struct time_sync_msg* msg = (struct time_sync_msg*)(entry->dgram + SVC_ROUTE_HDR_SIZE); + int sender_synced = (msg->flags & TIME_SYNC_FLAG_SYNCED) != 0; + + int64_t t4 = gettimeofday_us(); + int64_t peer_offset = t4 - msg->t_sender_us; + + struct NTP_NODE_TIME* np = &inst->ntp_node; + struct NTP_NODE_PEER* peer = find_peer(np, conn->peer_node_id); + if (!peer) peer = add_peer(np, conn->peer_node_id); + if (peer) peer->offset_us = peer_offset; + + check_drift(conn->peer_node_id, peer_offset); + + int was_unsynced = !ntp_time_is_synced(inst); + + if (was_unsynced && sender_synced) { + inst->ntp.offset_us = msg->t_sender_us - t4; + inst->ntp.synced = 1; + inst->ntp.last_sync_tb = get_time_tb(); + + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "NTP_node: synced from node %012llx, offset=%lldus", + (unsigned long long)conn->peer_node_id, (long long)inst->ntp.offset_us); + } + + queue_dgram_free(entry); + + // Ответный TIME_SYNC — всегда + send_time_sync(inst, conn->peer_node_id); + + // Если мы только что скорректировались → раздаём коррекцию всем peer'ам + if (was_unsynced && ntp_time_is_synced(inst)) { + ntp_node_sync_peers(inst); + } +} + +void ntp_node_sync_peers(struct UTUN_INSTANCE* inst) { + if (!inst || !ntp_time_is_synced(inst)) return; + + int sent = 0; + for (struct ETCP_CONN* conn = inst->connections; conn; conn = conn->next) { + if (!conn->initialized || !conn->links_up) continue; + if (send_time_sync(inst, conn->peer_node_id) == 0) sent++; + } + + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "NTP_node: broadcast TIME_SYNC to %d peers", sent); +} + +int ntp_node_time_init(struct UTUN_INSTANCE* inst) { + if (!inst) return -1; + struct NTP_NODE_TIME* np = &inst->ntp_node; + np->peer_count = 0; + + etcp_router_bind(inst, ETCP_RT_ID_NTP_TIME, ntp_node_recv_cb); + etcp_add_new_conn_cbk(inst, ntp_node_on_new_conn, inst); + + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "NTP_node: initialized (RT_ID=0x%02X)", ETCP_RT_ID_NTP_TIME); + return 0; +} + +void ntp_node_time_destroy(struct UTUN_INSTANCE* inst) { + if (!inst) return; + struct NTP_NODE_TIME* np = &inst->ntp_node; + + etcp_router_unbind(inst, ETCP_RT_ID_NTP_TIME); + etcp_remove_new_conn_cbk(inst, ntp_node_on_new_conn, inst); + + np->peer_count = 0; + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "NTP_node: destroyed"); +} diff --git a/src/ntp_node_time.h b/src/ntp_node_time.h new file mode 100644 index 00000000..a760ceff --- /dev/null +++ b/src/ntp_node_time.h @@ -0,0 +1,35 @@ +#ifndef NTP_NODE_TIME_H +#define NTP_NODE_TIME_H + +#include + +#ifdef __cplusplus +extern "C" { +#endif + +struct UTUN_INSTANCE; + +#define NTP_NODE_MAX_PEERS 32 + +#define NTP_NODE_DRIFT_WARN_US 2000000LL // 2 seconds +#define NTP_NODE_DRIFT_ERROR_US 10000000LL // 10 seconds + +struct NTP_NODE_PEER { + uint64_t node_id; + int64_t offset_us; // my_time_recv - peer_time (positive = my clock ahead) +}; + +struct NTP_NODE_TIME { + struct NTP_NODE_PEER peers[NTP_NODE_MAX_PEERS]; + int peer_count; +}; + +int ntp_node_time_init(struct UTUN_INSTANCE* inst); +void ntp_node_time_destroy(struct UTUN_INSTANCE* inst); +void ntp_node_sync_peers(struct UTUN_INSTANCE* inst); + +#ifdef __cplusplus +} +#endif + +#endif // NTP_NODE_TIME_H diff --git a/src/ntp_time.c b/src/ntp_time.c index 6c76add0..6ca755af 100644 --- a/src/ntp_time.c +++ b/src/ntp_time.c @@ -1,4 +1,5 @@ #include "ntp_time.h" +#include "ntp_node_time.h" #include "utun_instance.h" #include "../lib/debug_config.h" #include "../lib/socket_compat.h" @@ -212,6 +213,7 @@ static void ntp_time_sync_cb(void* arg) { if (ntp->synced) { DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "NTP: time corrected, offset=%lldus", (long long)ntp->offset_us); + ntp_node_sync_peers(instance); } else { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "NTP: all servers unreachable, retry in %ds", ntp->resync_interval_sec); diff --git a/src/utun_instance.c b/src/utun_instance.c index c75842ee..4d2741fb 100644 --- a/src/utun_instance.c +++ b/src/utun_instance.c @@ -371,6 +371,7 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) { // Cancel NTP timer ntp_time_destroy(instance); + ntp_node_time_destroy(instance); // Shutdown message transport server if (instance->msg_t) { @@ -617,6 +618,11 @@ int utun_instance_init(struct UTUN_INSTANCE *instance) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "NTP time sync init failed (non-fatal)"); } + // Initialize NTP node time exchange (non-fatal) + if (ntp_node_time_init(instance) != 0) { + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "NTP node time exchange init failed (non-fatal)"); + } + DEBUG_DEBUG(DEBUG_CATEGORY_ETCP, "Connections initialized successfully"); return 0; } diff --git a/src/utun_instance.h b/src/utun_instance.h index fb924d5b..9251c559 100644 --- a/src/utun_instance.h +++ b/src/utun_instance.h @@ -21,6 +21,7 @@ extern "C" { #include "etcp_router.h" #include "proxy/tcp_proxy_server.h" #include "ntp_time.h" +#include "ntp_node_time.h" #include "etcp_api.h" #include "config_parser.h" @@ -149,6 +150,7 @@ struct UTUN_INSTANCE { // NTP time synchronization struct NTP_TIME ntp; + struct NTP_NODE_TIME ntp_node; }; // Functions diff --git a/tools/chatgui/libutun/CMakeLists.txt b/tools/chatgui/libutun/CMakeLists.txt index 382b8588..d40e834c 100644 --- a/tools/chatgui/libutun/CMakeLists.txt +++ b/tools/chatgui/libutun/CMakeLists.txt @@ -93,6 +93,7 @@ set(UTUN_COMMON_SOURCES ${SRC_DIR}/nat_transport.c ${SRC_DIR}/dummynet.c ${SRC_DIR}/ntp_time.c + ${SRC_DIR}/ntp_node_time.c ${SRC_DIR}/etcp_router.c ${SRC_DIR}/proxy/tcp_proxy_client.c ${SRC_DIR}/proxy/tcp_proxy_server.c diff --git a/tools/chatgui/src/channeldelegate.cpp b/tools/chatgui/src/channeldelegate.cpp index 33cfe93a..b8a43820 100644 --- a/tools/chatgui/src/channeldelegate.cpp +++ b/tools/chatgui/src/channeldelegate.cpp @@ -9,15 +9,30 @@ void ChannelDelegate::paint(QPainter *painter, const QStyleOptionViewItem &optio const QModelIndex &index) const { painter->save(); + const int barWidth = 4; + const int margin = 8; + const int iconSize = 32; + + int peersOnline = index.data(ChannelPeersOnlineRole).toInt(); + int autoActive = index.data(ChannelAutoConnectActiveRole).toInt(); + + /* ── background ── */ if (option.state & QStyle::State_Selected) { painter->fillRect(option.rect, option.palette.highlight()); } else if (option.state & QStyle::State_MouseOver) { painter->fillRect(option.rect, option.palette.light()); } - const int margin = 8; - const int iconSize = 32; - QRect iconRect(option.rect.left() + margin, + /* ── vertical status bar (left edge, over background) ── */ + QColor barColor = Qt::transparent; + if (peersOnline > 0) barColor = QColor("#4CAF50"); + else if (autoActive && peersOnline == 0) barColor = QColor("#FFC107"); + if (barColor.alpha() > 0) { + QRect barRect(option.rect.left(), option.rect.top(), barWidth, option.rect.height()); + painter->fillRect(barRect, barColor); + } + + QRect iconRect(option.rect.left() + margin + barWidth, option.rect.top() + (option.rect.height() - iconSize) / 2, iconSize, iconSize); @@ -26,6 +41,24 @@ void ChannelDelegate::paint(QPainter *painter, const QStyleOptionViewItem &optio icon.paint(painter, iconRect); } + /* ── online count badge (top-left of avatar, above it) ── */ + if (peersOnline > 0) { + int badgeSize = 16; + QRect badgeRect(iconRect.left() - 4, iconRect.top() - 6, badgeSize, badgeSize); + painter->setRenderHint(QPainter::Antialiasing); + painter->setBrush(QColor("#4CAF50")); + painter->setPen(QPen(Qt::white, 1.5)); + painter->drawEllipse(badgeRect); + QFont badgeFont = option.font; + badgeFont.setBold(true); + badgeFont.setPixelSize(10); + painter->setFont(badgeFont); + painter->setPen(Qt::white); + QString num = QString::number(peersOnline); + if (peersOnline > 9) num = "9+"; + painter->drawText(badgeRect, Qt::AlignCenter, num); + } + QRect textRect(iconRect.right() + margin, option.rect.top() + margin, option.rect.right() - iconRect.right() - 2 * margin, option.rect.height() - 2 * margin); diff --git a/tools/chatgui/src/channeldelegate.h b/tools/chatgui/src/channeldelegate.h index 436545a3..30f59090 100644 --- a/tools/chatgui/src/channeldelegate.h +++ b/tools/chatgui/src/channeldelegate.h @@ -3,10 +3,12 @@ #include enum ChannelDataRole { - ChannelNameRole = Qt::DisplayRole, - ChannelIconRole = Qt::DecorationRole, - ChannelLastMessageRole = Qt::UserRole + 1, - ChannelChannelIdRole = Qt::UserRole + 2 + ChannelNameRole = Qt::DisplayRole, + ChannelIconRole = Qt::DecorationRole, + ChannelLastMessageRole = Qt::UserRole + 1, + ChannelChannelIdRole = Qt::UserRole + 2, + ChannelPeersOnlineRole = Qt::UserRole + 3, /* int: count of online peers */ + ChannelAutoConnectActiveRole = Qt::UserRole + 4 /* bool: 1 if auto-connect in progress */ }; class ChannelDelegate : public QStyledItemDelegate { diff --git a/tools/chatgui/src/channellist.cpp b/tools/chatgui/src/channellist.cpp index b2f2b933..880f891b 100644 --- a/tools/chatgui/src/channellist.cpp +++ b/tools/chatgui/src/channellist.cpp @@ -8,6 +8,28 @@ #include #include #include +#include + +/* QListView that prevents deselection */ +class ChannelListView : public QListView { +public: + using QListView::QListView; +protected: + void mousePressEvent(QMouseEvent* e) override { + if (indexAt(e->pos()).isValid()) + QListView::mousePressEvent(e); + } + QItemSelectionModel::SelectionFlags selectionCommand(const QModelIndex& idx, const QEvent* e) const override { + if (idx.isValid()) + return QItemSelectionModel::ClearAndSelect | QItemSelectionModel::Current; + return QItemSelectionModel::NoUpdate; + } + void selectionChanged(const QItemSelection& sel, const QItemSelection& desel) override { + QListView::selectionChanged(sel, desel); + if (selectionModel() && !selectionModel()->hasSelection() && currentIndex().isValid()) + selectionModel()->select(currentIndex(), QItemSelectionModel::ClearAndSelect); + } +}; static QIcon makeChannelIcon(const QColor &color, const QString &letter) { QPixmap pix(32, 32); @@ -29,7 +51,7 @@ static QIcon makeChannelIcon(const QColor &color, const QString &letter) { ChannelList::ChannelList(DbManager* db, QWidget *parent) : QWidget(parent) - , m_listView(new QListView(this)) + , m_listView(new ChannelListView(this)) , m_model(new QStandardItemModel(this)) , m_db(db) { @@ -161,3 +183,22 @@ void ChannelList::onContextMenu(const QPoint& pos) { if (menu.exec(m_listView->viewport()->mapToGlobal(pos)) == inviteAction) emit inviteRequested(cid); } + +void ChannelList::setChannelPeersOnline(const QString& channelId, int count) { + for (int row = 0; row < m_model->rowCount(); row++) { + QStandardItem* item = m_model->item(row); + if (item && item->data(ChannelChannelIdRole).toString() == channelId) { + item->setData(count, ChannelPeersOnlineRole); + item->setData(m_autoConnectActive ? 1 : 0, ChannelAutoConnectActiveRole); + break; + } + } +} + +void ChannelList::setAutoConnectActive(bool active) { + m_autoConnectActive = active; + for (int row = 0; row < m_model->rowCount(); row++) { + QStandardItem* item = m_model->item(row); + if (item) item->setData(active ? 1 : 0, ChannelAutoConnectActiveRole); + } +} diff --git a/tools/chatgui/src/channellist.h b/tools/chatgui/src/channellist.h index 176025c2..2a9cc6ae 100644 --- a/tools/chatgui/src/channellist.h +++ b/tools/chatgui/src/channellist.h @@ -14,6 +14,8 @@ public: void loadChannels(); void selectChannel(const QString& channelId); + void setChannelPeersOnline(const QString& channelId, int count); + void setAutoConnectActive(bool active); signals: void channelSelected(const QString& channelId); @@ -29,4 +31,5 @@ private: QListView *m_listView; QStandardItemModel *m_model; DbManager *m_db; + bool m_autoConnectActive = false; }; diff --git a/tools/chatgui/src/mainwindow.cpp b/tools/chatgui/src/mainwindow.cpp index 5f190148..adce67e0 100644 --- a/tools/chatgui/src/mainwindow.cpp +++ b/tools/chatgui/src/mainwindow.cpp @@ -46,6 +46,17 @@ static void onMembersChangedCallback(const char* ch_id, int ch_id_len) { s_mainWindow->onMembersChanged(ch_id); } +static void onAutoConnectStatusCallback(uint8_t status, uint16_t total_tried, + uint16_t node_count, uint16_t connected) { + if (s_mainWindow) + s_mainWindow->onAutoConnectStatus(status, total_tried, node_count, connected); +} + +static void onChannelPeersOnlineCallback(const char* ch_id, int ch_id_len, uint16_t online) { + if (s_mainWindow) + s_mainWindow->onChannelPeersOnline(ch_id, ch_id_len, online); +} + MainWindow::MainWindow(QWidget *parent, DbManager* db, const QString& cfgPath, const QString& dbPath, const QString& debugFile, const QString& debugLevel, const QString& debugCategories) @@ -76,6 +87,8 @@ MainWindow::MainWindow(QWidget *parent, DbManager* db, const QString& cfgPath, 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); }); + gui_bridge_set_auto_connect_status_cb(onAutoConnectStatusCallback); + gui_bridge_set_channel_peers_online_cb(onChannelPeersOnlineCallback); setupNode(); setupMessaging(); @@ -258,6 +271,20 @@ void MainWindow::onMembersChanged(const char* ch_id) { m_accountList->refresh(); } +void MainWindow::onAutoConnectStatus(uint8_t status, uint16_t /*total_tried*/, + uint16_t /*node_count*/, uint16_t /*connected*/) { + if (m_channelList) { + m_channelList->setAutoConnectActive(status == 0); + } +} + +void MainWindow::onChannelPeersOnline(const char* ch_id, int ch_id_len, uint16_t online) { + if (m_channelList) { + QString cid = QString::fromUtf8(ch_id, ch_id_len); + m_channelList->setChannelPeersOnline(cid, (int)online); + } +} + void MainWindow::reloadChannels() { if (m_channelList) { m_channelList->loadChannels(); diff --git a/tools/chatgui/src/mainwindow.h b/tools/chatgui/src/mainwindow.h index b730d8c6..c5a70ed9 100644 --- a/tools/chatgui/src/mainwindow.h +++ b/tools/chatgui/src/mainwindow.h @@ -33,6 +33,8 @@ private slots: public: void onMessageReceived(const char* ch_id); void onMembersChanged(const char* ch_id); + void onAutoConnectStatus(uint8_t status, uint16_t total, uint16_t totalNodes, uint16_t connected); + void onChannelPeersOnline(const char* ch_id, int ch_id_len, uint16_t online); void reloadChannels(); private: diff --git a/tools/chatgui/transport/chat_core.c b/tools/chatgui/transport/chat_core.c index 3b7bf8e8..86407860 100644 --- a/tools/chatgui/transport/chat_core.c +++ b/tools/chatgui/transport/chat_core.c @@ -876,6 +876,185 @@ void chat_core_connect_from_invite(struct chat_invite* inv) { } } +/* ═══ Auto-connect: direct ETCP connection from SQLite (no BGP) ═══ */ + +#define CA_CONNECT_TIMEOUT_MS 3000 + +struct ca_state; + +struct ca_ctx { + struct ca_state* state; + int addr_index; +}; + +struct ca_state { + struct UTUN_INSTANCE* inst; + int addr_count; + int pending_count; + int completed; /* 0=pending, 1=delivered */ + uint64_t node_id; + struct ETCP_CONN** conns; + void** timers; + void (*result_cb)(int result, uint64_t node_id, void* arg); + void* result_arg; + uint8_t pubkey[SC_PUBKEY_SIZE]; +}; + +static void ca_cleanup(struct ca_state* st) { + if (!st) return; + u_free(st->conns); + u_free(st->timers); + u_free(st); +} + +static void ca_ready_cb(struct ETCP_CONN* conn, void* arg) { + struct ca_ctx* ctx = (struct ca_ctx*)arg; + struct ca_state* st = ctx->state; + if (st->completed) { u_free(ctx); return; } + st->completed = 1; + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect SUCCESS idx=%d peer=0x%016llx", + CC_ID, ctx->addr_index, (unsigned long long)st->node_id); + for (int i = 0; i < st->addr_count; i++) { + if (st->timers[i]) { uasync_cancel_timeout(st->inst->ua, st->timers[i]); st->timers[i] = NULL; } + if (st->conns[i] && i != ctx->addr_index) { etcp_connection_close(st->conns[i]); st->conns[i] = NULL; } + } + u_free(ctx); + st->result_cb(CONN_MGR_OK, st->node_id, st->result_arg); + ca_cleanup(st); +} + +static void ca_timeout_cb(void* arg) { + struct ca_ctx* ctx = (struct ca_ctx*)arg; + struct ca_state* st = ctx->state; + if (st->completed) { u_free(ctx); return; } + if (st->conns[ctx->addr_index]) { etcp_connection_close(st->conns[ctx->addr_index]); st->conns[ctx->addr_index] = NULL; } + st->timers[ctx->addr_index] = NULL; + st->pending_count--; + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect TIMEOUT idx=%d pending=%d/%d peer=0x%016llx", + CC_ID, ctx->addr_index, st->pending_count, st->addr_count, (unsigned long long)st->node_id); + if (st->pending_count <= 0 && !st->completed) { + st->completed = 1; + st->result_cb(CONN_MGR_ERR_TIMEOUT, st->node_id, st->result_arg); + ca_cleanup(st); + } + u_free(ctx); +} + +void chat_core_connect_auto(uint64_t node_id, + void (*cb)(int result, uint64_t node_id, void* arg), + void* arg) { + if (!g_cc.initialized || !g_cc.inst || !cb) return; + + struct ETCP_SOCKET* best_socket = g_cc.inst->etcp_sockets; + while (best_socket && best_socket->local_addr.ss_family != AF_INET) best_socket = best_socket->next; + if (!best_socket) { + DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect no AF_INET socket for node 0x%016llx", + CC_ID, (unsigned long long)node_id); + cb(CONN_MGR_ERR_INTERNAL, node_id, arg); + return; + } + + uint8_t pubkey[SC_PUBKEY_SIZE] = {0}; + sqlite3_stmt* st = NULL; + if (sqlite3_prepare_v2(g_cc.db, + "SELECT x25519_pubkey FROM nodes WHERE node_id=?", + -1, &st, NULL) == SQLITE_OK) { + sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); + if (sqlite3_step(st) == SQLITE_ROW) { + const void* pk = sqlite3_column_blob(st, 0); + if (pk && sqlite3_column_bytes(st, 0) == SC_PUBKEY_SIZE) + memcpy(pubkey, pk, SC_PUBKEY_SIZE); + } + sqlite3_finalize(st); + } + if (pubkey[0] == 0 && pubkey[1] == 0) { + DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect no pubkey for node 0x%016llx", + CC_ID, (unsigned long long)node_id); + cb(CONN_MGR_ERR_NOT_FOUND, node_id, arg); + return; + } + + /* collect IPv4 addresses */ + struct { uint8_t addr[4]; uint16_t port; } addrs[16]; + int addr_count = 0; + if (sqlite3_prepare_v2(g_cc.db, + "SELECT address, port FROM node_addresses WHERE node_id=? AND family=4", + -1, &st, NULL) == SQLITE_OK) { + sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); + while (sqlite3_step(st) == SQLITE_ROW && addr_count < 16) { + const void* a = sqlite3_column_blob(st, 0); + int alen = sqlite3_column_bytes(st, 0); + if (a && alen == 4) { + memcpy(addrs[addr_count].addr, a, 4); + addrs[addr_count].port = (uint16_t)sqlite3_column_int(st, 1); + addr_count++; + } + } + sqlite3_finalize(st); + } + if (addr_count == 0) { + DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect no IPv4 addrs for node 0x%016llx", + CC_ID, (unsigned long long)node_id); + cb(CONN_MGR_ERR_NO_ADDRESSES, node_id, arg); + return; + } + + struct ca_state* pst = u_calloc(1, sizeof(struct ca_state)); + if (!pst) { cb(CONN_MGR_ERR_INTERNAL, node_id, arg); return; } + pst->inst = g_cc.inst; + pst->addr_count = addr_count; + pst->pending_count = addr_count; + pst->node_id = node_id; + pst->result_cb = cb; + pst->result_arg = arg; + memcpy(pst->pubkey, pubkey, SC_PUBKEY_SIZE); + pst->conns = u_calloc(addr_count, sizeof(struct ETCP_CONN*)); + pst->timers = u_calloc(addr_count, sizeof(void*)); + if (!pst->conns || !pst->timers) { ca_cleanup(pst); cb(CONN_MGR_ERR_INTERNAL, node_id, arg); return; } + + for (int i = 0; i < addr_count; i++) { + struct sockaddr_in sin; + memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET; + memcpy(&sin.sin_addr.s_addr, addrs[i].addr, 4); + sin.sin_port = htons(addrs[i].port); + struct sockaddr_storage sa; + memset(&sa, 0, sizeof(sa)); memcpy(&sa, &sin, sizeof(sin)); + + struct ETCP_CONN* conn = etcp_connection_create(g_cc.inst, NULL); + if (!conn) { + DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect etcp_connection_create failed idx=%d", CC_ID, i); + pst->conns[i] = NULL; pst->timers[i] = NULL; pst->pending_count--; continue; + } + sc_init_ctx(&conn->crypto_ctx, &g_cc.inst->my_keys); + sc_set_peer_public_key(&conn->crypto_ctx, pubkey, 0); + + struct ca_ctx* pctx = u_calloc(1, sizeof(struct ca_ctx)); + if (!pctx) { etcp_connection_close(conn); pst->conns[i] = NULL; pst->timers[i] = NULL; pst->pending_count--; continue; } + pctx->state = pst; pctx->addr_index = i; + + etcp_conn_set_ready_cbk(conn, ca_ready_cb, pctx); + + if (!etcp_link_new(conn, best_socket, &sa, 0)) { + DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect etcp_link_new failed idx=%d", CC_ID, i); + u_free(pctx); etcp_connection_close(conn); + pst->conns[i] = NULL; pst->timers[i] = NULL; pst->pending_count--; continue; + } + pst->conns[i] = conn; + pst->timers[i] = uasync_set_timeout(g_cc.inst->ua, CA_CONNECT_TIMEOUT_MS * 10, + pctx, ca_timeout_cb, "ca_timeout"); + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect attempt %d/%d to %d.%d.%d.%d:%d", + CC_ID, i + 1, addr_count, + addrs[i].addr[0], addrs[i].addr[1], addrs[i].addr[2], addrs[i].addr[3], + addrs[i].port); + } + + if (pst->pending_count <= 0 && !pst->completed) { + DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect all %d attempts failed to start", CC_ID, addr_count); + cb(CONN_MGR_ERR_UNREACHABLE, node_id, arg); + ca_cleanup(pst); + } +} + /* ─── подготовка инфраструктуры канала (msg-таблица + db_sync instance) ─── */ void chat_core_ensure_channel_ready(const char* ch_id) { diff --git a/tools/chatgui/transport/chat_core.h b/tools/chatgui/transport/chat_core.h index 6af3978d..7ffa24cf 100644 --- a/tools/chatgui/transport/chat_core.h +++ b/tools/chatgui/transport/chat_core.h @@ -54,6 +54,12 @@ struct chat_invite { void chat_core_connect_from_invite(struct chat_invite* inv); +/* Прямое ETCP-подключение к узлу без BGP: pubkey+адреса из SQLite, + * коллбэк вызывается с CONN_MGR_OK / CONN_MGR_ERR_TIMEOUT / CONN_MGR_ERR_NO_ADDRESSES */ +void chat_core_connect_auto(uint64_t node_id, + void (*cb)(int result, uint64_t node_id, void* arg), + void* arg); + /* ── Создание канала (GUI → uasync) ── */ struct chat_channel_create { diff --git a/tools/chatgui/transport/chat_sync.c b/tools/chatgui/transport/chat_sync.c index 239b1a35..723277bc 100644 --- a/tools/chatgui/transport/chat_sync.c +++ b/tools/chatgui/transport/chat_sync.c @@ -18,10 +18,275 @@ #include "../../../lib/platform_compat.h" #include +#include #include static struct chat_sync* g_cs = NULL; +/* ═══════════════════════════════════════════════════════════════════════ + * Auto-connect: parallel connect to SQLite peers via direct ETCP (no BGP) + * ══════════════════════════════════════════════════════════════════════ */ + +#define AC_MAX_PARALLEL 10 +#define AC_TARGET_SUCCESS 3 +#define AC_RETRY_INTERVAL_MS 10000 +#define AC_ID "auto_connect" + +struct auto_connect; + +struct ac_flight { + struct auto_connect* ac; + uint64_t node_id; +}; + +struct auto_connect { + struct UTUN_INSTANCE* inst; + uint64_t* node_ids; + int node_count; + int next_index; + int in_flight; + int success_count; + int total_tried; + void* retry_timer; + uint8_t active; +}; + +static struct auto_connect* g_ac = NULL; + +static void ac_result_cb(int result, uint64_t node_id, void* arg); +static void ac_retry_timer_cb(void* arg); +static int ac_connect_one(struct auto_connect* ac, uint64_t nid); +static void ac_launch_batch(struct auto_connect* ac); +static int ac_collect_nodes(struct UTUN_INSTANCE* inst, uint64_t** out, int* out_count); + +/* ── helper: count connected nodes (excl. self) ── */ + +static int ac_count_connected(struct UTUN_INSTANCE* inst) { + int cnt = 0; + struct ETCP_CONN* c = inst->connections; + while (c) { + if (c->peer_node_id && c->peer_node_id != inst->node_id) { + struct ETCP_LINK* l = c->links; + while (l) { if (l->initialized && l->link_status) { cnt++; break; } l = l->next; } + } + c = c->next; + } + return cnt; +} + +/* ── helper: check if node already has active ETCP connection ── */ + +static int ac_node_busy(struct UTUN_INSTANCE* inst, uint64_t nid) { + struct ETCP_CONN* c = inst->connections; + while (c) { + if (c->peer_node_id == nid) { + struct ETCP_LINK* l = c->links; + while (l) { if (l->initialized && l->link_status) return 1; l = l->next; } + return 1; /* conn exists even without link_status yet */ + } + c = c->next; + } + return 0; +} + +/* ── sort: read RTT from node_addresses, sort nodes ascending ── */ + +static void ac_sort_by_rtt(sqlite3* db, uint64_t* ids, int count) { + if (!db || count < 2) return; + struct ac_rtt_item { uint64_t id; int32_t rtt; }; + struct ac_rtt_item* items = u_malloc((size_t)count * sizeof(*items)); + if (!items) return; + for (int i = 0; i < count; i++) { + items[i].id = ids[i]; items[i].rtt = INT32_MAX; + } + sqlite3_stmt* st = NULL; + if (sqlite3_prepare_v2(db, + "SELECT node_id, rtt FROM node_addresses WHERE rtt IS NOT NULL ORDER BY rtt ASC", + -1, &st, NULL) == SQLITE_OK) { + while (sqlite3_step(st) == SQLITE_ROW) { + uint64_t nid = (uint64_t)sqlite3_column_int64(st, 0); + int32_t rtt = sqlite3_column_int(st, 1); + for (int i = 0; i < count; i++) if (items[i].id == nid) { items[i].rtt = rtt; break; } + } + sqlite3_finalize(st); + } + for (int i = 1; i < count; i++) { + struct ac_rtt_item key = items[i]; int j = i - 1; + while (j >= 0 && items[j].rtt > key.rtt) { items[j+1] = items[j]; j--; } + items[j+1] = key; + } + for (int i = 0; i < count; i++) ids[i] = items[i].id; + u_free(items); +} + +/* ── collect + sort nodes from SQLite peers (all channels) ── */ + +static int ac_collect_nodes(struct UTUN_INSTANCE* inst, uint64_t** out, int* out_count) { + int cap = 32, cnt = 0; + uint64_t* ids = u_malloc((size_t)cap * sizeof(uint64_t)); + if (!ids) return -1; + + uint8_t ch_buf[4096]; size_t ch_len; + if (chat_core_list_channels(ch_buf, sizeof(ch_buf), &ch_len) != 0 || ch_len < 2) + { if (cnt == 0) { u_free(ids); *out = NULL; *out_count = 0; return 0; } } + + uint16_t ch_cnt; memcpy(&ch_cnt, ch_buf, 2); + const uint8_t* p = ch_buf + 2; size_t rem = ch_len - 2; + for (uint16_t ci = 0; ci < ch_cnt && rem >= 1; ci++) { + uint8_t id_len = *p++; rem--; + if (rem < id_len) break; + char cid[64]; memcpy(cid, p, id_len); cid[id_len] = '\0'; p += id_len; rem -= id_len; + + uint8_t pb[2048]; size_t plen; + if (chat_core_list_peers(cid, pb, sizeof(pb), &plen) != 0 || plen < 2) continue; + uint16_t pc; memcpy(&pc, pb, 2); + for (uint16_t pj = 0; pj < pc && (size_t)(2 + (pj+1)*8) <= plen; pj++) { + uint64_t nid; memcpy(&nid, pb + 2 + pj*8, 8); + if (nid == 0 || nid == inst->node_id) continue; + if (ac_node_busy(inst, nid)) continue; + int dup = 0; + for (int j = 0; j < cnt; j++) if (ids[j] == nid) { dup = 1; break; } + if (dup) continue; + if (cnt >= cap) { cap *= 2; ids = u_realloc(ids, (size_t)cap * sizeof(uint64_t)); if (!ids) return -1; } + ids[cnt++] = nid; + } + } + + if (cnt > 0) { + sqlite3* db = inst->topo_groups ? inst->topo_groups->topo_sqlite_db : NULL; + if (db) ac_sort_by_rtt(db, ids, cnt); + } + *out = ids; *out_count = cnt; + return cnt; +} + +/* ── connect to node via direct ETCP (timeout handled by chat_core_connect_auto) ── */ + +static int ac_connect_one(struct auto_connect* ac, uint64_t nid) { + struct ac_flight* f = u_calloc(1, sizeof(*f)); + if (!f) return -1; + f->ac = ac; f->node_id = nid; + ac->in_flight++; + chat_core_connect_auto(nid, ac_result_cb, f); + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: connecting to node 0x%016llx (in_flight=%d)", + AC_ID, (unsigned long long)nid, ac->in_flight); + return 0; +} + +/* ── result callback (fired by chat_core_connect_auto) ── */ + +static void ac_result_cb(int result, uint64_t node_id, void* arg) { + struct ac_flight* f = (struct ac_flight*)arg; + struct auto_connect* ac = f->ac; + if (!ac->active) { u_free(f); return; } + ac->in_flight--; ac->total_tried++; + if (result == CONN_MGR_OK) ac->success_count++; + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: result node=0x%016llx result=%d succ=%d in_flight=%d total=%d", + AC_ID, (unsigned long long)node_id, result, ac->success_count, ac->in_flight, ac->total_tried); + uint8_t evt[7]; evt[0] = ac->success_count >= AC_TARGET_SUCCESS ? 1 : 0; + uint16_t t = (uint16_t)ac->total_tried, n = (uint16_t)ac->node_count, s = (uint16_t)ac->success_count; + memcpy(evt + 1, &t, 2); memcpy(evt + 3, &n, 2); memcpy(evt + 5, &s, 2); + gui_bridge_post(GUI_EVT_AUTO_CONNECT_STATUS, evt, 7); + u_free(f); + if (!ac->active) return; + while (ac->in_flight < AC_MAX_PARALLEL && ac->next_index < ac->node_count) { + uint64_t nid = ac->node_ids[ac->next_index++]; + ac_connect_one(ac, nid); + } + if (ac->in_flight == 0 && ac->success_count < AC_TARGET_SUCCESS) + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: all tries done, %d/%d connected, waiting retry", + AC_ID, ac->success_count, ac->total_tried); +} + +/* ── launch initial batch ── */ + +static void ac_launch_batch(struct auto_connect* ac) { + while (ac->in_flight < AC_MAX_PARALLEL && ac->next_index < ac->node_count) { + uint64_t nid = ac->node_ids[ac->next_index++]; + ac_connect_one(ac, nid); + } +} + +/* ── retry timer callback ── */ + +static void ac_retry_timer_cb(void* arg) { + struct auto_connect* ac = (struct auto_connect*)arg; + ac->retry_timer = NULL; + if (!ac->active) return; + int connected = ac_count_connected(ac->inst); + if (connected >= AC_TARGET_SUCCESS) { + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: %d connected, target reached, stopping", + AC_ID, connected); + chat_sync_auto_connect_stop(); + return; + } + if (ac->in_flight > 0) { + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: retry, %d in_flight, waiting", AC_ID, ac->in_flight); + } else { + if (ac->node_ids) { u_free(ac->node_ids); ac->node_ids = NULL; } + ac->node_count = 0; ac->next_index = 0; + ac_collect_nodes(ac->inst, &ac->node_ids, &ac->node_count); + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: retry, %d nodes to try, %d connected", + AC_ID, ac->node_count, connected); + if (ac->node_count > 0) { + uint8_t evt[7]; evt[0] = 0; + uint16_t t = 0, n = (uint16_t)ac->node_count, s = (uint16_t)ac->success_count; + memcpy(evt + 1, &t, 2); memcpy(evt + 3, &n, 2); memcpy(evt + 5, &s, 2); + gui_bridge_post(GUI_EVT_AUTO_CONNECT_STATUS, evt, 7); + ac_launch_batch(ac); + } + } + ac->retry_timer = uasync_set_timeout(ac->inst->ua, AC_RETRY_INTERVAL_MS * 10, ac, ac_retry_timer_cb, "ac_retry"); +} + +/* ═══════════════════════════════════════════════════════════════════════ + * Public auto_connect API + * ══════════════════════════════════════════════════════════════════════ */ + +void chat_sync_auto_connect_start(struct UTUN_INSTANCE* inst) { + if (!inst || !inst->ua) return; + chat_sync_auto_connect_stop(); + struct auto_connect* ac = u_calloc(1, sizeof(*ac)); + if (!ac) return; + ac->inst = inst; + ac->active = 1; + g_ac = ac; + int connected = ac_count_connected(inst); + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: start, %d already connected", AC_ID, connected); + if (connected >= AC_TARGET_SUCCESS) { + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: already have %d connections, idle", AC_ID, connected); + ac->retry_timer = uasync_set_timeout(inst->ua, AC_RETRY_INTERVAL_MS * 10, ac, ac_retry_timer_cb, "ac_retry"); + return; + } + ac_collect_nodes(inst, &ac->node_ids, &ac->node_count); + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: collected %d candidate nodes", AC_ID, ac->node_count); + if (ac->node_count > 0) { + uint8_t evt[7]; evt[0] = 0; uint16_t z = 0, n = (uint16_t)ac->node_count; + memcpy(evt + 1, &z, 2); memcpy(evt + 3, &n, 2); memcpy(evt + 5, &z, 2); + gui_bridge_post(GUI_EVT_AUTO_CONNECT_STATUS, evt, 7); + ac_launch_batch(ac); + } + ac->retry_timer = uasync_set_timeout(inst->ua, AC_RETRY_INTERVAL_MS * 10, ac, ac_retry_timer_cb, "ac_retry"); +} + +void chat_sync_auto_connect_stop(void) { + struct auto_connect* ac = g_ac; + if (!ac) return; + ac->active = 0; g_ac = NULL; + if (ac->retry_timer) { uasync_cancel_timeout(ac->inst->ua, ac->retry_timer); ac->retry_timer = NULL; } + uint8_t evt[7]; evt[0] = 2; + uint16_t z = 0; memcpy(evt + 1, &z, 2); memcpy(evt + 3, &z, 2); memcpy(evt + 5, &z, 2); + gui_bridge_post(GUI_EVT_AUTO_CONNECT_STATUS, evt, 7); + /* DON'T free ac — in-flight callbacks still reference it via ac_flight */ + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: stopped", AC_ID); +} + +void chat_sync_auto_connect_switch_group(struct UTUN_INSTANCE* inst, uint64_t new_group_id) { + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: switching to group 0x%016llx", AC_ID, (unsigned long long)new_group_id); + chat_sync_auto_connect_stop(); + chat_sync_auto_connect_start(inst); +} + struct channel_cache { char channel_id[64]; uint64_t* peer_ids; @@ -418,6 +683,39 @@ static void cs_join_timeout_cb(void* arg) { cs->pending_invite_ch_id = 0; } +/* ── helper: check if peer has active ETCP link ── */ + +static int cs_is_peer_online(struct UTUN_INSTANCE* inst, uint64_t peer_id) { + struct ETCP_CONN* c = inst->connections; + while (c) { + if (c->peer_node_id == peer_id) { + struct ETCP_LINK* l = c->links; + while (l) { if (l->initialized && l->link_status) return 1; l = l->next; } + } + c = c->next; + } + return 0; +} + +/* ── post per-channel online peers count to GUI ── */ + +static void cs_post_channel_online(struct chat_sync* cs, const char* ch_id) { + int online = 0; + struct channel_cache* ch = NULL; + for (int i = 0; i < cs->channel_count; i++) + if (strcmp(cs->channels[i].channel_id, ch_id) == 0) { ch = &cs->channels[i]; break; } + if (ch) { + for (int j = 0; j < ch->peer_count; j++) + if (cs_is_peer_online(cs->inst, ch->peer_ids[j])) online++; + } + size_t cl = strlen(ch_id); + if (cl > 63) cl = 63; + uint8_t evt[66]; evt[0] = (uint8_t)cl; + memcpy(evt + 1, ch_id, cl); + uint16_t oc = (uint16_t)online; memcpy(evt + 1 + cl, &oc, 2); + gui_bridge_post(GUI_EVT_CHANNEL_PEERS_ONLINE, evt, 1 + (int)cl + 2); +} + static void _on_member_sync_done(uint64_t peer, const char* ns, int result, void* arg) { struct channel_cache* ch = (struct channel_cache*)arg; if (result == MT_OK && ch) ch->synced = CS_SYNC_DONE; @@ -462,6 +760,7 @@ static void cs_on_conn_up(struct ETCP_CONN* conn, void* arg) { memcpy(msg + 5, ch->last_chain_hash, 32); cs_send(g_cs, ch->channel_id, peer, msg, 37); member_sync_start(g_cs->inst, peer, ch->channel_id, _on_member_sync_done, ch); + cs_post_channel_online(g_cs, ch->channel_id); } if (synced_cnt > 0) { DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: conn_up sync path peer=%016llx, syncing %d channels", @@ -474,6 +773,14 @@ 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; + uint16_t rtt = conn->rtt_avg_100; + if (rtt > 0 && peer != 0 && g_cs->inst->topo_groups && g_cs->inst->topo_groups->topo_sqlite_db) { + sqlite3* db = g_cs->inst->topo_groups->topo_sqlite_db; + sqlite3_stmt* st = NULL; + sqlite3_prepare_v2(db, "UPDATE node_addresses SET rtt=? WHERE node_id=?", -1, &st, NULL); + if (st) { sqlite3_bind_int(st, 1, (int)rtt); sqlite3_bind_int64(st, 2, (sqlite3_int64)peer); sqlite3_step(st); sqlite3_finalize(st); } + DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: conn_down saved rtt=%u for node %016llx", CS_ID, rtt, (unsigned long long)peer); + } cs_cancel_proto_timers(g_cs); if (g_cs->pending_invite_node_id == peer) { DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: conn down while waiting invite resp peer=%016llx, timers cancelled, state kept for retry on reconnect", @@ -489,6 +796,7 @@ static void cs_on_conn_down(struct ETCP_CONN* conn, void* arg) { uint8_t rem[9]; rem[0] = CS_MSG_PEER_REMOVE; memcpy(rem + 1, &peer, 8); cs_propagate(g_cs, g_cs->channels[i].channel_id, peer, rem, 9); + cs_post_channel_online(g_cs, g_cs->channels[i].channel_id); } } } @@ -604,6 +912,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; + chat_sync_auto_connect_stop(); member_sync_destroy(inst); cs->initialized = 0; g_cs = NULL; diff --git a/tools/chatgui/transport/chat_sync.h b/tools/chatgui/transport/chat_sync.h index b31ff567..e679d04c 100644 --- a/tools/chatgui/transport/chat_sync.h +++ b/tools/chatgui/transport/chat_sync.h @@ -61,6 +61,11 @@ void chat_sync_connect_from_invite(uint64_t channel_id, uint64_t node_id, const uint8_t* pubkey_bin, const uint8_t* addrs_data, int addr_count); +/* Авто-подключение к узлам группы (параллельно до 10, останавливается при 3 успешных) */ +void chat_sync_auto_connect_start(struct UTUN_INSTANCE* inst); +void chat_sync_auto_connect_stop(void); +void chat_sync_auto_connect_switch_group(struct UTUN_INSTANCE* inst, uint64_t new_group_id); + #ifdef __cplusplus } #endif diff --git a/tools/chatgui/transport/gui_bridge.h b/tools/chatgui/transport/gui_bridge.h index c9c9e122..c2e5e675 100644 --- a/tools/chatgui/transport/gui_bridge.h +++ b/tools/chatgui/transport/gui_bridge.h @@ -17,7 +17,9 @@ struct UASYNC; #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] */ +#define GUI_EVT_MY_NODE_ID 6 /* data: [node_id:8] */ +#define GUI_EVT_AUTO_CONNECT_STATUS 7 /* data: [status:1][total_tried:2][node_count:2][success:2] */ +#define GUI_EVT_CHANNEL_PEERS_ONLINE 8 /* data: [ch_id_len:1][ch_id:var][online_count:2] */ /* ── API ── */ @@ -56,6 +58,15 @@ void gui_bridge_set_members_changed_cb(gui_members_changed_fn cb); 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); +/* Callback для статуса авто-подключения (вызывается из GUI-потока) + * status: 0=in_progress, 1=done(3+), 2=stopped */ +typedef void (*gui_auto_connect_status_fn)(uint8_t status, uint16_t total_tried, uint16_t node_count, uint16_t connected); +void gui_bridge_set_auto_connect_status_cb(gui_auto_connect_status_fn cb); + +/* Callback для обновления числа онлайн-пиров канала (вызывается из GUI-потока) */ +typedef void (*gui_channel_peers_online_fn)(const char* ch_id, int ch_id_len, uint16_t online_count); +void gui_bridge_set_channel_peers_online_cb(gui_channel_peers_online_fn cb); + #ifdef __cplusplus } #endif diff --git a/tools/chatgui/transport/gui_bridge_impl.cpp b/tools/chatgui/transport/gui_bridge_impl.cpp index 27c9a2e9..4c71d7c8 100644 --- a/tools/chatgui/transport/gui_bridge_impl.cpp +++ b/tools/chatgui/transport/gui_bridge_impl.cpp @@ -29,6 +29,8 @@ static gui_msg_received_fn g_msg_received_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 gui_auto_connect_status_fn g_auto_connect_status_cb = nullptr; +static gui_channel_peers_online_fn g_channel_peers_online_cb = nullptr; static struct UASYNC* g_ua = nullptr; /* ── GuiBridgeReceiver implementation ── */ @@ -98,6 +100,27 @@ void GuiBridgeReceiver::processPost(int eventType, QByteArray data) { if (g_my_node_id_cb) g_my_node_id_cb(nid); } break; + case GUI_EVT_AUTO_CONNECT_STATUS: + if (dlen >= 7) { + uint8_t status = d[0]; + uint16_t total_tried; memcpy(&total_tried, d + 1, 2); + uint16_t node_count; memcpy(&node_count, d + 3, 2); + uint16_t connected; memcpy(&connected, d + 5, 2); + const char* st_str = status == 0 ? "in_progress" : status == 1 ? "done" : "stopped"; + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "gui_bridge: AUTO_CONNECT %s tried=%u/%u connected=%u", st_str, total_tried, node_count, connected); + if (g_auto_connect_status_cb) g_auto_connect_status_cb(status, total_tried, node_count, connected); + } + break; + case GUI_EVT_CHANNEL_PEERS_ONLINE: + if (dlen >= 3) { + uint8_t chLen = d[0]; + if (dlen >= 3 + chLen) { + uint16_t online; memcpy(&online, d + 1 + chLen, 2); + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "gui_bridge: CHANNEL_PEERS_ONLINE ch=%.*s online=%u", chLen, (const char*)d + 1, online); + if (g_channel_peers_online_cb) g_channel_peers_online_cb((const char*)d + 1, chLen, online); + } + } + break; default: DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "gui_bridge: unknown event type %d", eventType); break; @@ -164,6 +187,14 @@ void gui_bridge_set_my_node_id_cb(gui_my_node_id_fn cb) { g_my_node_id_cb = cb; } +void gui_bridge_set_auto_connect_status_cb(gui_auto_connect_status_fn cb) { + g_auto_connect_status_cb = cb; +} + +void gui_bridge_set_channel_peers_online_cb(gui_channel_peers_online_fn cb) { + g_channel_peers_online_cb = cb; +} + } /* extern "C" */ #include "gui_bridge_impl.moc" diff --git a/tools/chatgui/transport/node_config.cpp b/tools/chatgui/transport/node_config.cpp index 62dfd581..42517edf 100644 --- a/tools/chatgui/transport/node_config.cpp +++ b/tools/chatgui/transport/node_config.cpp @@ -72,6 +72,10 @@ static const QSet NETWORK_KEYS = { "id", "pubkey", "signing_key" }; +static const QSet NTP_KEYS = { + "enabled", "server", "interval" +}; + static const QStringList VALID_DEBUG_LEVELS = { "none", "error", "warn", "info", "debug", "trace" }; @@ -92,6 +96,7 @@ static bool isSectionValid(const QString& section) { if (section == "tcp_proxy_client") return true; if (section == "tcp_proxy_server") return true; if (section == "msg_transport") return true; + if (section == "ntp") return true; return false; } @@ -109,6 +114,7 @@ static const QSet* keysForSection(const QString& section) { if (section == "tcp_proxy_client") return &TCP_PROXY_CLIENT_KEYS; if (section == "tcp_proxy_server") return &TCP_PROXY_SERVER_KEYS; if (section == "msg_transport") return &MSG_TRANSPORT_KEYS; + if (section == "ntp") return &NTP_KEYS; return nullptr; // [debug] — free-form, no key validation } @@ -339,6 +345,16 @@ bool NodeConfig::save() { out << "debug_level=" << (m_debugLevel.isEmpty() ? "info" : m_debugLevel) << "\n"; out << "debug_categories=" << m_debugCategories << "\n\n"; + out << "[ntp]\n"; + out << "enabled=yes\n"; + out << "server=pool.ntp.org\n"; + out << "server=time.google.com\n"; + out << "server=time.cloudflare.com\n"; + out << "server=time.windows.com\n"; + out << "server=ntp1.vniiftri.ru\n"; + out << "server=ntp2.vniiftri.ru\n"; + out << "interval=3600\n\n"; + f.close(); return true; } diff --git a/tools/chatgui/transport/utun_node.cpp b/tools/chatgui/transport/utun_node.cpp index dbc5fdd3..d9991f24 100644 --- a/tools/chatgui/transport/utun_node.cpp +++ b/tools/chatgui/transport/utun_node.cpp @@ -241,7 +241,8 @@ void UtunNode::runLoop() { /* Initialize chat_core (DB) and chat_sync (channel/message P2P sync) */ chat_core_init(m_instance, m_dbPath.toUtf8().constData()); chat_sync_init(m_instance, nullptr); - DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "chat_core + chat_sync initialized"); + chat_sync_auto_connect_start(m_instance); + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "chat_core + chat_sync + auto_connect initialized"); DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "utun_node: entering poll loop"); while (!m_stop) {