Browse Source

db_sync: LMDB→SQLite, multi-instance, etcp_api, chatgui integration

- src/db_sync: LMDB replaced with SQLite (separate tables per instance: db_sync_name_id)
  TTL index (author, timestamp), WAL mode, per-instance peers/timers
  opaque DB_SYNC_INSTANCE* handle, insert callback, select API
  new fields: author_signature, delivered_peers, delivery_chain
  transport: etcp_api direct P2P (was etcp_router)

- src/db_sync.h: user-oriented API docs, multi-instance lifecycle

- lib/Makefile.am: liblmdb sources removed
- lib/sqlite3.c+h: SQLite amalgamation added

- topo_node: LMDB→SQLite (src/topo_node_sqlite.c moved from chatgui)

- chatgui: real db_sync.c replaces db_sync_stub.c
  chat_core: db_sync_insert_signed + dual-write insert callback
  utun_node: db_sync_init/destroy replaces chat_sync_init/destroy
  CMakeLists: stubs removed, real db_sync.c added
topo_upd
Evgeny 3 months ago
parent
commit
f617dcd68e
  1. 2
      AGENTS.md
  2. 10
      lib/Makefile.am
  3. 269376
      lib/sqlite3.c
  4. 14347
      lib/sqlite3.h
  5. 6
      src/Makefile.am
  6. 2
      src/config_parser.h
  7. 1179
      src/db_sync.c
  8. 68
      src/db_sync.h
  9. 43
      src/topo_group.c
  10. 35
      src/topo_group.h
  11. 22
      src/topo_node.h
  12. 66
      src/topo_node_lmdb.c
  13. 23
      src/topo_node_lmdb.h
  14. 4
      src/topo_node_sqlite.c
  15. 0
      src/topo_node_sqlite.h
  16. 60
      tests/test_db_sync.c
  17. 16
      tools/chatgui/CMakeLists.txt
  18. 5
      tools/chatgui/libutun/CMakeLists.txt
  19. 215
      tools/chatgui/transport/chat_core.c
  20. 16
      tools/chatgui/transport/chat_core.h
  21. 9
      tools/chatgui/transport/utun_node.cpp

2
AGENTS.md

@ -10,6 +10,8 @@
This file contains essential information for AI coding agents working in the uTun codebase. This file contains essential information for AI coding agents working in the uTun codebase.
Это devel. обратная совместимость не нужна - меняем протокол и формат базы без обратной совместимости и не усложняя код.
## Quick Reference ## Quick Reference
**Repository:** uTun - Secure VPN tunnel with ETCP protocol **Repository:** uTun - Secure VPN tunnel with ETCP protocol

10
lib/Makefile.am

@ -29,18 +29,16 @@ libuasync_a_SOURCES = \
serialize.h \ serialize.h \
tcp_io.c \ tcp_io.c \
tcp_io.h \ tcp_io.h \
liblmdb/mdb.c \ sqlite3.c \
liblmdb/midl.c \ sqlite3.h
liblmdb/lmdb.h \
liblmdb/midl.h
libuasync_a_CFLAGS = \ libuasync_a_CFLAGS = \
-D_ISOC99_SOURCE \ -D_ISOC99_SOURCE \
-DDEBUG_OUTPUT_STDERR \ -DDEBUG_OUTPUT_STDERR \
-DSQLITE_THREADSAFE=1 \
-g \ -g \
-I$(top_srcdir)/src \ -I$(top_srcdir)/src \
-I$(top_srcdir)/lib \ -I$(top_srcdir)/lib
-I$(srcdir)/liblmdb
# Clean build directory # Clean build directory
clean-local: clean-local:

269376
lib/sqlite3.c

File diff suppressed because it is too large Load Diff

14347
lib/sqlite3.h

File diff suppressed because it is too large Load Diff

6
src/Makefile.am

@ -12,7 +12,7 @@ utun_CORE_SOURCES = \
topo_group.c \ topo_group.c \
route_ping.c \ route_ping.c \
topo_node.c \ topo_node.c \
topo_node_lmdb.c \ topo_node_sqlite.c \
route_connectivity.c \ route_connectivity.c \
conn_mgr.c \ conn_mgr.c \
db_sync.c \ db_sync.c \
@ -65,7 +65,7 @@ libutun_a_SOURCES = \
topo_group.c \ topo_group.c \
route_ping.c \ route_ping.c \
topo_node.c \ topo_node.c \
topo_node_lmdb.c \ topo_node_sqlite.c \
route_connectivity.c \ route_connectivity.c \
conn_mgr.c \ conn_mgr.c \
db_sync.c \ db_sync.c \
@ -119,6 +119,8 @@ utun_CFLAGS = \
-I$(top_srcdir)/lib \ -I$(top_srcdir)/lib \
-I$(top_srcdir)/src/uip \ -I$(top_srcdir)/src/uip \
-g \ -g \
-DUSE_SQLITE \
-DSQLITE_THREADSAFE=1 \
$(DEBUG_FLAGS) $(DEBUG_FLAGS)
# Libraries # Libraries

2
src/config_parser.h

@ -114,7 +114,7 @@ struct global_config {
// Debug and logging configuration // Debug and logging configuration
char log_file[256]; // Path to log file (empty = stdout) char log_file[256]; // Path to log file (empty = stdout)
char db_path[256]; // Path to LMDB nodeinfo database (empty = disabled) char db_path[256]; // Path to SQLite database (empty = disabled)
int db_sync_enabled; // 1 = enable distributed DB sync (default: 0) int db_sync_enabled; // 1 = enable distributed DB sync (default: 0)
uint32_t db_sync_ttl; // TTL for unsent records in seconds (default: 86400) uint32_t db_sync_ttl; // TTL for unsent records in seconds (default: 86400)
char debug_level[16]; // debug level: error, warn, info, debug, trace char debug_level[16]; // debug level: error, warn, info, debug, trace

1179
src/db_sync.c

File diff suppressed because it is too large Load Diff

68
src/db_sync.h

@ -1,4 +1,33 @@
// db_sync.h — Distributed content-addressed table with LMDB storage and peer sync via etcp_router // db_sync.h — Distributed content-addressed table with SQLite + peer sync via etcp_api (direct P2P)
//
// Назначение: децентрализованная реплицируемая таблица JSON-записей между всеми узлами сети.
// Модуль поддерживает несколько независимых инстансов (таблиц), каждый идентифицируется парой (name, id).
//
// Использование:
// 1. Включить в конфиге: db_sync_enabled = 1
// Опционально: db_sync_ttl = 86400 (по умолчанию)
// 2. db_sync_init() — вызывается автоматически при старте utun_instance
// 3. struct DB_SYNC_INSTANCE* si = db_sync_instance_add(inst, "chats", 1);
// Создаёт/регистрирует инстанс. При поднятии соединений автоматически запускает sync.
// 4. db_sync_insert(si, json_data) — вставить запись. Автоматически PUSH всем synced-пирам.
// 5. db_sync_count(si) — количество записей в локальной БД для этого инстанса
// 6. db_sync_instance_remove(si) — удалить инстанс (таблица БД не удаляется)
// 7. db_sync_destroy() — вызывается автоматически при завершении
//
// Синхронизация:
// - При поднятии ETCP-соединения с пиром для каждого инстанса запускается полная синхронизация
// - Каждое сообщение содержит instance_hash (первые 64 бита SHA256(name||id_be)),
// что позволяет маршрутизировать сообщения к нужному инстансу на приёмной стороне
// - Если пир не имеет инстанса с таким хешем — возвращает DB_MSG_ERROR(0x01)
// - Если пир добавляет инстанс позже — сам инициирует sync в нашу сторону
// - Протокол: сравнение chain_hash для поиска расхождений (7 типов сообщений)
// - Новые записи немедленно рассылаются подключённым пирам через PUSH
//
// Нюансы:
// - Записи не редактируются и не удаляются явно — только TTL-очистка (per-instance)
// - Дубликаты определяются по (instance_hash, timestamp, datahash)
// - БД хранится в SQLite, путь: <db_path>/sync
// - Синхронизация идёт через etcp_api (service ID 0x20), только прямые P2P соединения
#ifndef DB_SYNC_H #ifndef DB_SYNC_H
#define DB_SYNC_H #define DB_SYNC_H
@ -6,11 +35,11 @@
extern "C" { extern "C" {
#endif #endif
#include <stdint.h> #include <stdint.h>
#include <stddef.h> #include <stddef.h>
struct UTUN_INSTANCE; struct UTUN_INSTANCE;
struct DB_SYNC_INSTANCE;
// etcp_router service ID // etcp_router service ID
#define ETCP_RT_ID_DB_SYNC 0x20 #define ETCP_RT_ID_DB_SYNC 0x20
@ -23,6 +52,11 @@ struct UTUN_INSTANCE;
#define DB_MSG_PUSH 0x05 #define DB_MSG_PUSH 0x05
#define DB_MSG_ACK_PUSH 0x06 #define DB_MSG_ACK_PUSH 0x06
#define DB_MSG_SYNC_DONE 0x07 #define DB_MSG_SYNC_DONE 0x07
#define DB_MSG_ERROR 0x08
// Error codes for DB_MSG_ERROR
#define DB_ERR_NOT_FOUND 0x01 // instance not found
#define DB_ERR_DISABLED 0x02 // instance exists but disabled
// Defaults // Defaults
#define DB_SYNC_DEFAULT_TTL 86400 #define DB_SYNC_DEFAULT_TTL 86400
@ -37,12 +71,36 @@ struct UTUN_INSTANCE;
#define DB_REFINE_HASHES 16 #define DB_REFINE_HASHES 16
#define DB_SEND_DATA_MAX 32 #define DB_SEND_DATA_MAX 32
// Global lifecycle
int db_sync_init(struct UTUN_INSTANCE* inst); int db_sync_init(struct UTUN_INSTANCE* inst);
void db_sync_destroy(struct UTUN_INSTANCE* inst); void db_sync_destroy(struct UTUN_INSTANCE* inst);
int db_sync_insert(struct UTUN_INSTANCE* inst, const char* json_data);
int db_sync_insert_len(struct UTUN_INSTANCE* inst, const char* json_data, size_t len);
uint32_t db_sync_count(struct UTUN_INSTANCE* inst);
// Instance management
struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const char* name, uint64_t id);
void db_sync_instance_remove(struct DB_SYNC_INSTANCE* si);
// Data operations (per-instance)
int db_sync_insert_len(struct DB_SYNC_INSTANCE* si, const char* json_data, size_t len);
int db_sync_insert(struct DB_SYNC_INSTANCE* si, const char* json_data);
int db_sync_insert_signed(struct DB_SYNC_INSTANCE* si, const char* json_data, size_t len,
const uint8_t* sig, size_t sig_len);
uint32_t db_sync_count(struct DB_SYNC_INSTANCE* si);
// Select: iterate records ordered by (timestamp, datahash), starting at offset, max limit (0=unlimited).
// Returns number of records passed to callback.
typedef void (*db_sync_select_cb)(void* arg, uint64_t id, uint64_t timestamp,
const char* data, size_t data_len, uint64_t author,
const uint8_t* author_sig, size_t sig_len,
int delivered_peers, const char* delivery_chain);
int db_sync_select(struct DB_SYNC_INSTANCE* si, uint32_t offset, uint32_t limit,
db_sync_select_cb cb, void* arg);
// Insert callback: fired after local or peer insert succeeds.
// author_node_id = self for local inserts, peer node_id for remote.
typedef void (*db_sync_insert_cb)(struct DB_SYNC_INSTANCE* si,
const char* json_data, size_t len,
uint64_t author_node_id, void* arg);
void db_sync_set_insert_cb(struct DB_SYNC_INSTANCE* si, db_sync_insert_cb cb, void* arg);
#ifdef __cplusplus #ifdef __cplusplus
} }

