Browse Source

перенос chat-логики из chatgui/transport в src/chat/

- Создан src/chat/ с chat_event.h/c — система нотификаций
- Перенесены chat_core, chat_msg, chat_channel, chat_profile, chat_status
- Перенесены chat_sync, member_sync, merkle_sync, db_sync
- gui_bridge_post() заменён на chat_event_post() с возможностью регистрации обработчика
- chat_event_set_handler() — chatgui регистрирует мост в gui_bridge
- Без обработчика: события логируются через DEBUG_DEBUG + log_dump
- chat интегрирован в utun_instance.c (при db_path в конфиге)
- Обновлены include-пути, Makefile.am, CMakeLists.txt
- Из chatgui удалены старые transport/ файлы, обновлён utun_node.cpp
topo_upd
Evgeny 2 months ago
parent
commit
ed6132fcda
  1. 26
      src/Makefile.am
  2. 62
      src/chat/chat_channel.c
  3. 30
      src/chat/chat_core.c
  4. 0
      src/chat/chat_core.h
  5. 4
      src/chat/chat_core_priv.h
  6. 29
      src/chat/chat_event.c
  7. 42
      src/chat/chat_event.h
  8. 22
      src/chat/chat_msg.c
  9. 26
      src/chat/chat_profile.c
  10. 20
      src/chat/chat_status.c
  11. 66
      src/chat/chat_sync.c
  12. 5
      src/chat/chat_sync.h
  13. 24
      src/chat/db_sync.c
  14. 0
      src/chat/db_sync.h
  15. 8
      src/chat/member_sync.c
  16. 0
      src/chat/member_sync.h
  17. 12
      src/chat/merkle_sync.c
  18. 0
      src/chat/merkle_sync.h
  19. 16
      src/utun_instance.c
  20. 4
      tests/Makefile.am
  21. 2
      tests/test_chat_sync_stress.c
  22. 2
      tests/test_db_sync.c
  23. 10
      tools/chatgui/CMakeLists.txt
  24. 1
      tools/chatgui/libutun/CMakeLists.txt
  25. 16
      tools/chatgui/transport/utun_node.cpp

26
src/Makefile.am

@ -18,7 +18,7 @@ utun_CORE_SOURCES = \
routing_layer/conn_mgr.c \ routing_layer/conn_mgr.c \
routing_layer/topo_recovery.c \ routing_layer/topo_recovery.c \
routing_layer/topo_group_connect.c \ routing_layer/topo_group_connect.c \
db_sync.c \ chat/db_sync.c \
routing_layer/routing.c \ routing_layer/routing.c \
tun_if.c \ tun_if.c \
tun_route.c \ tun_route.c \
@ -57,7 +57,16 @@ utun_CORE_SOURCES = \
lwip_tcp/lwip_pbuf.c \ lwip_tcp/lwip_pbuf.c \
lwip_tcp/lwip_tcp.c \ lwip_tcp/lwip_tcp.c \
lwip_tcp/lwip_tcp_in.c \ lwip_tcp/lwip_tcp_in.c \
lwip_tcp/lwip_tcp_out.c lwip_tcp/lwip_tcp_out.c \
chat/chat_event.c \
chat/chat_core.c \
chat/chat_msg.c \
chat/chat_channel.c \
chat/chat_profile.c \
chat/chat_status.c \
chat/chat_sync.c \
chat/member_sync.c \
chat/merkle_sync.c
# libutun: all core sources except main() # libutun: all core sources except main()
libutun_a_SOURCES = \ libutun_a_SOURCES = \
@ -75,7 +84,7 @@ libutun_a_SOURCES = \
routing_layer/conn_mgr.c \ routing_layer/conn_mgr.c \
routing_layer/topo_recovery.c \ routing_layer/topo_recovery.c \
routing_layer/topo_group_connect.c \ routing_layer/topo_group_connect.c \
db_sync.c \ chat/db_sync.c \
routing_layer/routing.c \ routing_layer/routing.c \
tun_if.c \ tun_if.c \
tun_route.c \ tun_route.c \
@ -114,7 +123,16 @@ libutun_a_SOURCES = \
lwip_tcp/lwip_pbuf.c \ lwip_tcp/lwip_pbuf.c \
lwip_tcp/lwip_tcp.c \ lwip_tcp/lwip_tcp.c \
lwip_tcp/lwip_tcp_in.c \ lwip_tcp/lwip_tcp_in.c \
lwip_tcp/lwip_tcp_out.c lwip_tcp/lwip_tcp_out.c \
chat/chat_event.c \
chat/chat_core.c \
chat/chat_msg.c \
chat/chat_channel.c \
chat/chat_profile.c \
chat/chat_status.c \
chat/chat_sync.c \
chat/member_sync.c \
chat/merkle_sync.c
libutun_a_CFLAGS = $(utun_CFLAGS) libutun_a_CFLAGS = $(utun_CFLAGS)
# Platform-specific TUN libs (Windows only) # Platform-specific TUN libs (Windows only)

62
tools/chatgui/transport/chat_channel.c → src/chat/chat_channel.c

