You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
 
 
 
 
 
 

418 lines
15 KiB

// utun_node.cpp — embedded uTun node running in a dedicated thread
#include "utun_node.h"
#include <cinttypes>
#include <QDir>
#include <QMap>
#include <QPair>
#include <csignal>
#include "../../lib/socket_compat.h"
#include "../db/db_manager.h"
#include "../src/invite_link.h"
// === C headers ===
extern "C" {
#include "utun_instance.h"
#include "etcp_router.h"
#include "etcp_api.h"
#include "etcp.h"
#include "config_parser.h"
#include "config_updater.h"
#include "control_server.h"
#include "secure_channel.h"
#include "../lib/u_async.h"
#include "../lib/ll_queue.h"
#include "../lib/memory_pool.h"
#include "../lib/debug_config.h"
#include "../lib/mem.h"
#include "chat/chat_sync.h"
#include "chat/chat_core.h"
#include "chat/chat_event.h"
#include "chat/db_sync.h"
#include "gui_bridge.h"
#include "topo_node_sqlite.h"
#include "topo_node.h"
#include "topo_group.h"
}
#define ETCP_RT_ID_CHAT 0x11
static UtunNode* g_currentNode = nullptr;
static void signal_handler(int sig) {
(void)sig;
if (g_currentNode) g_currentNode->requestStop();
}
UtunNode::UtunNode(QObject* parent) : QObject(parent) {}
UtunNode::~UtunNode() {
stop();
}
bool UtunNode::start(const QString& configPath) {
if (m_running) return false;
m_configPath = configPath;
m_stop = false;
m_running = true;
m_thread = std::thread(&UtunNode::runLoop, this);
for (int i = 0; i < 50 && m_running && !m_instance; ++i)
std::this_thread::sleep_for(std::chrono::milliseconds(10));
if (!m_instance) {
m_stop = true;
if (m_thread.joinable()) m_thread.join();
m_running = false;
return false;
}
emit started();
return true;
}
void UtunNode::setDebugFile(const QString& path) {
m_debugFile = path;
}
void UtunNode::setDbPath(const QString& path) {
m_dbPath = path;
QDir(path).mkpath(".");
}
void UtunNode::setDebugLevel(const QString& level) {
m_debugLevel = level;
}
void UtunNode::setDebugCategories(const QString& categories) {
m_debugCategories = categories;
}
void UtunNode::requestStop() {
m_stop = true;
if (m_ua) uasync_wakeup(m_ua);
}
void UtunNode::stop() {
if (!m_running) return;
requestStop();
if (m_thread.joinable()) m_thread.join();
finalize();
m_running = false;
}
void UtunNode::finalize() {
if (m_ua) { uasync_destroy(m_ua, 0); m_ua = nullptr; }
m_instance = nullptr;
m_running = false;
DEBUG_DEBUG(DEBUG_CATEGORY_GENERAL, "finalize: cleanup done");
}
void UtunNode::send(uint64_t dstNodeId, const QByteArray& data) {
if (!m_instance) return;
struct ll_entry* entry = queue_entry_new(data.size());
if (!entry) return;
memcpy(entry->data, data.constData(), data.size());
etcp_route_send(m_instance, TOPO_GROUP_UTUN, dstNodeId, entry, 0);
}
QString UtunNode::nodeIdHex() const {
if (!m_instance) return {};
char buf[32];
snprintf(buf, sizeof(buf), "%016" PRIx64, m_instance->node_id);
return QString(buf);
}
QString UtunNode::pubKeyHex() const {
if (!m_instance || !m_instance->config) return {};
return QString::fromLatin1(m_instance->config->global.my_public_key_hex, 64);
}
int64_t UtunNode::ntpOffsetUs() const {
return m_instance ? m_instance->ntp.offset_us : 0;
}
bool UtunNode::ntpSynced() const {
return m_instance ? m_instance->ntp.synced : false;
}
bool UtunNode::ntpEnabled() const {
return m_instance ? m_instance->ntp.enabled : false;
}
QList<NodeAddr> UtunNode::getInviteAddresses(DbManager* db) {
(void)db;
QList<NodeAddr> addrs;
QMap<QPair<QByteArray, quint16>, uint8_t> addrProto; // key=(addr,port) → proto bits
auto add_addr = [&](const QByteArray& addr, quint16 port, int family,
int sockId, uint8_t proto) {
QPair<QByteArray, quint16> key(addr, port);
uint8_t existing = addrProto.value(key, 0);
addrProto[key] = existing | proto;
// Store sockId/family for serialization (overwrite is fine — same port)
};
auto append_from_addr_proto = [&](const QByteArray& addr, quint16 port,
int family, int sockId, uint8_t proto) {
NodeAddr a;
a.family = family;
a.address = addr;
a.port = port;
a.socketId = sockId;
a.protocol = proto;
addrs.append(a);
};
if (m_instance) {
/* 1) UDP from active ETCP links */
struct ll_entry* entry = m_instance->connections->head;
while (entry) {
struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data;
struct ETCP_CONN* conn = ce->conn;
if (conn && conn->initialized && conn->links_up > 0) {
struct ETCP_LINK* link = conn->links;
while (link) {
if (link->link_state == 3 && link->nat_type == NAT_TYPE_DIRECT) {
struct sockaddr_storage* sa = &link->remote_addr;
if (sa->ss_family == AF_INET) {
struct sockaddr_in* sin = (struct sockaddr_in*)sa;
QByteArray addr((const char*)&sin->sin_addr, 4);
add_addr(addr, ntohs(sin->sin_port), 4,
link->conn ? (int)link->conn->sock_id : 0, INVITE_PROTO_UDP);
} else {
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)sa;
QByteArray addr((const char*)&sin6->sin6_addr, 16);
add_addr(addr, ntohs(sin6->sin6_port), 6,
link->conn ? (int)link->conn->sock_id : 0, INVITE_PROTO_UDP);
}
}
link = link->next;
}
}
entry = entry->next;
}
/* 2) UDP from local ETCP sockets (fallback) */
{
struct ETCP_SOCKET* sock = m_instance->etcp_sockets;
while (sock) {
struct sockaddr_storage* sa = NULL;
if (sock->nat_addr.ss_family != 0)
sa = &sock->nat_addr;
else if (sock->interface_addr.ss_family != 0)
sa = &sock->interface_addr;
if (sa) {
if (sa->ss_family == AF_INET) {
struct sockaddr_in* sin = (struct sockaddr_in*)sa;
QByteArray addr((const char*)&sin->sin_addr, 4);
add_addr(addr, ntohs(sin->sin_port), 4, (int)sock->sock_id, INVITE_PROTO_UDP);
} else {
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)sa;
QByteArray addr((const char*)&sin6->sin6_addr, 16);
add_addr(addr, ntohs(sin6->sin6_port), 6, (int)sock->sock_id, INVITE_PROTO_UDP);
}
}
sock = sock->next;
}
}
/* 3) TCP from local TCP sockets */
{
struct ETCP_SOCKET* s = m_instance->etcp_sockets;
while (s) {
if (!s->is_tcp) { s = s->next; continue; }
struct sockaddr_storage* sa = &s->interface_addr;
if (sa->ss_family == AF_INET) {
struct sockaddr_in* sin = (struct sockaddr_in*)sa;
QByteArray addr((const char*)&sin->sin_addr, 4);
add_addr(addr, ntohs(sin->sin_port), 4, (int)s->sock_id, INVITE_PROTO_TCP);
} else if (sa->ss_family == AF_INET6) {
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)sa;
QByteArray addr((const char*)&sin6->sin6_addr, 16);
add_addr(addr, ntohs(sin6->sin6_port), 6, (int)s->sock_id, INVITE_PROTO_TCP);
}
s = s->next;
}
}
/* 4) Merge: dedup by (addr,port), write result with merged proto */
for (auto it = addrProto.begin(); it != addrProto.end(); ++it) {
QByteArray addr = it.key().first;
quint16 port = it.key().second;
uint8_t proto = it.value();
int family = (addr.size() == 16) ? 6 : 4;
append_from_addr_proto(addr, port, family, 0, proto);
}
}
return addrs;
}
void UtunNode::recvCallback(struct ETCP_CONN* conn, struct ll_entry* entry) {
if (!g_currentNode || !entry) return;
uint64_t src = conn ? conn->peer_node_id : 0;
QByteArray data((const char*)entry->data, (int)(entry->len ? entry->len : 0));
QMetaObject::invokeMethod(g_currentNode, [=] {
emit g_currentNode->messageReceived(src, data);
}, Qt::QueuedConnection);
}
void UtunNode::runLoop() {
g_currentNode = this;
debug_config_init();
debug_enable_function_name(0);
/* ── debug level from config ── */
if (!m_debugLevel.isEmpty()) {
QString lvl = m_debugLevel.toLower();
debug_level_t dl = DEBUG_LEVEL_INFO;
if (lvl == "none") dl = DEBUG_LEVEL_NONE;
else if (lvl == "error") dl = DEBUG_LEVEL_ERROR;
else if (lvl == "warn") dl = DEBUG_LEVEL_WARN;
else if (lvl == "info") dl = DEBUG_LEVEL_INFO;
else if (lvl == "debug") dl = DEBUG_LEVEL_DEBUG;
else if (lvl == "trace") dl = DEBUG_LEVEL_TRACE;
g_debug_config.level = dl;
} else {
g_debug_config.level = DEBUG_LEVEL_INFO;
}
/* chatgui defaults: suppress spam, verbose for our categories */
g_debug_config.category_levels[DEBUG_CATEGORY_SYS] = DEBUG_LEVEL_DISABLED;
g_debug_config.category_levels[DEBUG_CATEGORY_BGP] = DEBUG_LEVEL_INFO;
/* per-category overrides from config: cat=level,cat=level,... */
if (!m_debugCategories.isEmpty()) {
QByteArray cats = m_debugCategories.toUtf8();
char* token = strtok(cats.data(), ",");
while (token) {
while (*token == ' ') token++;
char* eq = strchr(token, '=');
if (eq) {
*eq = '\0';
char* cat_name = token;
char* lvl_name = eq + 1;
while (*lvl_name == ' ') lvl_name++;
debug_category_t cat = get_category_by_name(cat_name);
debug_level_t lvl = debug_level_from_name(lvl_name);
if (cat != DEBUG_CATEGORY_NONE && lvl != DEBUG_LEVEL_NONE) {
g_debug_config.category_levels[cat] = lvl;
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "debug category: %s=%s",
debug_get_category_name(cat),
lvl_name);
}
}
token = strtok(NULL, ",");
}
}
/* file output (shared with GUI via debug_config) */
if (!m_debugFile.isEmpty()) {
debug_enable_file_output(m_debugFile.toUtf8().constData(), 0);
debug_enable_console(1);
g_debug_config.thread_marker = 2;
}
utun_instance_set_tun_init_enabled(0);
struct UASYNC* ua = uasync_create();
if (!ua) {
QMetaObject::invokeMethod(this, [this] { emit error("uasync_create failed"); });
g_currentNode = nullptr;
return;
}
m_ua = ua;
m_instance = utun_instance_create(ua, m_configPath.toUtf8().constData());
if (!m_instance) {
QMetaObject::invokeMethod(this, [this] { emit error("utun_instance_create failed"); });
uasync_destroy(ua, 0);
m_ua = nullptr;
g_currentNode = nullptr;
return;
}
m_instance->test_user_ptr = this;
/* enable db_sync for chat message P2P sync */
m_instance->config->global.db_sync_enabled = 1;
if (!m_dbPath.isEmpty())
snprintf(m_instance->config->global.db_path, sizeof(m_instance->config->global.db_path),
"%s", m_dbPath.toUtf8().constData());
/* topo_groups_init did not open topo_sqlite_db because db_path was empty at that time.
* Open the single shared SQLite connection here with FULLMUTEX (GUI reads from another thread).
* db_sync_init and chat_core_init will reuse this handle. */
if (m_instance->config->global.db_path[0] && !m_instance->topo_sqlite_db) {
char db_file[512];
snprintf(db_file, sizeof(db_file), "%s/chats.db", m_instance->config->global.db_path);
int rc = sqlite3_open_v2(db_file, &m_instance->topo_sqlite_db,
SQLITE_OPEN_READWRITE | SQLITE_OPEN_CREATE | SQLITE_OPEN_FULLMUTEX, NULL);
if (rc == SQLITE_OK && m_instance->topo_sqlite_db) {
sqlite3_exec(m_instance->topo_sqlite_db, "PRAGMA journal_mode=WAL", NULL, NULL, NULL);
sqlite3_exec(m_instance->topo_sqlite_db, "PRAGMA foreign_keys=ON", NULL, NULL, NULL);
sqlite3_exec(m_instance->topo_sqlite_db, "PRAGMA synchronous=NORMAL", NULL, NULL, NULL);
sqlite3_exec(m_instance->topo_sqlite_db, "PRAGMA wal_autocheckpoint=10000", NULL, NULL, NULL);
topo_node_sqlite_init(m_instance->topo_sqlite_db);
DEBUG_DEBUG(DEBUG_CATEGORY_GENERAL, "SQLite (shared): %s db=%p", db_file, (void*)m_instance->topo_sqlite_db);
} else {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "SQLite open failed: %s rc=%d", db_file, rc);
if (m_instance->topo_sqlite_db) { sqlite3_close(m_instance->topo_sqlite_db); m_instance->topo_sqlite_db = NULL; }
}
}
/* set ua for gui_bridge before init, so GUI can post */
gui_bridge_set_uasync(ua);
if (utun_instance_init(m_instance) != 0) {
QMetaObject::invokeMethod(this, [this] { emit error("utun_instance_init failed"); });
utun_instance_destroy(m_instance);
m_instance = nullptr;
uasync_destroy(ua, 0);
m_ua = nullptr;
g_currentNode = nullptr;
return;
}
{ int sock_count = 0; for (struct ETCP_SOCKET* s = m_instance->etcp_sockets; s; s = s->next) sock_count++;
DEBUG_DEBUG(DEBUG_CATEGORY_GENERAL, "utun_node: after instance_init, etcp_sockets=%d conns=%d",
sock_count, queue_entry_count(m_instance->connections)); }
etcp_router_bind(m_instance, ETCP_RT_ID_CHAT, recvCallback);
utun_add_nodeinfo_cbk(m_instance, gui_nodeinfo_cb_impl, nullptr);
/* Bridge chat events to GUI via gui_bridge */
chat_event_set_handler([](int type, const uint8_t* data, int len) {
gui_bridge_post(type, data, len);
});
/* Initialize chat_core (DB) and chat_sync (channel/message P2P sync) */
chat_core_init(m_instance, QString(m_dbPath + "/chats.db").toUtf8().constData());
chat_sync_init(m_instance);
chat_core_sync_my_addresses();
DEBUG_DEBUG(DEBUG_CATEGORY_GENERAL, "chat_core + chat_sync initialized");
DEBUG_DEBUG(DEBUG_CATEGORY_GENERAL, "utun_node: entering poll loop");
signal(SIGINT, signal_handler);
signal(SIGTERM, signal_handler);
while (!m_stop) {
uasync_poll(ua, 100);
}
DEBUG_DEBUG(DEBUG_CATEGORY_GENERAL, "runLoop: poll exit m_stop=%d", (int)m_stop);
if (m_instance) {
DEBUG_DEBUG(DEBUG_CATEGORY_GENERAL, "runLoop: destroying instance in worker thread");
utun_instance_destroy(m_instance);
m_instance = nullptr;
}
signal(SIGINT, SIG_DFL);
signal(SIGTERM, SIG_DFL);
DEBUG_DEBUG(DEBUG_CATEGORY_GENERAL, "runLoop: emitting stopped");
QMetaObject::invokeMethod(this, [this] { emit stopped(); });
}