43
src/topo_group.c

@ -16,15 +16,12 @@
#include "etcp_debug.h" #include "etcp_debug.h"
#include "config_parser.h" #include "config_parser.h"
#include "topo_node.h" #include "topo_node.h"
#include "topo_node_lmdb.h" #include "topo_node_sqlite.h"
#include "route_lib.h" #include "route_lib.h"
#include "topo_group.h" #include "topo_group.h"
#include "route_ping.h" #include "route_ping.h"
#include "route_connectivity.h" #include "route_connectivity.h"
#include "control_server.h" #include "control_server.h"
#ifdef USE_SQLITE
#include "../../tools/chatgui/transport/topo_node_sqlite.h"
#endif
#define TOPO_GROUP_UTUN 0x8000000000000000ULL #define TOPO_GROUP_UTUN 0x8000000000000000ULL
@ -224,9 +221,6 @@ static struct TOPO_GROUP* topo_group_create(struct UTUN_INSTANCE* instance, uint
sc_derive_ed25519_pubkey(instance->my_keys.private_key, group->ed25519_public_key); sc_derive_ed25519_pubkey(instance->my_keys.private_key, group->ed25519_public_key);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "Ed25519 pubkey derived from X25519 private key"); DEBUG_INFO(DEBUG_CATEGORY_BGP, "Ed25519 pubkey derived from X25519 private key");
if (group_type == TOPO_GROUP_TYPE_UTUN && instance->config && instance->config->global.db_path[0])
topo_node_lmdb_init(group, instance->config->global.db_path);
group->nodes = queue_new(instance->ua, BGP_NODES_HASH_SIZE, offsetof(struct TOPO_NODEQ, hash_node_id) - sizeof(struct ll_entry), 8, "group_nodes"); group->nodes = queue_new(instance->ua, BGP_NODES_HASH_SIZE, offsetof(struct TOPO_NODEQ, hash_node_id) - sizeof(struct ll_entry), 8, "group_nodes");
if (!group->nodes) { queue_entry_free(qe); DEBUG_ERROR(DEBUG_CATEGORY_BGP, "nodes queue creation failed"); return NULL; } if (!group->nodes) { queue_entry_free(qe); DEBUG_ERROR(DEBUG_CATEGORY_BGP, "nodes queue creation failed"); return NULL; }
@ -250,8 +244,6 @@ static void topo_group_destroy(struct TOPO_GROUP* group) {
if (group->local_node) { topo_node_free_lists(group, group->local_node); u_free(group->local_node); } if (group->local_node) { topo_node_free_lists(group, group->local_node); u_free(group->local_node); }
if (group->group_type == TOPO_GROUP_TYPE_UTUN) topo_node_lmdb_destroy(group);
if (group->nodes) queue_free(group->nodes); if (group->nodes) queue_free(group->nodes);
} }
@ -273,6 +265,20 @@ struct TOPO_GROUPS* topo_groups_init(struct UTUN_INSTANCE* instance) {
g->v4_subnet_pool = memory_pool_init(sizeof(struct TOPO_SUBNET4), "to_v4sub"); g->v4_subnet_pool = memory_pool_init(sizeof(struct TOPO_SUBNET4), "to_v4sub");
g->v6_subnet_pool = memory_pool_init(sizeof(struct TOPO_SUBNET6), "to_v6sub"); g->v6_subnet_pool = memory_pool_init(sizeof(struct TOPO_SUBNET6), "to_v6sub");
if (instance->config && instance->config->global.db_path[0]) {
int rc = sqlite3_open(instance->config->global.db_path, &g->topo_sqlite_db);
if (rc == SQLITE_OK && g->topo_sqlite_db) {
sqlite3_exec(g->topo_sqlite_db, "PRAGMA journal_mode=WAL", NULL, NULL, NULL);
sqlite3_exec(g->topo_sqlite_db, "PRAGMA foreign_keys=ON", NULL, NULL, NULL);
topo_node_sqlite_init(g->topo_sqlite_db);
DEBUG_INFO(DEBUG_CATEGORY_BGP, "SQLite opened: %s", instance->config->global.db_path);
} else {
DEBUG_ERROR(DEBUG_CATEGORY_BGP, "sqlite3_open(%s) failed: %s",
instance->config->global.db_path, g->topo_sqlite_db ? sqlite3_errmsg(g->topo_sqlite_db) : "null db");
if (g->topo_sqlite_db) { sqlite3_close(g->topo_sqlite_db); g->topo_sqlite_db = NULL; }
}
}
instance->topo_groups = g; instance->topo_groups = g;
struct TOPO_GROUP* default_group = topo_group_create(instance, TOPO_GROUP_UTUN, TOPO_GROUP_TYPE_UTUN); struct TOPO_GROUP* default_group = topo_group_create(instance, TOPO_GROUP_UTUN, TOPO_GROUP_TYPE_UTUN);
@ -313,6 +319,8 @@ void topo_groups_destroy(struct UTUN_INSTANCE* instance) {
memory_pool_destroy(g->v4_subnet_pool); memory_pool_destroy(g->v4_subnet_pool);
memory_pool_destroy(g->v6_subnet_pool); memory_pool_destroy(g->v6_subnet_pool);
if (g->topo_sqlite_db) { sqlite3_close(g->topo_sqlite_db); g->topo_sqlite_db = NULL; }
u_free(g); instance->topo_groups = NULL; u_free(g); instance->topo_groups = NULL;
} }
@ -337,13 +345,11 @@ struct TOPO_GROUP* topo_groups_create_group(struct TOPO_GROUPS* g, uint64_t grou
return group; return group;
} }
#ifdef USE_SQLITE
void topo_groups_set_sqlite_db(struct TOPO_GROUPS* g, sqlite3* db) { void topo_groups_set_sqlite_db(struct TOPO_GROUPS* g, sqlite3* db) {
if (!g) return; if (!g) return;
g->topo_sqlite_db = db; g->topo_sqlite_db = db;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "topo_sqlite_db set"); DEBUG_INFO(DEBUG_CATEGORY_BGP, "topo_sqlite_db set");
} }
#endif
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_group_set_nat_check_local(struct TOPO_GROUP* group, int allow) { if (!group) return; group->allow_nat_check_local = allow ? 1 : 0; }
@ -630,13 +636,12 @@ int topo_group_process_nodeinfo(struct TOPO_GROUP* group, struct ETCP_CONN* from
if (group->group_type != TOPO_GROUP_TYPE_CHAT && group->instance->rt) route_insert(group->instance->rt, nodeinfo1); if (group->group_type != TOPO_GROUP_TYPE_CHAT && group->instance->rt) route_insert(group->instance->rt, nodeinfo1);
if (node_id != group->instance->node_id) { if (node_id != group->instance->node_id) {
#ifdef USE_SQLITE sqlite3* sdb = group->instance->topo_groups->topo_sqlite_db;
if (group->group_type == TOPO_GROUP_TYPE_CHAT && group->instance->topo_groups->topo_sqlite_db) { if (sdb) {
topo_node_sqlite_node_put(group->instance->topo_groups->topo_sqlite_db, nodeinfo1); topo_node_sqlite_node_put(sdb, nodeinfo1);
if (group->channel_id[0]) topo_node_sqlite_member_put(group->instance->topo_groups->topo_sqlite_db, group->channel_id, node_id, NULL, NULL); if (group->group_type == TOPO_GROUP_TYPE_CHAT && group->channel_id[0])
} else topo_node_sqlite_member_put(sdb, group->channel_id, node_id, NULL, NULL);
#endif }
topo_node_lmdb_put(group, nodeinfo1);
} }
if (group->instance->control_srv) control_server_notify_node_change(group->instance->control_srv, nodeinfo1); if (group->instance->control_srv) control_server_notify_node_change(group->instance->control_srv, nodeinfo1);
@ -687,10 +692,8 @@ int topo_group_process_withdraw(struct TOPO_GROUP* group, struct ETCP_CONN* send
struct ll_entry* entry = queue_find_data_by_index(group->nodes, &node_id); struct ll_entry* entry = queue_find_data_by_index(group->nodes, &node_id);
if (entry) { queue_remove_data(group->nodes, entry); queue_entry_free(entry); } if (entry) { queue_remove_data(group->nodes, entry); queue_entry_free(entry); }
topo_group_broadcast_withdraw(group, node_id, wd_source, sender); topo_group_broadcast_withdraw(group, node_id, wd_source, sender);
#ifdef USE_SQLITE
if (group->group_type == TOPO_GROUP_TYPE_CHAT && group->instance->topo_groups->topo_sqlite_db) if (group->group_type == TOPO_GROUP_TYPE_CHAT && group->instance->topo_groups->topo_sqlite_db)
topo_node_sqlite_member_del(group->instance->topo_groups->topo_sqlite_db, group->channel_id, node_id); topo_node_sqlite_member_del(group->instance->topo_groups->topo_sqlite_db, group->channel_id, node_id);
#endif
} }
return 0; return 0;
} }