@ -5,19 +5,19 @@
*/ */
#include "chat_core_priv.h" #include "chat_core_priv.h"
#include "gui_bridge.h" #include "chat_event.h"
#include "../../../src/utun_instance.h" #include "../utun_instance.h"
#include "../../../src/ntp_time.h" #include "../ntp_time.h"
#include "topo_node_sqlite.h" #include "../routing_layer/topo_node_sqlite.h"
#include "topo_group.h" #include "../routing_layer/topo_group.h"
#include "topo_group_connect.h" #include "../routing_layer/topo_group_connect.h"
#include "member_sync.h" #include "member_sync.h"
#include "secure_channel.h" #include "../transport_layer/secure_channel.h"
#include "../../../lib/mem.h" #include "../../lib/mem.h"
#include "../../../lib/ll_queue.h" #include "../../lib/ll_queue.h"
#include "../../../lib/platform_compat.h" #include "../../lib/platform_compat.h"
#include "etcp.h" #include "../transport_layer/etcp.h"
#include <openssl/sha.h> #include <openssl/sha.h>
#include <openssl/evp.h> #include <openssl/evp.h>
@ -91,10 +91,10 @@ void chat_core_ensure_channel_ready(const char* ch_id) {
if (si) { if (si) {
si_register(si, ch_id); si_register(si, ch_id);
db_sync_set_insert_cb(si, on_msg_inserted, u_strdup(ch_id)); db_sync_set_insert_cb(si, on_msg_inserted, u_strdup(ch_id));
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: channel ready ch=%s tbl=%s hash=0x%016llx", DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: channel ready ch=%s tbl=%s hash=0x%016llx",
CC_ID, ch_id, tbl_msg, (unsigned long long)ch_hash); CC_ID, ch_id, tbl_msg, (unsigned long long)ch_hash);
} else { } else {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: db_sync_instance_add failed for ch=%s", DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: db_sync_instance_add failed for ch=%s",
CC_ID, ch_id); CC_ID, ch_id);
} }
} }
@ -103,13 +103,13 @@ void chat_core_ensure_channel_ready(const char* ch_id) {
void chat_core_create_channel(struct chat_channel_create* req) { void chat_core_create_channel(struct chat_channel_create* req) {
if (!g_cc.initialized) { if (!g_cc.initialized) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: create_channel NOT INITIALIZED ch=%s", DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: create_channel NOT INITIALIZED ch=%s",
CC_ID, req ? req->channel_id : "(null)"); CC_ID, req ? req->channel_id : "(null)");
return; return;
} }
if (!req) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: create_channel req=NULL", CC_ID); return; } if (!req) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: create_channel req=NULL", CC_ID); return; }
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: create_channel BEGIN ch=%s name=%s", DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: create_channel BEGIN ch=%s name=%s",
CC_ID, req->channel_id, req->name); CC_ID, req->channel_id, req->name);
int rc = topo_node_sqlite_channel_put(g_cc.db, int rc = topo_node_sqlite_channel_put(g_cc.db,
@ -119,10 +119,10 @@ void chat_core_create_channel(struct chat_channel_create* req) {
req->signature); req->signature);
if (rc != 0) { if (rc != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: channel_put FAILED ch=%s rc=%d", DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: channel_put FAILED ch=%s rc=%d",
CC_ID, req->channel_id, rc); CC_ID, req->channel_id, rc);
} else { } else {
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: channel_put OK ch=%s name=%s owner=0x%016llx", DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: channel_put OK ch=%s name=%s owner=0x%016llx",
CC_ID, req->channel_id, req->name, (unsigned long long)req->owner_node_id); CC_ID, req->channel_id, req->name, (unsigned long long)req->owner_node_id);
chat_core_ensure_channel_ready(req->channel_id); chat_core_ensure_channel_ready(req->channel_id);
@ -148,17 +148,17 @@ void chat_core_create_channel(struct chat_channel_create* req) {
g_cc.inst->my_keys.public_key, g_cc.inst->my_ed25519_pubkey, g_cc.inst->my_keys.public_key, g_cc.inst->my_ed25519_pubkey,
join_sig, join_ts, NULL, 0, juser3, my_addrs, my_addr_cnt); join_sig, join_ts, NULL, 0, juser3, my_addrs, my_addr_cnt);
if (mrc != 0) { if (mrc != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: member_sync_put(self) FAILED ch=%s rc=%d", DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: member_sync_put(self) FAILED ch=%s rc=%d",
CC_ID, req->channel_id, mrc); CC_ID, req->channel_id, mrc);
} }
} }
gui_bridge_post(GUI_EVT_CHANNEL_UPDATED, data, 1 + ch_id_len); chat_event_post(CHAT_EVT_CHANNEL_UPDATED, data, 1 + ch_id_len);
} }
} }
void chat_core_create_channel_trampoline(void* arg) { void chat_core_create_channel_trampoline(void* arg) {
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "chat_core: TRAMPOLINE invoked arg=%p", arg); DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "chat_core: TRAMPOLINE invoked arg=%p", arg);
struct chat_channel_create* req = (struct chat_channel_create*)arg; struct chat_channel_create* req = (struct chat_channel_create*)arg;
chat_core_create_channel(req); chat_core_create_channel(req);
u_free(req); u_free(req);
@ -171,10 +171,10 @@ void chat_core_create_channel_auto_trampoline(void* arg) {
void chat_core_create_channel_auto(const char* name) { void chat_core_create_channel_auto(const char* name) {
if (!g_cc.initialized || !name) { if (!g_cc.initialized || !name) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: create_channel_auto — not initialized or NULL name", CC_ID); DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: create_channel_auto — not initialized or NULL name", CC_ID);
return; return;
} }
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: create_channel_auto BEGIN name=%s", CC_ID, name); DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: create_channel_auto BEGIN name=%s", CC_ID, name);
uint8_t x25519_pub[32] = {0}, x25519_priv[32] = {0}; uint8_t x25519_pub[32] = {0}, x25519_priv[32] = {0};
{ EVP_PKEY* pkey = NULL; EVP_PKEY_CTX* ctx = EVP_PKEY_CTX_new_id(EVP_PKEY_X25519, NULL); { EVP_PKEY* pkey = NULL; EVP_PKEY_CTX* ctx = EVP_PKEY_CTX_new_id(EVP_PKEY_X25519, NULL);
@ -226,31 +226,31 @@ void chat_core_create_channel_auto(const char* name) {
/* ─── управление подключением к каналу (выбор в GUI) ─── */ /* ─── управление подключением к каналу (выбор в GUI) ─── */
void chat_core_connect_channel(const char* ch_id) { void chat_core_connect_channel(const char* ch_id) {
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTIVITY, "%s: connect_channel called ch_id=%s initialized=%d", CC_ID, ch_id ? ch_id : "(null)", g_cc.initialized); DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: connect_channel called ch_id=%s initialized=%d", CC_ID, ch_id ? ch_id : "(null)", g_cc.initialized);
if (!g_cc.initialized || !ch_id || !ch_id[0]) { DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: connect_channel — skip (not ready)", CC_ID); return; } if (!g_cc.initialized || !ch_id || !ch_id[0]) { DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: connect_channel — skip (not ready)", CC_ID); return; }
if (!g_cc.inst->topo_groups) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: connect_channel — no topo_groups", CC_ID); return; } if (!g_cc.inst->topo_groups) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: connect_channel — no topo_groups", CC_ID); return; }
uint64_t gid = strtoull(ch_id, NULL, 10); uint64_t gid = strtoull(ch_id, NULL, 10);
if (gid == 0) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: connect_channel — invalid ch_id=%s", CC_ID, ch_id); return; } if (gid == 0) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: connect_channel — invalid ch_id=%s", CC_ID, ch_id); return; }
chat_core_ensure_channel_ready(ch_id); chat_core_ensure_channel_ready(ch_id);
struct TOPO_GROUP* group = topo_groups_find(g_cc.inst->topo_groups, gid); struct TOPO_GROUP* group = topo_groups_find(g_cc.inst->topo_groups, gid);
if (!group) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: connect_channel — group not found ch=%s gid=%016llx total_groups=%d", CC_ID, ch_id, (unsigned long long)gid, queue_entry_count(g_cc.inst->topo_groups->group_list)); return; } if (!group) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: connect_channel — group not found ch=%s gid=%016llx total_groups=%d", CC_ID, ch_id, (unsigned long long)gid, queue_entry_count(g_cc.inst->topo_groups->group_list)); return; }
if (group->connect) { if (group->connect) {
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: connect_channel — already in progress ch=%s active=%d", DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: connect_channel — already in progress ch=%s active=%d",
CC_ID, ch_id, topo_group_connect_active_count(group)); CC_ID, ch_id, topo_group_connect_active_count(group));
return; return;
} }
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: connect_channel ch=%s gid=%016llx group_type=%d — starting connect_init", CC_ID, ch_id, (unsigned long long)gid, group->group_type); DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: connect_channel ch=%s gid=%016llx group_type=%d — starting connect_init", CC_ID, ch_id, (unsigned long long)gid, group->group_type);
topo_group_connect_init(group); topo_group_connect_init(group);
} }
void chat_core_connect_channel_trampoline(void* arg) { void chat_core_connect_channel_trampoline(void* arg) {
char* ch_id = (char*)arg; char* ch_id = (char*)arg;
DEBUG_TRACE(DEBUG_CATEGORY_CONNECTIVITY, "%s: connect_channel_trampoline ch_id=%s", CC_ID, ch_id ? ch_id : "(null)"); DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: connect_channel_trampoline ch_id=%s", CC_ID, ch_id ? ch_id : "(null)");
if (!ch_id) return; if (!ch_id) return;
chat_core_connect_channel(ch_id); chat_core_connect_channel(ch_id);
u_free(ch_id); u_free(ch_id);

30
tools/chatgui/transport/chat_core.c → src/chat/chat_core.c

@ -11,16 +11,16 @@
*/ */
#include "chat_core_priv.h" #include "chat_core_priv.h"
#include "gui_bridge.h" #include "chat_event.h"
#include "../../../src/utun_instance.h" #include "../utun_instance.h"
#include "../../../src/ntp_time.h" #include "../ntp_time.h"
#include "topo_node_sqlite.h" #include "../routing_layer/topo_node_sqlite.h"
#include "../../../lib/mem.h" #include "../../lib/mem.h"
#include <string.h> #include <string.h>
#include <stdio.h> #include <stdio.h>
#include "../../../lib/platform_compat.h" #include "../../lib/platform_compat.h"
/* ─── глобальное состояние (extern в chat_core_priv.h) ─── */ /* ─── глобальное состояние (extern в chat_core_priv.h) ─── */
@ -54,7 +54,7 @@ int db_exec(const char* sql) {
char* err = NULL; char* err = NULL;
int rc = sqlite3_exec(g_cc.db, sql, NULL, NULL, &err); int rc = sqlite3_exec(g_cc.db, sql, NULL, NULL, &err);
if (rc != SQLITE_OK) { if (rc != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: sql error: %s", CC_ID, err); DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: sql error: %s", CC_ID, err);
sqlite3_free(err); sqlite3_free(err);
} }
return rc; return rc;
@ -91,12 +91,12 @@ int chat_core_init(struct UTUN_INSTANCE* inst, const char* db_path) {
if (inst->topo_sqlite_db) { if (inst->topo_sqlite_db) {
g_cc.db = inst->topo_sqlite_db; g_cc.shared_db = 1; g_cc.db = inst->topo_sqlite_db; g_cc.shared_db = 1;
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: using shared SQLite db=%p", CC_ID, (void*)g_cc.db); DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: using shared SQLite db=%p", CC_ID, (void*)g_cc.db);
} else { } else {
int rc = sqlite3_open_v2(db_path, &g_cc.db, int rc = sqlite3_open_v2(db_path, &g_cc.db,
SQLITE_OPEN_READWRITE | SQLITE_OPEN_CREATE | SQLITE_OPEN_FULLMUTEX, NULL); SQLITE_OPEN_READWRITE | SQLITE_OPEN_CREATE | SQLITE_OPEN_FULLMUTEX, NULL);
if (rc != SQLITE_OK) { if (rc != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: cannot open DB %s: %s", DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: cannot open DB %s: %s",
CC_ID, db_path, sqlite3_errmsg(g_cc.db)); CC_ID, db_path, sqlite3_errmsg(g_cc.db));
sqlite3_close(g_cc.db); g_cc.db = NULL; return -1; sqlite3_close(g_cc.db); g_cc.db = NULL; return -1;
} }
@ -167,7 +167,7 @@ int chat_core_init(struct UTUN_INSTANCE* inst, const char* db_path) {
if (sqlite3_prepare_v2(g_cc.db, "SELECT COUNT(*) FROM accounts", if (sqlite3_prepare_v2(g_cc.db, "SELECT COUNT(*) FROM accounts",
-1, &cnt_st, NULL) == SQLITE_OK) { -1, &cnt_st, NULL) == SQLITE_OK) {
if (sqlite3_step(cnt_st) == SQLITE_ROW && sqlite3_column_int(cnt_st, 0) == 0) { if (sqlite3_step(cnt_st) == SQLITE_ROW && sqlite3_column_int(cnt_st, 0) == 0) {
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: seeding 5 stub accounts", CC_ID); DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: seeding 5 stub accounts", CC_ID);
struct { int64_t id; const char* name; const char* letter; const char* color; } accs[] = { struct { int64_t id; const char* name; const char* letter; const char* color; } accs[] = {
{1, "Alice", "A", "#4A90E2"}, {2, "Bob", "B", "#7ED321"}, {1, "Alice", "A", "#4A90E2"}, {2, "Bob", "B", "#7ED321"},
{3, "Carol", "C", "#F5A623"}, {4, "Dave", "D", "#9013FE"}, {3, "Carol", "C", "#F5A623"}, {4, "Dave", "D", "#9013FE"},
@ -203,16 +203,16 @@ int chat_core_init(struct UTUN_INSTANCE* inst, const char* db_path) {
if (ch && ch[0]) { chat_core_ensure_channel_ready(ch); loaded++; } if (ch && ch[0]) { chat_core_ensure_channel_ready(ch); loaded++; }
} }
sqlite3_finalize(stmt); sqlite3_finalize(stmt);
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: loaded %d channels from DB", CC_ID, loaded); DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: loaded %d channels from DB", CC_ID, loaded);
} }
} }
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: initialized, db=%s node_id=0x%016llx", DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: initialized, db=%s node_id=0x%016llx",
CC_ID, db_path, (unsigned long long)g_cc.my_node_id); CC_ID, db_path, (unsigned long long)g_cc.my_node_id);
{ uint8_t nid[8]; memcpy(nid, &g_cc.my_node_id, 8); gui_bridge_post(GUI_EVT_MY_NODE_ID, nid, 8); } { uint8_t nid[8]; memcpy(nid, &g_cc.my_node_id, 8); chat_event_post(CHAT_EVT_MY_NODE_ID, nid, 8); }
gui_bridge_post(GUI_EVT_DB_READY, NULL, 0); chat_event_post(CHAT_EVT_DB_READY, NULL, 0);
return 0; return 0;
} }
@ -228,7 +228,7 @@ void chat_core_destroy(struct UTUN_INSTANCE* inst) {
if (g_cc.db && !g_cc.shared_db) { sqlite3_close(g_cc.db); } if (g_cc.db && !g_cc.shared_db) { sqlite3_close(g_cc.db); }
g_cc.db = NULL; g_cc.inst = NULL; g_cc.db = NULL; g_cc.inst = NULL;
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: destroyed", CC_ID); DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: destroyed", CC_ID);
} }
sqlite3* chat_core_get_db(void) { return g_cc.db; } sqlite3* chat_core_get_db(void) { return g_cc.db; }

