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.
333 lines
12 KiB
333 lines
12 KiB
// utun_node.cpp — embedded uTun node running in a dedicated thread |
|
#include "utun_node.h" |
|
#include <cinttypes> |
|
#include <QDir> |
|
#include "../../lib/socket_compat.h" |
|
|
|
#include "../db/db_manager.h" |
|
#include "../src/debug_ui.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_sync.h" |
|
#include "chat_core.h" |
|
#include "db_sync.h" |
|
#include "gui_bridge.h" |
|
#include "../src/topo_node_sqlite.h" |
|
} |
|
|
|
#define ETCP_RT_ID_CHAT 0x11 |
|
|
|
static thread_local UtunNode* g_currentNode = nullptr; |
|
|
|
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::stop() { |
|
if (!m_running) return; |
|
m_stop = true; |
|
if (m_thread.joinable()) m_thread.join(); |
|
m_running = false; |
|
} |
|
|
|
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, 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; |
|
|
|
if (m_instance) { |
|
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) { |
|
NodeAddr a; |
|
a.socketId = link->conn ? (int)link->conn->sock_id : 0; |
|
struct sockaddr_storage* sa = &link->remote_addr; |
|
if (sa->ss_family == AF_INET) { |
|
struct sockaddr_in* sin = (struct sockaddr_in*)sa; |
|
a.family = 4; |
|
a.address = QByteArray((const char*)&sin->sin_addr, 4); |
|
a.port = ntohs(sin->sin_port); |
|
} else { |
|
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)sa; |
|
a.family = 6; |
|
a.address = QByteArray((const char*)&sin6->sin6_addr, 16); |
|
a.port = ntohs(sin6->sin6_port); |
|
} |
|
addrs.append(a); |
|
} |
|
link = link->next; |
|
} |
|
} |
|
entry = entry->next; |
|
} |
|
} |
|
|
|
if (addrs.isEmpty() && m_instance) { |
|
struct ETCP_SOCKET* sock = m_instance->etcp_sockets; |
|
while (sock) { |
|
struct sockaddr_storage* sa; |
|
if (sock->nat_addr.ss_family != 0) |
|
sa = &sock->nat_addr; |
|
else if (sock->interface_addr.ss_family != 0) |
|
sa = &sock->interface_addr; |
|
else { sock = sock->next; continue; } |
|
|
|
NodeAddr a; |
|
a.socketId = (int)sock->sock_id; |
|
if (sa->ss_family == AF_INET) { |
|
struct sockaddr_in* sin = (struct sockaddr_in*)sa; |
|
a.family = 4; |
|
a.address = QByteArray((const char*)&sin->sin_addr, 4); |
|
a.port = ntohs(sin->sin_port); |
|
} else { |
|
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)sa; |
|
a.family = 6; |
|
a.address = QByteArray((const char*)&sin6->sin6_addr, 16); |
|
a.port = ntohs(sin6->sin6_port); |
|
} |
|
addrs.append(a); |
|
sock = sock->next; |
|
} |
|
} |
|
|
|
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_MEMORY] = DEBUG_LEVEL_DISABLED; |
|
g_debug_config.category_levels[DEBUG_CATEGORY_TIMING] = 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_INFO(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; |
|
guiDebugFile(m_debugFile.toUtf8().constData()); |
|
} |
|
|
|
utun_instance_set_tun_init_enabled(0); |
|
utun_instance_set_topo_group_enabled(0); |
|
|
|
struct UASYNC* ua = uasync_create(); |
|
if (!ua) { |
|
QMetaObject::invokeMethod(this, [this] { emit error("uasync_create failed"); }); |
|
g_currentNode = nullptr; |
|
return; |
|
} |
|
|
|
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); |
|
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_INFO(DEBUG_CATEGORY_DEBUG, "SQLite (shared): %s db=%p", db_file, (void*)m_instance->topo_sqlite_db); |
|
} else { |
|
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "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); |
|
g_currentNode = nullptr; |
|
return; |
|
} |
|
|
|
etcp_router_bind(m_instance, ETCP_RT_ID_CHAT, recvCallback); |
|
|
|
/* 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, nullptr); |
|
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) { |
|
uasync_poll(ua, 100); |
|
} |
|
|
|
etcp_router_bind(m_instance, ETCP_RT_ID_CHAT, nullptr); |
|
chat_sync_destroy(m_instance); |
|
chat_core_destroy(m_instance); |
|
utun_instance_destroy(m_instance); |
|
m_instance = nullptr; |
|
uasync_destroy(ua, 0); |
|
g_currentNode = nullptr; |
|
QMetaObject::invokeMethod(this, [this] { emit stopped(); }); |
|
}
|
|
|