35
src/topo_group.h

@ -1,3 +1,27 @@
/**
* @file topo_group.h
* @brief BGP-подобный обмен топологией узлов между пирами через ETCP.
*
* Механика:
* - Узлы обмениваются информацией друг о друге (pubkey, адреса, подсети)
* через NODEINFO/WITHDRAW сообщения, маршрутизируемые по ETCP.
* - Каждый узел хранит полную таблицу известных узлов (структура TOPO_NODEQ).
* - При подключении нового пира — полная синхронизация таблицы.
* - При отключении пира — withdraw всех узлов, достижимых только через него.
*
* Группы (TOPO_GROUP):
* - TOPO_GROUP_TYPE_UTUN — основная VPN-сеть (обмен маршрутами, подсетями)
* - TOPO_GROUP_TYPE_CHAT — чат-группа (только узлы, без подсетей)
* - Каждая группа изолирована: узлы из utun-группы не видны в чат-группе и наоборот
*
* Дополнительные функции:
* - NAT-детекция и проверка связности (NAT_INFO, NAT_CHECK_REQ, PING)
* - Поиск оптимального маршрута до узла (topo_group_find_conn_for_node)
*
* Хранение: все узлы сохраняются в SQLite (таблицы nodes, node_addresses,
* channels, peers_<channel_id>). База открывается в topo_groups_init()
* если в конфиге задан db_path.
*/
#ifndef TOPO_GROUP_H #ifndef TOPO_GROUP_H
#define TOPO_GROUP_H #define TOPO_GROUP_H
@ -10,13 +34,10 @@ extern "C" {
#include <stddef.h> #include <stddef.h>
#include "../lib/ll_queue.h" #include "../lib/ll_queue.h"
#include "../lib/memory_pool.h" #include "../lib/memory_pool.h"
#include "../lib/liblmdb/lmdb.h" #include "../lib/sqlite3.h"
#include "route_lib.h" #include "route_lib.h"
#include "topo_node.h" #include "topo_node.h"
#include "secure_channel.h" #include "secure_channel.h"
#ifdef USE_SQLITE
#include <sqlite3.h>
#endif
struct UTUN_INSTANCE; struct UTUN_INSTANCE;
@ -125,8 +146,6 @@ struct TOPO_GROUP {
uint8_t allow_nat_check_local; uint8_t allow_nat_check_local;
uint8_t ed25519_public_key[SC_PUBKEY_SIZE]; uint8_t ed25519_public_key[SC_PUBKEY_SIZE];
char channel_id[64]; // channel_id для групп типа CHAT char channel_id[64]; // channel_id для групп типа CHAT
MDB_env* nodeinfo_env;
MDB_dbi nodeinfo_dbi;
}; };
/** /**
@ -141,9 +160,7 @@ struct TOPO_GROUPS {
struct memory_pool* v6_addr_pool; struct memory_pool* v6_addr_pool;
struct memory_pool* v4_subnet_pool; struct memory_pool* v4_subnet_pool;
struct memory_pool* v6_subnet_pool; struct memory_pool* v6_subnet_pool;
#ifdef USE_SQLITE
sqlite3* topo_sqlite_db; sqlite3* topo_sqlite_db;
#endif
topo_node_updated_fn node_updated_cb; /* optional — set by chatgui's member_sync */ topo_node_updated_fn node_updated_cb; /* optional — set by chatgui's member_sync */
}; };
@ -294,9 +311,7 @@ struct nat_check_arg {
uint16_t nat_port; uint16_t nat_port;
}; };
#ifdef USE_SQLITE
void topo_groups_set_sqlite_db(struct TOPO_GROUPS* g, sqlite3* db); void topo_groups_set_sqlite_db(struct TOPO_GROUPS* g, sqlite3* db);
#endif
#ifdef __cplusplus #ifdef __cplusplus

22
src/topo_node.h

@ -1,4 +1,24 @@
#ifndef TOPO_NODE_H /**
* @file topo_node.h
* @brief Модель данных узла сети — структуры, сериализация, wire-формат.
*
* Этот модуль определяет:
* - Как узел представлен в памяти (TOPO_NODE, TOPO_NODEQ, адреса, подсети)
* - Как узел передаётся по сети (packed-структуры TOPOMSG_*)
* - Как узел сериализуется/десериализуется (topo_node_serialize/deserialize)
*
* Модуль не содержит сетевой логики и BGP-обмена — этим занимается topo_group.
*
* Хранение: данные узлов сохраняются в SQLite (таблицы nodes, node_addresses)
* через модуль topo_node_sqlite.
*
* Основные понятия:
* TOPO_NODE — идентичность узла (pubkeys, имя, версия, списки адресов)
* TOPO_NODEQ — запись в очереди группы (плюс пути, hop_list, связность)
* TOPO_NODEPATH — путь до узла через конкретное ETCP-соединение
* TOPOMSG_NODE — wire-формат, передаваемый по NODEINFO
*/
#ifndef TOPO_NODE_H
#define TOPO_NODE_H #define TOPO_NODE_H
#ifdef __cplusplus #ifdef __cplusplus