0
tools/chatgui/transport/chat_core.h → src/chat/chat_core.h

4
tools/chatgui/transport/chat_core_priv.h → src/chat/chat_core_priv.h

@ -2,7 +2,7 @@
* chat_core_priv.h — внутренний заголовок для под-модулей chat_core * chat_core_priv.h — внутренний заголовок для под-модулей chat_core
* *
* Предоставляет доступ к глобальному состоянию g_cc и общим хелперам. * Предоставляет доступ к глобальному состоянию g_cc и общим хелперам.
* Не включать извне transport/ — только для chat_core*.c. * Не включать извне chat/ — только для chat_core*.c.
*/ */
#ifndef CHAT_CORE_PRIV_H #ifndef CHAT_CORE_PRIV_H
@ -11,7 +11,7 @@
#include "chat_core.h" #include "chat_core.h"
#include "db_sync.h" #include "db_sync.h"
#include "../../../lib/debug_config.h" #include "../../lib/debug_config.h"
#include <sqlite3.h> #include <sqlite3.h>
#include <string.h> #include <string.h>

29
src/chat/chat_event.c

@ -0,0 +1,29 @@
/*
* chat_event.c — реализация системы нотификаций
*
* Пока обработчик не зарегистрирован — события логируются.
* При регистрации обработчика — форвардятся ему синхронно.
*/
#include "chat_event.h"
#include "../lib/debug_config.h"
static chat_event_handler_fn g_handler = NULL;
void chat_event_set_handler(chat_event_handler_fn handler) {
g_handler = handler;
}
void chat_event_post(int type, const uint8_t* data, int len) {
if (g_handler) { g_handler(type, data, len); return; }
static const char* names[] = {
[1]="MSG_RECEIVED", [2]="CONNECT_RESULT", [3]="NEW_PEER",
[4]="CHANNEL_UPDATED", [5]="MEMBERS_CHANGED", [6]="MY_NODE_ID",
[7]="AUTO_CONNECT_STATUS", [8]="CHANNEL_PEERS_ONLINE",
[9]="DB_READY", [10]="STATUS_REFRESH", [11]="NODE_CHANGED",
};
const char* n = (type >= 1 && type <= 11) ? names[type] : "?";
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "chat_event: %s(%d) data=%d bytes", n, type, len);
if (data && len > 0) log_dump(DEBUG_LEVEL_DEBUG, DEBUG_CATEGORY_DEBUG, " ", data, len);
}

42
src/chat/chat_event.h