66
src/topo_node_lmdb.c

@ -1,66 +0,0 @@
#include <stdlib.h>
#include <string.h>
#include <stdint.h>
#include <sys/stat.h>
#include "topo_node_lmdb.h"
#include "topo_group.h"
#include "../lib/debug_config.h"
#include "../lib/mem.h"
#define NODEINFO_DB_MAPSIZE (10UL * 1024 * 1024)
int topo_node_lmdb_init(struct TOPO_GROUP* bgp, const char* db_path) {
if (!bgp || !db_path || !db_path[0]) return -1;
int rc = mdb_env_create(&bgp->nodeinfo_env);
if (rc != MDB_SUCCESS) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "LMDB mdb_env_create failed: %s", mdb_strerror(rc)); return -1; }
rc = mdb_env_set_mapsize(bgp->nodeinfo_env, NODEINFO_DB_MAPSIZE);
if (rc != MDB_SUCCESS) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "LMDB mdb_env_set_mapsize failed: %s", mdb_strerror(rc)); mdb_env_close(bgp->nodeinfo_env); bgp->nodeinfo_env = NULL; return -1; }
rc = mdb_env_set_maxdbs(bgp->nodeinfo_env, 2);
if (rc != MDB_SUCCESS) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "LMDB mdb_env_set_maxdbs failed: %s", mdb_strerror(rc)); mdb_env_close(bgp->nodeinfo_env); bgp->nodeinfo_env = NULL; return -1; }
rc = mdb_env_open(bgp->nodeinfo_env, db_path, 0, 0644);
if (rc != MDB_SUCCESS) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "LMDB mdb_env_open(%s) failed: %s", db_path, mdb_strerror(rc)); mdb_env_close(bgp->nodeinfo_env); bgp->nodeinfo_env = NULL; return -1; }
MDB_txn* txn = NULL; rc = mdb_txn_begin(bgp->nodeinfo_env, NULL, 0, &txn);
if (rc != MDB_SUCCESS) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "LMDB mdb_txn_begin failed: %s", mdb_strerror(rc)); mdb_env_close(bgp->nodeinfo_env); bgp->nodeinfo_env = NULL; return -1; }
rc = mdb_dbi_open(txn, "nodeinfo", MDB_CREATE, &bgp->nodeinfo_dbi);
if (rc != MDB_SUCCESS) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "LMDB mdb_dbi_open(nodeinfo) failed: %s", mdb_strerror(rc)); mdb_txn_abort(txn); mdb_env_close(bgp->nodeinfo_env); bgp->nodeinfo_env = NULL; return -1; }
rc = mdb_txn_commit(txn);
if (rc != MDB_SUCCESS) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "LMDB mdb_txn_commit failed: %s", mdb_strerror(rc)); mdb_env_close(bgp->nodeinfo_env); bgp->nodeinfo_env = NULL; return -1; }
DEBUG_INFO(DEBUG_CATEGORY_BGP, "LMDB nodeinfo database opened: %s", db_path);
return 0;
}
void topo_node_lmdb_destroy(struct TOPO_GROUP* bgp) {
if (!bgp) return;
if (bgp->nodeinfo_env) { mdb_env_close(bgp->nodeinfo_env); bgp->nodeinfo_env = NULL; DEBUG_INFO(DEBUG_CATEGORY_BGP, "LMDB nodeinfo database closed"); }
}
int topo_node_lmdb_put(struct TOPO_GROUP* bgp, struct TOPO_NODEQ* nq) {
if (!bgp || !bgp->nodeinfo_env || !nq || !nq->node) return -1;
uint8_t buf[8192];
int ser_len = topo_node_serialize(bgp, nq, buf, sizeof(buf));
if (ser_len < 0) return -1;
MDB_val key, data;
key.mv_size = 8; key.mv_data = &nq->node->node_id;
data.mv_size = (size_t)ser_len; data.mv_data = buf;
MDB_txn* txn = NULL; int rc = mdb_txn_begin(bgp->nodeinfo_env, NULL, 0, &txn);
if (rc != MDB_SUCCESS) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "LMDB put mdb_txn_begin failed: %s", mdb_strerror(rc)); return -1; }
rc = mdb_put(txn, bgp->nodeinfo_dbi, &key, &data, 0);
if (rc != MDB_SUCCESS) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "LMDB mdb_put node %016llx failed: %s", (unsigned long long)nq->node->node_id, mdb_strerror(rc)); mdb_txn_abort(txn); return -1; }
rc = mdb_txn_commit(txn);
if (rc != MDB_SUCCESS) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "LMDB put mdb_txn_commit failed: %s", mdb_strerror(rc)); return -1; }
return 0;
}
int topo_node_lmdb_del(struct TOPO_GROUP* bgp, uint64_t node_id) {
if (!bgp || !bgp->nodeinfo_env) return -1;
MDB_val key; key.mv_size = 8; key.mv_data = &node_id;
MDB_txn* txn = NULL; int rc = mdb_txn_begin(bgp->nodeinfo_env, NULL, 0, &txn);
if (rc != MDB_SUCCESS) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "LMDB del mdb_txn_begin failed: %s", mdb_strerror(rc)); return -1; }
rc = mdb_del(txn, bgp->nodeinfo_dbi, &key, NULL);
if (rc != MDB_SUCCESS && rc != MDB_NOTFOUND) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "LMDB mdb_del node %016llx failed: %s", (unsigned long long)node_id, mdb_strerror(rc)); mdb_txn_abort(txn); return -1; }
rc = mdb_txn_commit(txn);
if (rc != MDB_SUCCESS) { DEBUG_ERROR(DEBUG_CATEGORY_BGP, "LMDB del mdb_txn_commit failed: %s", mdb_strerror(rc)); return -1; }
return 0;
}

23
src/topo_node_lmdb.h

@ -1,23 +0,0 @@
#ifndef TOPO_NODE_LMDB_H
#define TOPO_NODE_LMDB_H
#ifdef __cplusplus
extern "C" {
#endif
#include <stdint.h>
#include "topo_node.h"
struct TOPO_GROUP;
int topo_node_lmdb_init(struct TOPO_GROUP* group, const char* db_path);
void topo_node_lmdb_destroy(struct TOPO_GROUP* group);
int topo_node_lmdb_put(struct TOPO_GROUP* group, struct TOPO_NODEQ* nq);
int topo_node_lmdb_del(struct TOPO_GROUP* group, uint64_t node_id);
#ifdef __cplusplus
}
#endif
#endif

4
tools/chatgui/transport/topo_node_sqlite.c → src/topo_node_sqlite.c

@ -1,6 +1,6 @@
#include "topo_node_sqlite.h" #include "topo_node_sqlite.h"
#include "../../../src/topo_node.h" #include "topo_node.h"
#include "../../../lib/debug_config.h" #include "../lib/debug_config.h"
#include <string.h> #include <string.h>
#include <stdio.h> #include <stdio.h>

0
tools/chatgui/transport/topo_node_sqlite.h → src/topo_node_sqlite.h

60
tests/test_db_sync.c

@ -36,6 +36,8 @@
static struct UTUN_INSTANCE* inst_a = NULL; static struct UTUN_INSTANCE* inst_a = NULL;
static struct UTUN_INSTANCE* inst_b = NULL; static struct UTUN_INSTANCE* inst_b = NULL;
static struct DB_SYNC_INSTANCE* si_a = NULL;
static struct DB_SYNC_INSTANCE* si_b = NULL;
static struct UASYNC* ua = NULL; static struct UASYNC* ua = NULL;
static int test_phase = 0; // 0=running, 1=success, 2=failure static int test_phase = 0; // 0=running, 1=success, 2=failure
static void* timeout_id = NULL; static void* timeout_id = NULL;
@ -65,12 +67,10 @@ static int create_temp_configs(void) {
snprintf(config_a, sizeof(config_a), "%s/a.conf", temp_dir); snprintf(config_a, sizeof(config_a), "%s/a.conf", temp_dir);
snprintf(config_b, sizeof(config_b), "%s/b.conf", temp_dir); snprintf(config_b, sizeof(config_b), "%s/b.conf", temp_dir);
// Create LMDB directories (parent first, then child) // Create SQLite db directories (parent dir for sync file)
char db_path[320]; char db_path[320];
snprintf(db_path, sizeof(db_path), "%s/db_a", temp_dir); utun_mkdir(db_path, 0755); snprintf(db_path, sizeof(db_path), "%s/db_a", temp_dir); utun_mkdir(db_path, 0755);
snprintf(db_path, sizeof(db_path), "%s/db_a/sync", temp_dir); utun_mkdir(db_path, 0755);
snprintf(db_path, sizeof(db_path), "%s/db_b", temp_dir); utun_mkdir(db_path, 0755); snprintf(db_path, sizeof(db_path), "%s/db_b", temp_dir); utun_mkdir(db_path, 0755);
snprintf(db_path, sizeof(db_path), "%s/db_b/sync", temp_dir); utun_mkdir(db_path, 0755);
if (write_file(config_a, if (write_file(config_a,
"[global]\n" "[global]\n"
@ -120,16 +120,14 @@ static int create_temp_configs(void) {
static void cleanup_temp_configs(void) { static void cleanup_temp_configs(void) {
unlink(config_a); unlink(config_b); unlink(config_a); unlink(config_b);
char db_a[320], db_b[320]; char pa[320];
snprintf(db_a, sizeof(db_a), "%s/db_a/sync/data.mdb", temp_dir); snprintf(pa, sizeof(pa), "%s/db_a/sync", temp_dir); unlink(pa);
snprintf(db_b, sizeof(db_b), "%s/db_b/sync/data.mdb", temp_dir); snprintf(pa, sizeof(pa), "%s/db_a/sync-wal", temp_dir); unlink(pa);
unlink(db_a); unlink(db_b); snprintf(pa, sizeof(pa), "%s/db_a/sync-shm", temp_dir); unlink(pa);
snprintf(db_a, sizeof(db_a), "%s/db_a/sync/lock.mdb", temp_dir); snprintf(pa, sizeof(pa), "%s/db_b/sync", temp_dir); unlink(pa);
snprintf(db_b, sizeof(db_b), "%s/db_b/sync/lock.mdb", temp_dir); snprintf(pa, sizeof(pa), "%s/db_b/sync-wal", temp_dir); unlink(pa);
unlink(db_a); unlink(db_b); snprintf(pa, sizeof(pa), "%s/db_b/sync-shm", temp_dir); unlink(pa);
char pa[320]; snprintf(pa, sizeof(pa), "%s/db_a/sync", temp_dir); test_rmdir(pa);
snprintf(pa, sizeof(pa), "%s/db_a", temp_dir); test_rmdir(pa); snprintf(pa, sizeof(pa), "%s/db_a", temp_dir); test_rmdir(pa);
snprintf(pa, sizeof(pa), "%s/db_b/sync", temp_dir); test_rmdir(pa);
snprintf(pa, sizeof(pa), "%s/db_b", temp_dir); test_rmdir(pa); snprintf(pa, sizeof(pa), "%s/db_b", temp_dir); test_rmdir(pa);
test_rmdir(temp_dir); test_rmdir(temp_dir);
} }
@ -159,33 +157,33 @@ static int cond_links_init(void) {
} }
static int cond_count_a(uint32_t expected) { static int cond_count_a(uint32_t expected) {
if (!inst_a) return 0; if (!si_a) return 0;
return db_sync_count(inst_a) == expected; return db_sync_count(si_a) == expected;
} }
static int cond_count_b(uint32_t expected) { static int cond_count_b(uint32_t expected) {
if (!inst_b) return 0; if (!si_b) return 0;
return db_sync_count(inst_b) == expected; return db_sync_count(si_b) == expected;
} }
static uint32_t count_a_target, count_b_target; static uint32_t count_a_target, count_b_target;
static int _cond_count_a(void) { return cond_count_a(count_a_target); } static int _cond_count_a(void) { return cond_count_a(count_a_target); }
static int _cond_count_b(void) { return cond_count_b(count_b_target); } static int _cond_count_b(void) { return cond_count_b(count_b_target); }
static int insert_many(struct UTUN_INSTANCE* inst, int start, int count) { static int insert_many(struct DB_SYNC_INSTANCE* si, int start, int count) {
char buf[128]; char buf[128];
for (int i = start; i < start + count && test_phase == 0; i++) { for (int i = start; i < start + count && test_phase == 0; i++) {
snprintf(buf, sizeof(buf), "{\"idx\":%d,\"val\":\"data_%d\",\"pad\":\"%s\"}", i, i, snprintf(buf, sizeof(buf), "{\"idx\":%d,\"val\":\"data_%d\",\"pad\":\"%s\"}", i, i,
"xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx"); "xxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxxx");
int ret = db_sync_insert(inst, buf); int ret = db_sync_insert(si, buf);
if (ret < 0) { fprintf(stderr, "insert_many failed at %d ret=%d\n", i, ret); return -1; } if (ret < 0) { fprintf(stderr, "insert_many failed at %d ret=%d\n", i, ret); return -1; }
} }
return 0; return 0;
} }
static int insert_many_batch(struct UTUN_INSTANCE* inst, int start, int count) { static int insert_many_batch(struct DB_SYNC_INSTANCE* si, int start, int count) {
char buf[256]; char buf[256];
for (int i = start; i < start + count && test_phase == 0; i++) { for (int i = start; i < start + count && test_phase == 0; i++) {
snprintf(buf, sizeof(buf), "{\"n\":%d,\"text\":\"record_number_%d_abcdefghijklmnopqrstuvwxyz\"}", i, i); snprintf(buf, sizeof(buf), "{\"n\":%d,\"text\":\"record_number_%d_abcdefghijklmnopqrstuvwxyz\"}", i, i);
db_sync_insert(inst, buf); db_sync_insert(si, buf);
} }
return 0; return 0;
} }
@ -211,28 +209,34 @@ int main(void) {
fprintf(stderr, "instance init failed\n"); cleanup_temp_configs(); return 1; fprintf(stderr, "instance init failed\n"); cleanup_temp_configs(); return 1;
} }
si_a = db_sync_instance_add(inst_a, "test", 1);
si_b = db_sync_instance_add(inst_b, "test", 1);
if (!si_a || !si_b) { fprintf(stderr, "db_sync_instance_add failed\n"); return 1; }
timeout_id = uasync_set_timeout(ua, TEST_TIMEOUT_TB, NULL, test_timeout, "test_timeout"); timeout_id = uasync_set_timeout(ua, TEST_TIMEOUT_TB, NULL, test_timeout, "test_timeout");
// ===== Phase 1: базовый CRUD ===== // ===== Phase 1: базовый CRUD =====
printf("Phase 1: basic CRUD...\n"); printf("Phase 1: basic CRUD...\n");
if (db_sync_count(inst_a) != 0) { fprintf(stderr, "FAIL: initial count not 0 (got %u)\n", db_sync_count(inst_a)); test_phase=2; } if (db_sync_count(si_a) != 0) { fprintf(stderr, "FAIL: initial count not 0 (got %u)\n", db_sync_count(si_a)); test_phase=2; }
if (db_sync_insert(inst_a, "{\"key\":\"val1\"}") != 0) { fprintf(stderr,"FAIL: insert 1\n"); test_phase=2; } if (db_sync_insert(si_a, "{\"key\":\"val1\"}") != 0) { fprintf(stderr,"FAIL: insert 1\n"); test_phase=2; }
if (db_sync_insert(inst_a, "{\"key\":\"val2\"}") != 0) { fprintf(stderr,"FAIL: insert 2\n"); test_phase=2; } if (db_sync_insert(si_a, "{\"key\":\"val2\"}") != 0) { fprintf(stderr,"FAIL: insert 2\n"); test_phase=2; }
if (db_sync_insert(inst_a, "{\"key\":\"val3\"}") != 0) { fprintf(stderr,"FAIL: insert 3\n"); test_phase=2; } if (db_sync_insert(si_a, "{\"key\":\"val3\"}") != 0) { fprintf(stderr,"FAIL: insert 3\n"); test_phase=2; }
if (db_sync_count(inst_a) != 3) { fprintf(stderr,"FAIL: count not 3 (got %u)\n", db_sync_count(inst_a)); test_phase=2; } if (db_sync_count(si_a) != 3) { fprintf(stderr,"FAIL: count not 3 (got %u)\n", db_sync_count(si_a)); test_phase=2; }
// Dedup by (timestamp,datahash): same content at different time = new record // Dedup by (timestamp,datahash): same content at different time = new record
if (db_sync_insert(inst_a, "{\"key\":\"val4\"}") != 0) { fprintf(stderr,"FAIL: insert 4\n"); test_phase=2; } if (db_sync_insert(si_a, "{\"key\":\"val4\"}") != 0) { fprintf(stderr,"FAIL: insert 4\n"); test_phase=2; }
if (db_sync_count(inst_a) != 4) { fprintf(stderr,"FAIL: count not 4 (got %u)\n", db_sync_count(inst_a)); test_phase=2; } if (db_sync_count(si_a) != 4) { fprintf(stderr,"FAIL: count not 4 (got %u)\n", db_sync_count(si_a)); test_phase=2; }
if (test_phase == 0) printf("Phase 1: PASS (count=4)\n"); if (test_phase == 0) printf("Phase 1: PASS (count=4)\n");
// ===== Phase 2: initial sync A↔B ===== // ===== Phase 2: initial sync A↔B =====
printf("Phase 2: initial sync...\n"); printf("Phase 2: initial sync...\n");
if (test_phase == 0) { count_b_target = 4; if (!wait_for("B count=4", _cond_count_b, PHASE_TIMEOUT_TB)) test_phase=2; } if (test_phase == 0) { count_b_target = 4; if (!wait_for("B count=4", _cond_count_b, PHASE_TIMEOUT_TB)) test_phase=2; }
if (test_phase == 0) printf("Phase 2: PASS (B synced %u records)\n", db_sync_count(inst_b)); if (test_phase == 0) printf("Phase 2: PASS (B synced %u records)\n", db_sync_count(si_b));
// ===== Cleanup ===== // ===== Cleanup =====
if (timeout_id && ua) { uasync_cancel_timeout(ua, timeout_id); timeout_id = NULL; } if (timeout_id && ua) { uasync_cancel_timeout(ua, timeout_id); timeout_id = NULL; }
if (si_a) { db_sync_instance_remove(si_a); si_a = NULL; }
if (si_b) { db_sync_instance_remove(si_b); si_b = NULL; }
if (inst_a) { inst_a->running = 0; utun_instance_destroy(inst_a); inst_a = NULL; } if (inst_a) { inst_a->running = 0; utun_instance_destroy(inst_a); inst_a = NULL; }
if (inst_b) { inst_b->running = 0; utun_instance_destroy(inst_b); inst_b = NULL; } if (inst_b) { inst_b->running = 0; utun_instance_destroy(inst_b); inst_b = NULL; }
if (ua) { uasync_destroy(ua, 0); ua = NULL; } if (ua) { uasync_destroy(ua, 0); ua = NULL; }

16
tools/chatgui/CMakeLists.txt

@ -75,16 +75,14 @@ add_executable(chatgui
transport/chat_core.c transport/chat_core.c
transport/chat_sync.c transport/chat_sync.c
transport/member_sync.c transport/member_sync.c
transport/db_sync_stub.c
transport/topo_node_sqlite.c
db/db_manager.cpp db/db_manager.cpp
db/sqlite3.c ../../lib/sqlite3.c
resources/chatgui.qrc resources/chatgui.qrc
) )
target_include_directories(chatgui PRIVATE ${CMAKE_SOURCE_DIR}/db) target_include_directories(chatgui PRIVATE ${CMAKE_SOURCE_DIR}/../../lib ${CMAKE_SOURCE_DIR}/../../src ${CMAKE_SOURCE_DIR}/db)
target_compile_definitions(chatgui PRIVATE SQLITE_THREADSAFE=1 USE_SQLITE) target_compile_definitions(chatgui PRIVATE SQLITE_THREADSAFE=1 USE_SQLITE)
set_source_files_properties(db/sqlite3.c transport/chat_core.c transport/chat_sync.c transport/member_sync.c transport/db_sync_stub.c PROPERTIES LANGUAGE C) set_source_files_properties(../../lib/sqlite3.c transport/chat_core.c transport/chat_sync.c transport/member_sync.c PROPERTIES LANGUAGE C)
if(WIN32) if(WIN32)
target_link_libraries(chatgui PRIVATE ${QT_LIBS} rlottie::rlottie ZLIB::ZLIB OpenSSL::Crypto ZXing::ZXing utun pthread) target_link_libraries(chatgui PRIVATE ${QT_LIBS} rlottie::rlottie ZLIB::ZLIB OpenSSL::Crypto ZXing::ZXing utun pthread)
else() else()
@ -95,16 +93,16 @@ endif()
# Tests # Tests
# ============================================================================ # ============================================================================
add_executable(test_member_sync add_executable(test_member_sync
db/sqlite3.c ../../lib/sqlite3.c
../../tests/test_member_sync.c ../../tests/test_member_sync.c
) )
target_include_directories(test_member_sync PRIVATE target_include_directories(test_member_sync PRIVATE
transport ${CMAKE_SOURCE_DIR} ${CMAKE_SOURCE_DIR}/db transport ${CMAKE_SOURCE_DIR} ${CMAKE_SOURCE_DIR}/../../lib
${CMAKE_SOURCE_DIR}/../../src ${CMAKE_SOURCE_DIR}/../../lib ${CMAKE_SOURCE_DIR}/../../src
zxing-cpp/core/src ../../tests zxing-cpp/core/src ../../tests
) )
target_compile_definitions(test_member_sync PRIVATE USE_SQLITE SQLITE_THREADSAFE=1) target_compile_definitions(test_member_sync PRIVATE USE_SQLITE SQLITE_THREADSAFE=1)
set_source_files_properties(db/sqlite3.c ../../tests/test_member_sync.c PROPERTIES LANGUAGE C) set_source_files_properties(../../lib/sqlite3.c ../../tests/test_member_sync.c PROPERTIES LANGUAGE C)
if(WIN32) if(WIN32)
target_link_libraries(test_member_sync PRIVATE OpenSSL::Crypto pthread) target_link_libraries(test_member_sync PRIVATE OpenSSL::Crypto pthread)
else() else()

5
tools/chatgui/libutun/CMakeLists.txt

@ -23,6 +23,7 @@ set(UASYNC_SOURCES
${LIB_DIR}/radix.c ${LIB_DIR}/radix.c
${LIB_DIR}/serialize.c ${LIB_DIR}/serialize.c
${LIB_DIR}/tcp_io.c ${LIB_DIR}/tcp_io.c
${LIB_DIR}/sqlite3.c
${LIB_DIR}/liblmdb/mdb.c ${LIB_DIR}/liblmdb/mdb.c
${LIB_DIR}/liblmdb/midl.c ${LIB_DIR}/liblmdb/midl.c
) )
@ -59,10 +60,10 @@ set(UTUN_COMMON_SOURCES
${SRC_DIR}/topo_group.c ${SRC_DIR}/topo_group.c
${SRC_DIR}/route_ping.c ${SRC_DIR}/route_ping.c
${SRC_DIR}/topo_node.c ${SRC_DIR}/topo_node.c
${SRC_DIR}/topo_node_lmdb.c ${SRC_DIR}/topo_node_sqlite.c
${SRC_DIR}/route_connectivity.c ${SRC_DIR}/route_connectivity.c
${SRC_DIR}/db_sync.c
${SRC_DIR}/conn_mgr.c ${SRC_DIR}/conn_mgr.c
${TRANSPORT_DIR}/db_sync_stub.c
${TRANSPORT_DIR}/chat_sync.c ${TRANSPORT_DIR}/chat_sync.c
${TRANSPORT_DIR}/member_sync.c ${TRANSPORT_DIR}/member_sync.c
${SRC_DIR}/routing.c ${SRC_DIR}/routing.c

215
tools/chatgui/transport/chat_core.c

@ -6,7 +6,7 @@
*/ */
#include "chat_core.h" #include "chat_core.h"
#include "chat_sync.h" #include "db_sync.h"
#include "gui_bridge.h" #include "gui_bridge.h"
#include "topo_node_sqlite.h" #include "topo_node_sqlite.h"
#include "member_sync.h" #include "member_sync.h"
@ -41,12 +41,21 @@ static struct chat_core_ctx {
sqlite3* db; sqlite3* db;
uint64_t my_node_id; uint64_t my_node_id;
/* db_sync instances per channel */
struct DB_SYNC_INSTANCE** si;
char** si_ch_id;
int si_count, si_capacity;
/* курсоры (макс 16) */ /* курсоры (макс 16) */
sqlite3_stmt* cursors[16]; sqlite3_stmt* cursors[16];
uint32_t next_cursor_id; uint32_t next_cursor_id;
uint8_t initialized; uint8_t initialized;
} g_cc; } g_cc;
static struct DB_SYNC_INSTANCE* si_find(const char* ch_id);
static void si_register(struct DB_SYNC_INSTANCE* si, const char* ch_id);
static void on_db_sync_insert(struct DB_SYNC_INSTANCE* si, const char* data, size_t len, uint64_t author, void* arg);
/* ─── утилиты ─── */ /* ─── утилиты ─── */
static void sanitize_ch_id(const char* ch_id, char* out, size_t out_sz) { static void sanitize_ch_id(const char* ch_id, char* out, size_t out_sz) {
@ -286,96 +295,46 @@ void chat_core_update_my_name_trampoline(void* arg) {
/* ─── отправка сообщения (GUI → uasync) ─── */ /* ─── отправка сообщения (GUI → uasync) ─── */
void chat_core_submit_message(struct chat_msg_submit* req) { void chat_core_submit_message(struct chat_msg_submit* req) {
if (!g_cc.initialized) { if (!g_cc.initialized) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: submit before init", CC_ID); return; }
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: submit_message called before init", CC_ID);
return;
}
if (!req) return; if (!req) return;
const char* ch_id = req->channel_id; struct DB_SYNC_INSTANCE* si = si_find(req->channel_id);
const uint8_t* data = req->data; if (!si) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: no db_sync instance for ch=%s", CC_ID, req->channel_id); return; }
uint32_t data_len = req->data_len;
uint64_t ts = req->timestamp; /* Build JSON: {"n":<node_id>,"ch":"<ch_id>","ct":"<>","d":"<>"} */
uint64_t dh; char json[4096];
compute_datahash(data, data_len, &dh); snprintf(json, sizeof(json),
"{\"n\":%llu,\"ch\":\"%s\",\"ct\":\"%s\",\"d\":\"%.*s\"}",
uint8_t prev_ch[32]; get_last_chain_hash(ch_id, prev_ch); (unsigned long long)g_cc.my_node_id, req->channel_id,
uint8_t chain_hash[32]; req->content_type, (int)req->data_len, (const char*)req->data);
compute_chain_hash(prev_ch, (int64_t)ts, dh, chain_hash); /* FIXME: JSON escaping for d (quotes/backslashes in data) */
/* For now, naive format; will break on special chars. TODO: use base64. */
char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl));
char sql[256]; snprintf(sql, sizeof(sql), int ret = db_sync_insert_signed(si, json, strlen(json), NULL, 0);
"INSERT OR IGNORE INTO \"%s\"" if (ret != 0) {
" (node_id, content_type, data, timestamp, datahash, chain_hash," DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: db_sync_insert_signed failed ret=%d", CC_ID, ret);
" signature, is_outgoing, is_read)" /* Insert directly to msg_ table as fallback (dual-write not triggered by callback) */
" VALUES(?,?,?,?,?,?,?,1,1)", tbl); uint64_t dh; compute_datahash((const uint8_t*)req->data, req->data_len, &dh);
char tbl[80]; msg_table_name(req->channel_id, tbl, sizeof(tbl));
sqlite3_stmt* stmt = NULL; char sql[512]; snprintf(sql, sizeof(sql), "INSERT OR IGNORE INTO \"%s\" (node_id,content_type,data,timestamp,datahash,chain_hash,signature,is_outgoing,is_read) VALUES(?,?,?,?,?,?,?,1,1)", tbl);
if (sqlite3_prepare_v2(g_cc.db, sql, -1, &stmt, NULL) != SQLITE_OK) { sqlite3_stmt* st=NULL;
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: prepare insert failed: %s", if (sqlite3_prepare_v2(g_cc.db,sql,-1,&st,NULL)==SQLITE_OK) {
CC_ID, sqlite3_errmsg(g_cc.db)); sqlite3_bind_int64(st,1,(sqlite3_int64)g_cc.my_node_id);
return; sqlite3_bind_text(st,2,req->content_type,-1,SQLITE_STATIC);
} sqlite3_bind_blob(st,3,req->data,(int)req->data_len,SQLITE_STATIC);
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)g_cc.my_node_id); sqlite3_bind_int64(st,4,(sqlite3_int64)req->timestamp);
sqlite3_bind_text(stmt, 2, req->content_type, -1, SQLITE_STATIC); sqlite3_bind_int64(st,5,(sqlite3_int64)dh);
sqlite3_bind_blob(stmt, 3, data, (int)data_len, SQLITE_STATIC); static const uint8_t ch32[32]={0}; sqlite3_bind_blob(st,6,ch32,32,SQLITE_STATIC);
sqlite3_bind_int64(stmt, 4, (sqlite3_int64)ts); static const uint8_t sig64[64]={0}; sqlite3_bind_blob(st,7,sig64,64,SQLITE_STATIC);
sqlite3_bind_int64(stmt, 5, (sqlite3_int64)dh); sqlite3_step(st); sqlite3_finalize(st);
sqlite3_bind_blob(stmt, 6, chain_hash, 32, SQLITE_STATIC);
{
static const char zero64[64] = {0};
sqlite3_bind_blob(stmt, 7, zero64, 64, SQLITE_STATIC);
}
int rc = sqlite3_step(stmt);
sqlite3_finalize(stmt);
if (rc == SQLITE_DONE) {
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);
} else if (rc == SQLITE_CONSTRAINT) {
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: msg duplicate ch=%s ts=%llu",
CC_ID, ch_id, (unsigned long long)ts);
} else {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: insert msg failed: %s",
CC_ID, sqlite3_errmsg(g_cc.db));
} }
} /* notify GUI anyway */
uint8_t evt[65]; uint8_t cl=(uint8_t)strlen(req->channel_id); evt[0]=cl; memcpy(evt+1,req->channel_id,cl);
void chat_core_submit_trampoline(void* arg) { gui_bridge_post(GUI_EVT_MSG_RECEIVED, evt, 1+cl);
chat_core_submit_message((struct chat_msg_submit*)arg);
u_free(arg);
}
void chat_core_push_message(struct chat_msg_submit* req) {
if (!g_cc.initialized) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: push_message called before init", CC_ID);
return;
} }
if (!req) return;
uint64_t dh;
compute_datahash(req->data, req->data_len, &dh);
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: pushing msg ch=%s ts=%llu dh=0x%016llx",
CC_ID, req->channel_id, (unsigned long long)req->timestamp, (unsigned long long)dh);
chat_sync_push(g_cc.inst, req->channel_id, g_cc.my_node_id,
req->content_type, req->data, req->data_len,
req->timestamp, dh);
} }
void chat_core_push_trampoline(void* arg) { void chat_core_submit_trampoline(void* arg) { chat_core_submit_message((struct chat_msg_submit*)arg); u_free(arg); }
chat_core_push_message((struct chat_msg_submit*)arg);
u_free(arg);
}
/* ─── DB-операции для chat_sync ─── */ /* ─── DB-операции для chat_sync ─── */
@ -538,35 +497,6 @@ void chat_core_cursor_close(uint32_t cursor_id) {
} }
} }
void chat_core_mark_sent(const char* ch_id, uint64_t ts, uint64_t dh,
uint64_t node_id) {
if (!g_cc.initialized) return;
char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl));
char sql[256]; snprintf(sql, sizeof(sql),
"UPDATE \"%s\" SET sync_flags=sync_flags|1"
" WHERE timestamp=? AND datahash=? AND node_id=? AND is_outgoing=1", tbl);
sqlite3_stmt* stmt = NULL;
if (sqlite3_prepare_v2(g_cc.db, sql, -1, &stmt, NULL) != SQLITE_OK) return;
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)ts);
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)dh);
sqlite3_bind_int64(stmt, 3, (sqlite3_int64)node_id);
sqlite3_step(stmt); sqlite3_finalize(stmt);
}
void chat_core_ttl_delete(const char* ch_id, uint64_t node_id,
uint64_t cutoff_us) {
if (!g_cc.initialized) return;
char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl));
char sql[256]; snprintf(sql, sizeof(sql),
"DELETE FROM \"%s\" WHERE node_id=? AND is_outgoing=1 AND (sync_flags&1)=0 AND timestamp<?",
tbl);
sqlite3_stmt* stmt = NULL;
if (sqlite3_prepare_v2(g_cc.db, sql, -1, &stmt, NULL) != SQLITE_OK) return;
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)node_id);
sqlite3_bind_int64(stmt, 2, (sqlite3_int64)cutoff_us);
sqlite3_step(stmt); sqlite3_finalize(stmt);
}
int chat_core_list_channels(uint8_t* buf, size_t buf_size, size_t* out_len) { int chat_core_list_channels(uint8_t* buf, size_t buf_size, size_t* out_len) {
if (!g_cc.initialized || !buf || !out_len) return -1; if (!g_cc.initialized || !buf || !out_len) return -1;
sqlite3_stmt* stmt = NULL; sqlite3_stmt* stmt = NULL;
@ -955,6 +885,21 @@ void chat_core_create_channel(struct chat_channel_create* req) {
tbl_msg, tbl_msg); tbl_msg, tbl_msg);
db_exec(sql); db_exec(sql);
/* Register with db_sync for message sync */
{
/* Compute a deterministic hash from channel_id for db_sync instance id */
uint64_t ch_hash = 0;
{ const uint8_t* chd = (const uint8_t*)req->channel_id; size_t chl = strlen(req->channel_id); uint8_t sh[32]; SHA256(chd, chl, sh); memcpy(&ch_hash, sh, 8); }
struct DB_SYNC_INSTANCE* si = db_sync_instance_add(g_cc.inst, "chats", ch_hash);
if (si) {
si_register(si, req->channel_id);
db_sync_set_insert_cb(si, on_db_sync_insert, u_strdup(req->channel_id));
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: db_sync instance added for ch=%s hash=0x%016llx", CC_ID, req->channel_id, (unsigned long long)ch_hash);
} else {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: db_sync_instance_add failed for ch=%s", CC_ID, req->channel_id);
}
}
/* записываем канал в БД */ /* записываем канал в БД */
int rc = topo_node_sqlite_channel_put(g_cc.db, int rc = topo_node_sqlite_channel_put(g_cc.db,
req->channel_id, req->name, req->is_dm, req->owner_node_id, req->channel_id, req->name, req->is_dm, req->owner_node_id,
@ -1015,3 +960,45 @@ void chat_core_create_channel_trampoline(void* arg) {
chat_core_create_channel(req); chat_core_create_channel(req);
u_free(req); u_free(req);
} }
/* ─── db_sync helpers ─── */
static struct DB_SYNC_INSTANCE* si_find(const char* ch_id) {
for (int i=0; i<g_cc.si_count; i++) if (strcmp(g_cc.si_ch_id[i], ch_id)==0) return g_cc.si[i];
return NULL;
}
static void si_register(struct DB_SYNC_INSTANCE* si, const char* ch_id) {
if (g_cc.si_count>=g_cc.si_capacity) { int nc=g_cc.si_capacity?g_cc.si_capacity*2:8; g_cc.si=u_realloc(g_cc.si,nc*sizeof(void*)); g_cc.si_ch_id=u_realloc(g_cc.si_ch_id,nc*sizeof(char*)); g_cc.si_capacity=nc; }
g_cc.si[g_cc.si_count]=si; g_cc.si_ch_id[g_cc.si_count]=u_strdup(ch_id); g_cc.si_count++;
}
static void on_db_sync_insert(struct DB_SYNC_INSTANCE* si, const char* data, size_t len, uint64_t author, void* arg) {
const char* ch_id = (const char*)arg;
uint64_t dh; compute_datahash((const uint8_t*)data, len, &dh);
uint64_t ts = get_time_us();
char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl));
char sql[512]; snprintf(sql, sizeof(sql),
"INSERT OR IGNORE INTO \"%s\" (node_id,content_type,data,timestamp,datahash,chain_hash,signature,is_outgoing,is_read,sync_flags)"
" VALUES(?,?,?,?,?,?,?,?,?,?)", tbl);
sqlite3_stmt* st=NULL;
if (sqlite3_prepare_v2(g_cc.db,sql,-1,&st,NULL)!=SQLITE_OK) return;
sqlite3_bind_int64(st,1,(sqlite3_int64)author);
sqlite3_bind_text(st,2,"text/plain",-1,SQLITE_STATIC);
sqlite3_bind_blob(st,3,data,(int)len,SQLITE_STATIC);
sqlite3_bind_int64(st,4,(sqlite3_int64)ts);
sqlite3_bind_int64(st,5,(sqlite3_int64)dh);
{ static const uint8_t z32[32]={0}; sqlite3_bind_blob(st,6,z32,32,SQLITE_STATIC); }
{ static const uint8_t z64[64]={0}; sqlite3_bind_blob(st,7,z64,64,SQLITE_STATIC); }
sqlite3_bind_int(st,8,author==g_cc.my_node_id?1:0);
sqlite3_bind_int(st,9,1); sqlite3_bind_int(st,10,0);
int rc=sqlite3_step(st); sqlite3_finalize(st);
if (rc==SQLITE_DONE) {
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);
}
}
/* ─── stubs for chat_sync compatibility (db_sync handles this now) ─── */
void chat_core_mark_sent(const char* ch_id, uint64_t ts, uint64_t dh, uint64_t peer_id) { (void)ch_id; (void)ts; (void)dh; (void)peer_id; }
void chat_core_ttl_delete(const char* ch_id, uint64_t node_id, uint64_t cutoff_us) { (void)ch_id; (void)node_id; (void)cutoff_us; }