@ -0,0 +1,42 @@
/*
* chat_event.h — система нотификаций о событиях чата
*
* Производители (chat_core, chat_sync, ...) вызывают chat_event_post().
* Потребители (chatgui, headless, будущий HTTP API) регистрируют один обработчик.
*
* Без обработчика — данные логируются через DEBUG_DEBUG + log_dump.
* Все вызовы — из uasync-потока.
*/
#ifndef CHAT_EVENT_H
#define CHAT_EVENT_H
#include <stdint.h>
#include <stddef.h>
#ifdef __cplusplus
extern "C" {
#endif
/* Типы событий (совпадают с GUI_EVT_* в gui_bridge.h для прямой совместимости) */
#define CHAT_EVT_MSG_RECEIVED 1 /* [ch_id_len:1][ch_id:var][author_node_id:8] */
#define CHAT_EVT_CONNECT_RESULT 2 /* [node_id:8][result:4][channel_id:8] */
#define CHAT_EVT_NEW_PEER 3 /* [node_id:8] */
#define CHAT_EVT_CHANNEL_UPDATED 4 /* [ch_id_len:1][ch_id:var] */
#define CHAT_EVT_MEMBERS_CHANGED 5 /* [ch_id_len:1][ch_id:var] */
#define CHAT_EVT_MY_NODE_ID 6 /* [node_id:8] */
#define CHAT_EVT_AUTO_CONNECT_STATUS 7 /* [status:1][total_tried:2][node_count:2][success:2] */
#define CHAT_EVT_CHANNEL_PEERS_ONLINE 8 /* [ch_id_len:1][ch_id:var][online_count:2] */
#define CHAT_EVT_DB_READY 9 /* data: none */
#define CHAT_EVT_STATUS_REFRESH 10 /* data: status text */
#define CHAT_EVT_NODE_CHANGED 11 /* [node_id:8] */
typedef void (*chat_event_handler_fn)(int type, const uint8_t* data, int len);
void chat_event_set_handler(chat_event_handler_fn handler);
void chat_event_post(int type, const uint8_t* data, int len);
#ifdef __cplusplus
}
#endif
#endif /* CHAT_EVENT_H */

22
tools/chatgui/transport/chat_msg.c → src/chat/chat_msg.c

@ -5,12 +5,12 @@
*/ */
#include "chat_core_priv.h" #include "chat_core_priv.h"
#include "gui_bridge.h" #include "chat_event.h"
#include "../../../src/utun_instance.h" #include "../utun_instance.h"
#include "secure_channel.h" #include "../transport_layer/secure_channel.h"
#include "../../../lib/mem.h" #include "../../lib/mem.h"
#include "../../../lib/platform_compat.h" #include "../../lib/platform_compat.h"
#include <openssl/sha.h> #include <openssl/sha.h>
@ -29,14 +29,14 @@ static int msg_get_prev_chain_hash(const char* tbl, uint8_t out[32]) {
} }
void chat_core_submit_message(struct chat_msg_submit* req) { void chat_core_submit_message(struct chat_msg_submit* req) {
if (!g_cc.initialized) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: submit before init", CC_ID); return; } if (!g_cc.initialized) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: submit before init", CC_ID); return; }
if (!req) return; if (!req) return;
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: SUBMIT ch=%s ct=%s len=%u", DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: SUBMIT ch=%s ct=%s len=%u",
CC_ID, req->channel_id, req->content_type, req->data_len); CC_ID, req->channel_id, req->content_type, req->data_len);
struct DB_SYNC_INSTANCE* si = si_find(req->channel_id); struct DB_SYNC_INSTANCE* si = si_find(req->channel_id);
if (!si) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: no db_sync instance for ch=%s", CC_ID, req->channel_id); return; } if (!si) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: no db_sync instance for ch=%s", CC_ID, req->channel_id); return; }
char json[4096]; char json[4096];
snprintf(json, sizeof(json), snprintf(json, sizeof(json),
@ -50,13 +50,13 @@ void chat_core_submit_message(struct chat_msg_submit* req) {
size_t jl = strlen(json); memcpy(sig_msg + soff, json, jl); soff += jl; size_t jl = strlen(json); memcpy(sig_msg + soff, json, jl); soff += jl;
uint8_t sig[64]; uint8_t sig[64];
if (sc_ed25519_sign(g_cc.inst->my_ed25519_privkey, sig_msg, soff, sig) != SC_OK) { if (sc_ed25519_sign(g_cc.inst->my_ed25519_privkey, sig_msg, soff, sig) != SC_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: Ed25519 sign failed", CC_ID); DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: Ed25519 sign failed", CC_ID);
return; return;
} }
int ret = db_sync_insert_signed(si, json, strlen(json), sig, 64, ts); int ret = db_sync_insert_signed(si, json, strlen(json), sig, 64, ts);
if (ret != 0) { if (ret != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: db_sync_insert_signed failed ret=%d", CC_ID, ret); DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: db_sync_insert_signed failed ret=%d", CC_ID, ret);
return; return;
} }
} }
@ -218,7 +218,7 @@ void on_msg_inserted(struct DB_SYNC_INSTANCE* si, uint64_t record_ts, const char
const char* ch_id = (const char*)arg; const char* ch_id = (const char*)arg;
uint8_t evt[73]; uint8_t cl = (uint8_t)strlen(ch_id); evt[0] = cl; memcpy(evt + 1, ch_id, cl); uint8_t evt[73]; uint8_t cl = (uint8_t)strlen(ch_id); evt[0] = cl; memcpy(evt + 1, ch_id, cl);
memcpy(evt + 1 + cl, &author, 8); memcpy(evt + 1 + cl, &author, 8);
gui_bridge_post(GUI_EVT_MSG_RECEIVED, evt, 1 + cl + 8); chat_event_post(CHAT_EVT_MSG_RECEIVED, evt, 1 + cl + 8);
sqlite3_stmt* st = NULL; sqlite3_stmt* st = NULL;
sqlite3_prepare_v2(g_cc.db, "UPDATE channels SET last_msg_at = ? WHERE channel_id = ?", -1, &st, NULL); sqlite3_prepare_v2(g_cc.db, "UPDATE channels SET last_msg_at = ? WHERE channel_id = ?", -1, &st, NULL);

26
tools/chatgui/transport/chat_profile.c → src/chat/chat_profile.c

@ -5,18 +5,18 @@
*/ */
#include "chat_core_priv.h" #include "chat_core_priv.h"
#include "gui_bridge.h" #include "chat_event.h"
#include "../../../lib/json_flat.h" #include "../../lib/json_flat.h"
#include "topo_node_sqlite.h" #include "../routing_layer/topo_node_sqlite.h"
#include "member_sync.h" #include "member_sync.h"
#include "../../../src/utun_instance.h" #include "../utun_instance.h"
#include "etcp.h" #include "../transport_layer/etcp.h"
#include "secure_channel.h" #include "../transport_layer/secure_channel.h"
#include "../../../src/ntp_time.h" #include "../ntp_time.h"
#include "../../../lib/u_async.h" #include "../../lib/u_async.h"
#include "../../../lib/mem.h" #include "../../lib/mem.h"
#include "../../../lib/platform_compat.h" #include "../../lib/platform_compat.h"
#include <openssl/evp.h> #include <openssl/evp.h>
@ -68,13 +68,13 @@ void chat_core_update_my_name(const char* name) {
char juser[256]; snprintf(juser, sizeof(juser), "{\"name\":\"%s\"}", name); char juser[256]; snprintf(juser, sizeof(juser), "{\"name\":\"%s\"}", name);
int r1 = topo_node_sqlite_member_put(g_cc.db, ch, myid, new_sig, join_ts, NULL, 0, int r1 = topo_node_sqlite_member_put(g_cc.db, ch, myid, new_sig, join_ts, NULL, 0,
g_cc.inst->my_keys.public_key, g_cc.inst->my_ed25519_pubkey, juser, NULL); g_cc.inst->my_keys.public_key, g_cc.inst->my_ed25519_pubkey, juser, NULL);
if (r1 != 0) DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: member_put(my_name) FAILED ch=%s rc=%d", CC_ID, ch, r1); if (r1 != 0) DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: member_put(my_name) FAILED ch=%s rc=%d", CC_ID, ch, r1);
int r2 = member_sync_put(g_cc.inst, ch, myid, g_cc.inst->my_keys.public_key, int r2 = member_sync_put(g_cc.inst, ch, myid, g_cc.inst->my_keys.public_key,
g_cc.inst->my_ed25519_pubkey, new_sig, join_ts, NULL, 0, juser, NULL, 0); g_cc.inst->my_ed25519_pubkey, new_sig, join_ts, NULL, 0, juser, NULL, 0);
if (r2 != 0) DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: member_sync_put(my_name) FAILED ch=%s rc=%d", CC_ID, ch, r2); if (r2 != 0) DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: member_sync_put(my_name) FAILED ch=%s rc=%d", CC_ID, ch, r2);
uint8_t evt[65]; uint8_t cl = (uint8_t)strlen(ch); uint8_t evt[65]; uint8_t cl = (uint8_t)strlen(ch);
evt[0] = cl; memcpy(evt + 1, ch, cl); evt[0] = cl; memcpy(evt + 1, ch, cl);
gui_bridge_post(GUI_EVT_MEMBERS_CHANGED, evt, 1 + cl); chat_event_post(CHAT_EVT_MEMBERS_CHANGED, evt, 1 + cl);
} }
sqlite3_finalize(ps); sqlite3_finalize(ps);
} }