16
tools/chatgui/transport/chat_core.h

@ -31,20 +31,13 @@ struct chat_msg_submit {
void chat_core_submit_message(struct chat_msg_submit* req); void chat_core_submit_message(struct chat_msg_submit* req);
/* ── Push-only (без вставки в БД, только рассылка пирам) ── */ /* ── Stubs for chat_sync compatibility (todo: remove when chat_sync is fully replaced) ── */
void chat_core_push_message(struct chat_msg_submit* req); void chat_core_mark_sent(const char* ch_id, uint64_t ts, uint64_t dh, uint64_t peer_id);
void chat_core_ttl_delete(const char* ch_id, uint64_t node_id, uint64_t cutoff_us);
/* ── DB-операции для chat_sync ── */ /* ── DB-операции (оставлены для интроспекции) ── */
uint32_t chat_core_count(const char* ch_id);
int chat_core_chain_hash_at(const char* ch_id, uint32_t pos, uint8_t* hash_out);
int chat_core_insert_record(const char* ch_id, const uint8_t* rec, size_t len);
uint32_t chat_core_cursor_open(const char* ch_id);
int chat_core_cursor_next(uint32_t cursor_id, uint8_t* buf, size_t buf_size, size_t* out_len);
void chat_core_cursor_close(uint32_t cursor_id);
void chat_core_mark_sent(const char* ch_id, uint64_t ts, uint64_t dh, uint64_t node_id);
void chat_core_ttl_delete(const char* ch_id, uint64_t node_id, uint64_t cutoff_us);
int chat_core_list_channels(uint8_t* buf, size_t buf_size, size_t* out_len); int chat_core_list_channels(uint8_t* buf, size_t buf_size, size_t* out_len);
int chat_core_list_peers(const char* ch_id, uint8_t* buf, size_t buf_size, size_t* out_len); int chat_core_list_peers(const char* ch_id, uint8_t* buf, size_t buf_size, size_t* out_len);
int chat_core_load_nodeinfo(uint64_t node_id, uint8_t* buf, size_t buf_size, size_t* out_len); int chat_core_load_nodeinfo(uint64_t node_id, uint8_t* buf, size_t buf_size, size_t* out_len);
@ -86,6 +79,5 @@ void chat_core_update_my_name_trampoline(void* arg);
/* Трамплин для gui_bridge_post_uasync (GUI → uasync) */ /* Трамплин для gui_bridge_post_uasync (GUI → uasync) */
void chat_core_submit_trampoline(void* arg); void chat_core_submit_trampoline(void* arg);
void chat_core_push_trampoline(void* arg);
#endif /* CHAT_CORE_H */ #endif /* CHAT_CORE_H */

9
tools/chatgui/transport/utun_node.cpp

@ -22,6 +22,7 @@ extern "C" {
#include "../lib/mem.h" #include "../lib/mem.h"
#include "chat_sync.h" #include "chat_sync.h"
#include "chat_core.h" #include "chat_core.h"
#include "db_sync.h"
#include "gui_bridge.h" #include "gui_bridge.h"
} }
@ -237,10 +238,10 @@ void UtunNode::runLoop() {
etcp_router_bind(m_instance, ETCP_RT_ID_CHAT, recvCallback); etcp_router_bind(m_instance, ETCP_RT_ID_CHAT, recvCallback);
/* Initialize chat_core (DB) and chat_sync */ /* Initialize chat_core (DB) and db_sync (message sync) */
chat_core_init(m_instance, m_dbPath.toUtf8().constData()); chat_core_init(m_instance, m_dbPath.toUtf8().constData());
chat_sync_init(m_instance, NULL); db_sync_init(m_instance);
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "chat_core + chat_sync initialized"); DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "chat_core + db_sync initialized");
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "utun_node: entering poll loop"); DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "utun_node: entering poll loop");
while (!m_stop) { while (!m_stop) {
@ -248,7 +249,7 @@ void UtunNode::runLoop() {
} }
etcp_router_bind(m_instance, ETCP_RT_ID_CHAT, nullptr); etcp_router_bind(m_instance, ETCP_RT_ID_CHAT, nullptr);
chat_sync_destroy(m_instance); db_sync_destroy(m_instance);
chat_core_destroy(m_instance); chat_core_destroy(m_instance);
utun_instance_destroy(m_instance); utun_instance_destroy(m_instance);
m_instance = nullptr; m_instance = nullptr;

Loading…
Cancel
Save