20
tools/chatgui/transport/chat_status.c → src/chat/chat_status.c

@ -1,18 +1,18 @@
/* /*
* chat_status.c — сбор статуса (NTP + ETCP connections) и отправка в GUI * chat_status.c — сбор статуса (NTP + ETCP connections) и отправка
* *
* Вынесено из chat_core.c для уменьшения размера модуля. * Вынесено из chat_core.c для уменьшения размера модуля.
*/ */
#include "chat_core_priv.h" #include "chat_core_priv.h"
#include "gui_bridge.h" #include "chat_event.h"
#include "../../../src/utun_instance.h" #include "../utun_instance.h"
#include "../../../src/ntp_time.h" #include "../ntp_time.h"
#include "../../../lib/ll_queue.h" #include "../../lib/ll_queue.h"
#include "../../../lib/platform_compat.h" #include "../../lib/platform_compat.h"
#include "etcp.h" #include "../transport_layer/etcp.h"
#include "etcp_connections.h" #include "../transport_layer/etcp_connections.h"
static const char* nat_type_str(uint8_t t) { static const char* nat_type_str(uint8_t t) {
switch (t) { case 0: return "UNKNOWN"; case 1: return "EIM"; case 2: return "STRICT"; case 3: return "DIRECT"; default: return "?"; } switch (t) { case 0: return "UNKNOWN"; case 1: return "EIM"; case 2: return "STRICT"; case 3: return "DIRECT"; default: return "?"; }
@ -36,7 +36,7 @@ static void chat_core_collect_status(void) {
if (!g_cc.initialized || !g_cc.inst) { if (!g_cc.initialized || !g_cc.inst) {
off = snprintf(buf, sizeof(buf), "uTun not initialized\n"); off = snprintf(buf, sizeof(buf), "uTun not initialized\n");
gui_bridge_post(GUI_EVT_STATUS_REFRESH, (const uint8_t*)buf, off); chat_event_post(CHAT_EVT_STATUS_REFRESH, (const uint8_t*)buf, off);
return; return;
} }
@ -109,7 +109,7 @@ static void chat_core_collect_status(void) {
entry = entry->next; entry = entry->next;
} }
gui_bridge_post(GUI_EVT_STATUS_REFRESH, (const uint8_t*)buf, off); chat_event_post(CHAT_EVT_STATUS_REFRESH, (const uint8_t*)buf, off);
} }
void chat_core_collect_status_trampoline(void* arg) { void chat_core_collect_status_trampoline(void* arg) {

66
tools/chatgui/transport/chat_sync.c → src/chat/chat_sync.c

@ -1,24 +1,24 @@
#include "chat_sync.h" #include "chat_sync.h"
#include "chat_core.h" #include "chat_core.h"
#include "gui_bridge.h" #include "chat_event.h"
#include "topo_node_sqlite.h" #include "../routing_layer/topo_node_sqlite.h"
#include "topo_group.h" #include "../routing_layer/topo_group.h"
#include "conn_mgr.h" #include "../routing_layer/conn_mgr.h"
#include "topo_group_connect.h" #include "../routing_layer/topo_group_connect.h"
#include "../../../lib/json_flat.h" #include "../../lib/json_flat.h"
#include "member_sync.h" #include "member_sync.h"
#include "merkle_sync.h" #include "merkle_sync.h"
#include "../../../src/utun_instance.h" #include "../utun_instance.h"
#include "etcp_api.h" #include "../transport_layer/etcp_api.h"
#include "etcp.h" #include "../transport_layer/etcp.h"
#include "secure_channel.h" #include "../transport_layer/secure_channel.h"
#include "../../../src/ntp_time.h" #include "../ntp_time.h"
#include "../../../lib/u_async.h" #include "../../lib/u_async.h"
#include "../../../lib/ll_queue.h" #include "../../lib/ll_queue.h"
#include "../../../lib/debug_config.h" #include "../../lib/debug_config.h"
#include "../../../lib/mem.h" #include "../../lib/mem.h"
#include "../../../lib/platform_compat.h" #include "../../lib/platform_compat.h"
#include <string.h> #include <string.h>
#include <sqlite3.h> #include <sqlite3.h>
@ -170,7 +170,7 @@ static void chat_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: RECV ERROR from=%016llx ch=%s", CS_ID, (unsigned long long)peer, ch_id); DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: RECV ERROR from=%016llx ch=%s", CS_ID, (unsigned long long)peer, ch_id);
if (g_cs->info_req_timer) { uasync_cancel_timeout(g_cs->inst->ua, g_cs->info_req_timer); g_cs->info_req_timer = NULL; } if (g_cs->info_req_timer) { uasync_cancel_timeout(g_cs->inst->ua, g_cs->info_req_timer); g_cs->info_req_timer = NULL; }
uint8_t err[20]; memcpy(err, &peer, 8); int r = CONN_MGR_ERR_NOT_FOUND; memcpy(err + 8, &r, 4); uint8_t err[20]; memcpy(err, &peer, 8); int r = CONN_MGR_ERR_NOT_FOUND; memcpy(err + 8, &r, 4);
memcpy(err + 12, &g_cs->pending_invite_ch_id, 8); gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20); memcpy(err + 12, &g_cs->pending_invite_ch_id, 8); chat_event_post(CHAT_EVT_CONNECT_RESULT, err, 20);
g_cs->pending_invite_node_id = 0; g_cs->pending_invite_ch_id = 0; g_cs->pending_invite_node_id = 0; g_cs->pending_invite_ch_id = 0;
break; break;
} }
@ -192,7 +192,7 @@ static void cs_info_req_timeout_cb(void* arg) {
uint8_t err[20]; memcpy(err, &cs->pending_invite_node_id, 8); uint8_t err[20]; memcpy(err, &cs->pending_invite_node_id, 8);
int r = CONN_MGR_ERR_TIMEOUT; memcpy(err + 8, &r, 4); int r = CONN_MGR_ERR_TIMEOUT; memcpy(err + 8, &r, 4);
memcpy(err + 12, &cs->pending_invite_ch_id, 8); memcpy(err + 12, &cs->pending_invite_ch_id, 8);
gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20); chat_event_post(CHAT_EVT_CONNECT_RESULT, err, 20);
cs->pending_invite_node_id = 0; cs->pending_invite_node_id = 0;
cs->pending_invite_ch_id = 0; cs->pending_invite_ch_id = 0;
} }
@ -207,7 +207,7 @@ static void cs_join_timeout_cb(void* arg) {
uint8_t err[20]; memcpy(err, &peer, 8); uint8_t err[20]; memcpy(err, &peer, 8);
int r = CONN_MGR_ERR_TIMEOUT; memcpy(err + 8, &r, 4); int r = CONN_MGR_ERR_TIMEOUT; memcpy(err + 8, &r, 4);
memcpy(err + 12, &cs->pending_invite_ch_id, 8); memcpy(err + 12, &cs->pending_invite_ch_id, 8);
gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20); chat_event_post(CHAT_EVT_CONNECT_RESULT, err, 20);
cs->pending_invite_node_id = 0; cs->pending_invite_node_id = 0;
cs->pending_invite_ch_id = 0; cs->pending_invite_ch_id = 0;
} }
@ -265,7 +265,7 @@ static void cs_post_channel_online(struct chat_sync* cs, const char* ch_id) {
uint8_t evt[66]; evt[0] = (uint8_t)cl; uint8_t evt[66]; evt[0] = (uint8_t)cl;
memcpy(evt + 1, ch_id, cl); memcpy(evt + 1, ch_id, cl);
uint16_t oc = (uint16_t)online; memcpy(evt + 1 + cl, &oc, 2); 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); chat_event_post(CHAT_EVT_CHANNEL_PEERS_ONLINE, evt, 1 + (int)cl + 2);
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: PEERS_ONLINE ch=%.*s total_online=%d", CS_ID, (int)cl, ch_id, online); DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: PEERS_ONLINE ch=%.*s total_online=%d", CS_ID, (int)cl, ch_id, online);
} }
@ -315,7 +315,7 @@ static void cs_on_remote_status_changed(uint64_t peer, int online) {
if (cl > 63) cl = 63; if (cl > 63) cl = 63;
uint8_t evt[65]; evt[0] = (uint8_t)cl; uint8_t evt[65]; evt[0] = (uint8_t)cl;
memcpy(evt + 1, g_cs->channels[i].channel_id, cl); memcpy(evt + 1, g_cs->channels[i].channel_id, cl);
gui_bridge_post(GUI_EVT_MEMBERS_CHANGED, evt, 1 + (int)cl); chat_event_post(CHAT_EVT_MEMBERS_CHANGED, evt, 1 + (int)cl);
cs_post_channel_online(g_cs, g_cs->channels[i].channel_id); cs_post_channel_online(g_cs, g_cs->channels[i].channel_id);
break; break;
} }
@ -337,7 +337,7 @@ static void cs_on_peer_status_changed(uint64_t peer, int online) {
if (cl > 63) cl = 63; if (cl > 63) cl = 63;
uint8_t evt[65]; evt[0] = (uint8_t)cl; uint8_t evt[65]; evt[0] = (uint8_t)cl;
memcpy(evt + 1, g_cs->channels[i].channel_id, cl); memcpy(evt + 1, g_cs->channels[i].channel_id, cl);
gui_bridge_post(GUI_EVT_MEMBERS_CHANGED, evt, 1 + (int)cl); chat_event_post(CHAT_EVT_MEMBERS_CHANGED, evt, 1 + (int)cl);
cs_post_channel_online(g_cs, g_cs->channels[i].channel_id); cs_post_channel_online(g_cs, g_cs->channels[i].channel_id);
/* 3. Рассылка дельты всем synced-соседям */ /* 3. Рассылка дельты всем synced-соседям */
@ -523,9 +523,7 @@ static void refresh_timer_cb(void* arg) {
/* ── Public API ── */ /* ── Public API ── */
int chat_sync_init(struct UTUN_INSTANCE* inst, int chat_sync_init(struct UTUN_INSTANCE* inst) {
void (*gui_cb)(void*, int, const uint8_t*, int)) {
(void)gui_cb;
if (!inst) return -1; if (!inst) return -1;
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: init", CS_ID); DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: init", CS_ID);
struct chat_sync* cs = u_calloc(1, sizeof(*cs)); struct chat_sync* cs = u_calloc(1, sizeof(*cs));
@ -674,7 +672,7 @@ void chat_sync_connect_from_invite(uint64_t channel_id, uint64_t node_id,
if (!g_cs || !g_cs->inst || !g_cs->inst->ua) { if (!g_cs || !g_cs->inst || !g_cs->inst->ua) {
int r = -7; int r = -7;
uint8_t err[12]; memcpy(err, &node_id, 8); memcpy(err + 8, &r, 4); uint8_t err[12]; memcpy(err, &node_id, 8); memcpy(err + 8, &r, 4);
gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 12); chat_event_post(CHAT_EVT_CONNECT_RESULT, err, 12);
return; return;
} }
@ -701,7 +699,7 @@ void chat_sync_connect_from_invite(uint64_t channel_id, uint64_t node_id,
struct cm_invite_wrap { struct chat_invite inv; }* w = u_malloc(sizeof(struct cm_invite_wrap)); struct cm_invite_wrap { struct chat_invite inv; }* w = u_malloc(sizeof(struct cm_invite_wrap));
w->inv = *inv; u_free(inv); w->inv = *inv; u_free(inv);
gui_bridge_post_uasync_fn( uasync_post(g_cs->inst->ua,
(void(*)(void*))cm_invite_trampoline, w); (void(*)(void*))cm_invite_trampoline, w);
} }
@ -1114,7 +1112,7 @@ static void cs_handle_channel_join(struct chat_sync* cs, uint64_t peer,
cs_refresh_channels(cs); cs_refresh_channels(cs);
{ uint8_t evt[65]; uint8_t cl=(uint8_t)strlen(ch_id); evt[0]=cl; memcpy(evt+1,ch_id,cl); { 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_MEMBERS_CHANGED, evt, 1+cl); } chat_event_post(CHAT_EVT_MEMBERS_CHANGED, evt, 1+cl); }
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: JOIN calling cs_post_channel_online ch=%s", CS_ID, ch_id); DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: JOIN calling cs_post_channel_online ch=%s", CS_ID, ch_id);
cs_post_channel_online(cs, ch_id); cs_post_channel_online(cs, ch_id);
@ -1251,9 +1249,9 @@ static void cs_handle_welcome(struct chat_sync* cs, uint64_t peer,
uint8_t evt[65]; uint8_t ch_id_len = (uint8_t)strlen(ch_id); uint8_t evt[65]; uint8_t ch_id_len = (uint8_t)strlen(ch_id);
evt[0] = ch_id_len; memcpy(evt + 1, ch_id, ch_id_len); evt[0] = ch_id_len; memcpy(evt + 1, ch_id, ch_id_len);
gui_bridge_post(GUI_EVT_CHANNEL_UPDATED, evt, 1 + ch_id_len); chat_event_post(CHAT_EVT_CHANNEL_UPDATED, evt, 1 + ch_id_len);
gui_bridge_post(GUI_EVT_MEMBERS_CHANGED, evt, 1 + ch_id_len); chat_event_post(CHAT_EVT_MEMBERS_CHANGED, evt, 1 + ch_id_len);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: WELCOME calling cs_post_channel_online ch=%s", CS_ID, ch_id); DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: WELCOME calling cs_post_channel_online ch=%s", CS_ID, ch_id);
cs_post_channel_online(cs, ch_id); cs_post_channel_online(cs, ch_id);
@ -1264,8 +1262,8 @@ static void cs_handle_welcome(struct chat_sync* cs, uint64_t peer,
if (cs->pending_invite_ch_id != 0) { if (cs->pending_invite_ch_id != 0) {
uint64_t ch_id_num = strtoull(ch_id, NULL, 10); uint64_t ch_id_num = strtoull(ch_id, NULL, 10);
uint8_t cevt[20]; memcpy(cevt, &peer, 8); int r = 0; memcpy(cevt + 8, &r, 4); memcpy(cevt + 12, &ch_id_num, 8); uint8_t cevt[20]; memcpy(cevt, &peer, 8); int r = 0; memcpy(cevt + 8, &r, 4); memcpy(cevt + 12, &ch_id_num, 8);
gui_bridge_post(GUI_EVT_CONNECT_RESULT, cevt, 20); chat_event_post(CHAT_EVT_CONNECT_RESULT, cevt, 20);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: WELCOME invite success ch=%s peer=%016llx — posted GUI_EVT_CONNECT_RESULT", DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: WELCOME invite success ch=%s peer=%016llx — posted CHAT_EVT_CONNECT_RESULT",
CS_ID, ch_id, (unsigned long long)peer); CS_ID, ch_id, (unsigned long long)peer);
cs->pending_invite_node_id = 0; cs->pending_invite_ch_id = 0; if (cs->join_timer) { uasync_cancel_timeout(cs->inst->ua, cs->join_timer); cs->join_timer = NULL; } cs->pending_invite_node_id = 0; cs->pending_invite_ch_id = 0; if (cs->join_timer) { uasync_cancel_timeout(cs->inst->ua, cs->join_timer); cs->join_timer = NULL; }
} }
@ -1346,7 +1344,7 @@ static void cs_handle_peer_upsert(struct chat_sync* cs, uint64_t peer,
cs_refresh_channels(cs); cs_refresh_channels(cs);
{ uint8_t mevt[65]; uint8_t ml = (uint8_t)strlen(ch_id); mevt[0] = ml; memcpy(mevt + 1, ch_id, ml); gui_bridge_post(GUI_EVT_MEMBERS_CHANGED, mevt, 1 + ml); } { uint8_t mevt[65]; uint8_t ml = (uint8_t)strlen(ch_id); mevt[0] = ml; memcpy(mevt + 1, ch_id, ml); chat_event_post(CHAT_EVT_MEMBERS_CHANGED, mevt, 1 + ml); }
/* propagate to others (except sender and the subject node) */ /* propagate to others (except sender and the subject node) */
cs_propagate(cs, ch_id, peer, pl, len); cs_propagate(cs, ch_id, peer, pl, len);
@ -1366,7 +1364,7 @@ static void cs_handle_peer_remove(struct chat_sync* cs, uint64_t peer,
topo_node_sqlite_member_del(db, ch_id, node_id); topo_node_sqlite_member_del(db, ch_id, node_id);
cs_refresh_channels(cs); cs_refresh_channels(cs);
{ uint8_t mevt[65]; uint8_t ml = (uint8_t)strlen(ch_id); mevt[0] = ml; memcpy(mevt + 1, ch_id, ml); gui_bridge_post(GUI_EVT_MEMBERS_CHANGED, mevt, 1 + ml); } { uint8_t mevt[65]; uint8_t ml = (uint8_t)strlen(ch_id); mevt[0] = ml; memcpy(mevt + 1, ch_id, ml); chat_event_post(CHAT_EVT_MEMBERS_CHANGED, mevt, 1 + ml); }
cs_propagate(cs, ch_id, peer, pl, len); cs_propagate(cs, ch_id, peer, pl, len);

5
tools/chatgui/transport/chat_sync.h → src/chat/chat_sync.h

@ -42,14 +42,13 @@ struct UASYNC;
/* ── Public API ── */ /* ── Public API ── */
int chat_sync_init(struct UTUN_INSTANCE* inst, int chat_sync_init(struct UTUN_INSTANCE* inst);
void (*gui_result_cb)(void*, int, const uint8_t*, int));
void chat_sync_destroy(struct UTUN_INSTANCE* inst); void chat_sync_destroy(struct UTUN_INSTANCE* inst);
/* Прямое ETCP-подключение к узлу (pubkey+адреса из БД) */ /* Прямое ETCP-подключение к узлу (pubkey+адреса из БД) */
void chat_sync_connect_node(struct UTUN_INSTANCE* inst, uint64_t node_id); void chat_sync_connect_node(struct UTUN_INSTANCE* inst, uint64_t node_id);
/* Подключиться к пиру по данным invite-ссылки (вызывается из GUI-потока) */ /* Подключиться к пиру по данным invite-ссылки (вызывается из uasync-потока) */
void chat_sync_connect_from_invite(uint64_t channel_id, uint64_t node_id, void chat_sync_connect_from_invite(uint64_t channel_id, uint64_t node_id,
const uint8_t* pubkey_bin, const uint8_t* pubkey_bin,
const uint8_t* addrs_data, int addr_count); const uint8_t* addrs_data, int addr_count);

24
src/db_sync.c → src/chat/db_sync.c

@ -1,18 +1,18 @@
// db_sync.c — Distributed content-addressed table with SQLite + peer sync via etcp_api (direct P2P) // db_sync.c — Distributed content-addressed table with SQLite + peer sync via etcp_api (direct P2P)
#include "db_sync.h" #include "db_sync.h"
#include "etcp_api.h" #include "../transport_layer/etcp_api.h"
#include "etcp.h" #include "../transport_layer/etcp.h"
#include "utun_instance.h" #include "../utun_instance.h"
#include "topo_group.h" #include "../routing_layer/topo_group.h"
#include "topo_node_sqlite.h" #include "../routing_layer/topo_node_sqlite.h"
#include "secure_channel.h" #include "../transport_layer/secure_channel.h"
#include "../lib/debug_config.h" #include "../../lib/debug_config.h"
#include "../lib/mem.h" #include "../../lib/mem.h"
#include "../lib/u_async.h" #include "../../lib/u_async.h"
#include "../lib/sha256.h" #include "../../lib/sha256.h"
#include "../lib/sqlite3.h" #include "../../lib/sqlite3.h"
#include "../lib/platform_compat.h" #include "../../lib/platform_compat.h"
#include <string.h> #include <string.h>
#include <openssl/evp.h> #include <openssl/evp.h>

0
src/db_sync.h → src/chat/db_sync.h

8
tools/chatgui/transport/member_sync.c → src/chat/member_sync.c

@ -1,10 +1,10 @@
#include "member_sync.h" #include "member_sync.h"
#include "topo_node_sqlite.h" #include "../routing_layer/topo_node_sqlite.h"
#include "chat_core.h" #include "chat_core.h"
#include "../../../src/utun_instance.h" #include "../utun_instance.h"
#include "../../../lib/debug_config.h" #include "../../lib/debug_config.h"
#include "../../../lib/mem.h" #include "../../lib/mem.h"
#include <string.h> #include <string.h>
#include <stdlib.h> #include <stdlib.h>

0
tools/chatgui/transport/member_sync.h → src/chat/member_sync.h

12
tools/chatgui/transport/merkle_sync.c → src/chat/merkle_sync.c

@ -1,11 +1,11 @@
#include "merkle_sync.h" #include "merkle_sync.h"
#include "../../../src/utun_instance.h" #include "../utun_instance.h"
#include "etcp_api.h" #include "../transport_layer/etcp_api.h"
#include "etcp.h" #include "../transport_layer/etcp.h"
#include "../../../lib/debug_config.h" #include "../../lib/debug_config.h"
#include "../../../lib/mem.h" #include "../../lib/mem.h"
#include "../../../lib/u_async.h" #include "../../lib/u_async.h"
#include <string.h> #include <string.h>
#include <stdlib.h> #include <stdlib.h>

0
tools/chatgui/transport/merkle_sync.h → src/chat/merkle_sync.h

16
src/utun_instance.c

@ -13,7 +13,9 @@
#include "etcp_connections.h" #include "etcp_connections.h"
#include "etcp.h" #include "etcp.h"
#include "conn_mgr.h" #include "conn_mgr.h"
#include "db_sync.h" #include "chat/db_sync.h"
#include "chat/chat_core.h"
#include "chat/chat_sync.h"
#include "stcp_server.h" #include "stcp_server.h"
#include "control_server.h" #include "control_server.h"
@ -418,6 +420,10 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) {
instance->etcp_sockets = NULL; instance->etcp_sockets = NULL;
DEBUG_INFO(DEBUG_CATEGORY_MEMORY, "[INSTANCE_DESTROY] ETCP sockets cleanup complete"); DEBUG_INFO(DEBUG_CATEGORY_MEMORY, "[INSTANCE_DESTROY] ETCP sockets cleanup complete");
// Cleanup chat before db_sync (chat releases db_sync instances)
chat_sync_destroy(instance);
chat_core_destroy(instance);
// Cleanup db_sync before connections cleanup (db_sync_destroy iterates connections) // Cleanup db_sync before connections cleanup (db_sync_destroy iterates connections)
db_sync_destroy(instance); db_sync_destroy(instance);
@ -583,6 +589,14 @@ int utun_instance_init(struct UTUN_INSTANCE *instance) {
// db_sync — распределённая таблица с репликацией // db_sync — распределённая таблица с репликацией
db_sync_init(instance); db_sync_init(instance);
// chat — сообщения, каналы, P2P-синхронизация
if (instance->config->global.db_path[0] && instance->config->global.db_sync_enabled) {
char db_file[512];
snprintf(db_file, sizeof(db_file), "%s/chats.db", instance->config->global.db_path);
if (chat_core_init(instance, db_file) == 0)
chat_sync_init(instance);
}
// Set TUN interface in routing module // Set TUN interface in routing module
if (instance->tun) { if (instance->tun) {
routing_set_tun(instance); routing_set_tun(instance);

4
tests/Makefile.am

@ -283,8 +283,8 @@ test_db_sync_SOURCES = test_db_sync.c
test_db_sync_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_db_sync_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib
test_db_sync_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_db_sync_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS)
test_merkle_sync_SOURCES = test_merkle_sync.c $(top_srcdir)/tools/chatgui/transport/merkle_sync.c test_merkle_sync_SOURCES = test_merkle_sync.c
test_merkle_sync_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib -I$(top_srcdir)/tools/chatgui/transport test_merkle_sync_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/src/chat -I$(top_srcdir)/lib
test_merkle_sync_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) test_merkle_sync_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS)
test_chat_sync_stress_SOURCES = test_chat_sync_stress.c test_chat_sync_stress_SOURCES = test_chat_sync_stress.c

2
tests/test_chat_sync_stress.c

@ -16,7 +16,7 @@
#include "../src/utun_instance.h" #include "../src/utun_instance.h"
#include "../src/config_parser.h" #include "../src/config_parser.h"
#include "../src/config_updater.h" #include "../src/config_updater.h"
#include "../src/db_sync.h" #include "../src/chat/db_sync.h"
#include "secure_channel.h" #include "secure_channel.h"
#include "../lib/u_async.h" #include "../lib/u_async.h"
#include "../lib/debug_config.h" #include "../lib/debug_config.h"

2
tests/test_db_sync.c

@ -23,7 +23,7 @@
#include "../src/tun_if.h" #include "../src/tun_if.h"
#include "secure_channel.h" #include "secure_channel.h"
#include "secure_channel.h" #include "secure_channel.h"
#include "../src/db_sync.h" #include "../src/chat/db_sync.h"
#include "../lib/u_async.h" #include "../lib/u_async.h"
#include "../lib/debug_config.h" #include "../lib/debug_config.h"
#include "../lib/mem.h" #include "../lib/mem.h"

10
tools/chatgui/CMakeLists.txt

@ -85,14 +85,6 @@ add_executable(chatgui
transport/node_config.cpp transport/node_config.cpp
transport/config_updater.cpp transport/config_updater.cpp
transport/gui_bridge_impl.cpp transport/gui_bridge_impl.cpp
transport/chat_core.c
transport/chat_msg.c
transport/chat_channel.c
transport/chat_profile.c
transport/chat_status.c
transport/chat_sync.c
transport/member_sync.c
transport/merkle_sync.c
transport/miniaudio_impl.c transport/miniaudio_impl.c
db/db_manager.cpp db/db_manager.cpp
../../lib/sqlite3.c ../../lib/sqlite3.c
@ -101,7 +93,7 @@ add_executable(chatgui
target_include_directories(chatgui PRIVATE ${CMAKE_SOURCE_DIR}/../../lib ${CMAKE_SOURCE_DIR}/../../src ${CMAKE_SOURCE_DIR}/db ${CMAKE_SOURCE_DIR}/../../lib/libopus/include) target_include_directories(chatgui PRIVATE ${CMAKE_SOURCE_DIR}/../../lib ${CMAKE_SOURCE_DIR}/../../src ${CMAKE_SOURCE_DIR}/db ${CMAKE_SOURCE_DIR}/../../lib/libopus/include)
target_compile_definitions(chatgui PRIVATE SQLITE_THREADSAFE=1 USE_SQLITE) target_compile_definitions(chatgui PRIVATE SQLITE_THREADSAFE=1 USE_SQLITE)
set_source_files_properties(../../lib/sqlite3.c transport/chat_core.c transport/chat_msg.c transport/chat_channel.c transport/chat_profile.c transport/chat_status.c transport/chat_sync.c transport/member_sync.c transport/merkle_sync.c transport/miniaudio_impl.c PROPERTIES LANGUAGE C) set_source_files_properties(../../lib/sqlite3.c transport/miniaudio_impl.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()

1
tools/chatgui/libutun/CMakeLists.txt

@ -68,7 +68,6 @@ target_include_directories(utun PUBLIC
${SRC_DIR}/transport_layer ${SRC_DIR}/transport_layer
${SRC_DIR}/routing_layer ${SRC_DIR}/routing_layer
${LIB_DIR} ${LIB_DIR}
${TRANSPORT_DIR}
${SRC_DIR}/uip ${SRC_DIR}/uip
${CMAKE_SOURCE_DIR}/db ${CMAKE_SOURCE_DIR}/db
) )

16
tools/chatgui/transport/utun_node.cpp

@ -22,10 +22,11 @@ extern "C" {
#include "../lib/memory_pool.h" #include "../lib/memory_pool.h"
#include "../lib/debug_config.h" #include "../lib/debug_config.h"
#include "../lib/mem.h" #include "../lib/mem.h"
#include "chat_sync.h" #include "chat/chat_sync.h"
#include "chat_core.h" #include "chat/chat_core.h"
#include "db_sync.h" #include "chat/chat_event.h"
#include "gui_bridge.h" #include "db_sync.h"
#include "gui_bridge.h"
#include "topo_node_sqlite.h" #include "topo_node_sqlite.h"
#include "topo_node.h" #include "topo_node.h"
#include "topo_group.h" #include "topo_group.h"
@ -342,9 +343,14 @@ void UtunNode::runLoop() {
etcp_router_bind(m_instance, ETCP_RT_ID_CHAT, recvCallback); etcp_router_bind(m_instance, ETCP_RT_ID_CHAT, recvCallback);
/* 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) */ /* Initialize chat_core (DB) and chat_sync (channel/message P2P sync) */
chat_core_init(m_instance, QString(m_dbPath + "/chats.db").toUtf8().constData()); chat_core_init(m_instance, QString(m_dbPath + "/chats.db").toUtf8().constData());
chat_sync_init(m_instance, nullptr); chat_sync_init(m_instance);
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "chat_core + chat_sync initialized"); DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "chat_core + chat_sync initialized");
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "utun_node: entering poll loop"); DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "utun_node: entering poll loop");

Loading…
Cancel
Save