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

878 lines
38 KiB

/*
* chat_core.c — центральный API чата в потоке uasync
*
* Все DB-операции — через своё sqlite3-соединение (тот же файл БД, WAL).
* Сетевые вызовы — напрямую в ETCP.
*/
#include "chat_core.h"
#include "db_sync.h"
#include "gui_bridge.h"
#include "../../../lib/json_flat.h"
#include "topo_node_sqlite.h"
#include "topo_group.h"
#include "member_sync.h"
#include "../../../src/utun_instance.h"
#include "etcp_api.h"
#include "etcp.h"
#include "etcp_connections.h"
#include "secure_channel.h"
#include "../../../src/ntp_time.h"
#include "../../../lib/u_async.h"
#include "../../../lib/ll_queue.h"
#include "../../../lib/mem.h"
#include "../../../lib/debug_config.h"
#include <sqlite3.h>
#include <openssl/sha.h>
#include <openssl/evp.h>
#include <string.h>
#include <stdio.h>
#include "../../../lib/platform_compat.h"
#define CC_ID "chat_core"
/* ─── глобальное состояние ─── */
static struct chat_core_ctx {
struct UTUN_INSTANCE* inst;
sqlite3* db;
uint8_t shared_db;
uint64_t my_node_id;
/* db_sync instances per channel */
struct DB_SYNC_INSTANCE** si;
char** si_ch_id;
int si_count, si_capacity;
uint8_t initialized;
} 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_msg_inserted(struct DB_SYNC_INSTANCE* si, uint64_t record_ts, 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) {
size_t i = 0;
while (*ch_id && i < out_sz - 1) {
char c = *ch_id++;
if ((c >= 'a' && c <= 'z') || (c >= 'A' && c <= 'Z')
|| (c >= '0' && c <= '9') || c == '_')
out[i++] = c;
else
out[i++] = '_';
}
out[i] = '\0';
}
static void msg_table_name(const char* ch_id, char* buf, size_t sz) {
char san[64]; sanitize_ch_id(ch_id, san, sizeof(san));
snprintf(buf, sz, "msg_%s", san);
}
static void peers_table_name(const char* ch_id, char* buf, size_t sz) {
char san[64]; sanitize_ch_id(ch_id, san, sizeof(san));
snprintf(buf, sz, "peers_%s", san);
}
static int db_exec(const char* sql) {
char* err = NULL;
int rc = sqlite3_exec(g_cc.db, sql, NULL, NULL, &err);
if (rc != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: sql error: %s", CC_ID, err);
sqlite3_free(err);
}
return rc;
}
/* ─── жизненный цикл ─── */
int chat_core_init(struct UTUN_INSTANCE* inst, const char* db_path) {
if (!inst || !db_path) return -1;
memset(&g_cc, 0, sizeof(g_cc));
if (inst->topo_sqlite_db) {
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);
} else {
int rc = sqlite3_open_v2(db_path, &g_cc.db,
SQLITE_OPEN_READWRITE | SQLITE_OPEN_CREATE | SQLITE_OPEN_FULLMUTEX, NULL);
if (rc != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: cannot open DB %s: %s",
CC_ID, db_path, sqlite3_errmsg(g_cc.db));
sqlite3_close(g_cc.db); g_cc.db = NULL; return -1;
}
sqlite3_exec(g_cc.db, "PRAGMA journal_mode=WAL", NULL, NULL, NULL);
sqlite3_exec(g_cc.db, "PRAGMA foreign_keys=ON", NULL, NULL, NULL);
}
g_cc.inst = inst;
g_cc.my_node_id = inst->node_id;
sqlite3_int64 now_sec = (sqlite3_int64)ntp_time_get_seconds(inst);
/* записать себя в nodes + local_identity */
{
const char* my_name = inst->name[0] ? inst->name : "Me";
topo_node_sqlite_node_update_verified(g_cc.db, g_cc.my_node_id,
my_name, inst->my_keys.public_key, inst->my_ed25519_pubkey,
(uint64_t)now_sec, (time_t)now_sec);
topo_node_sqlite_node_set_online(g_cc.db, g_cc.my_node_id, 1);
sqlite3_stmt* st = NULL;
sqlite3_prepare_v2(g_cc.db,
"INSERT OR REPLACE INTO local_identity(id,node_id,name,x25519_pubkey,ed25519_pubkey,created_at,updated_at)"
" VALUES(1,?,?,?,?,?,?)", -1, &st, NULL);
if (st) {
sqlite3_bind_int64(st, 1, (sqlite3_int64)g_cc.my_node_id);
sqlite3_bind_text(st, 2, my_name, -1, SQLITE_STATIC);
sqlite3_bind_blob(st, 3, inst->my_keys.public_key, 32, SQLITE_STATIC);
sqlite3_bind_blob(st, 4, inst->my_ed25519_pubkey, 32, SQLITE_STATIC);
sqlite3_bind_int64(st, 5, now_sec);
sqlite3_bind_int64(st, 6, now_sec);
sqlite3_step(st); sqlite3_finalize(st);
}
}
g_cc.initialized = 1;
chat_core_sync_my_addresses();
db_exec(
"CREATE TABLE IF NOT EXISTS local_identity ("
" id INTEGER PRIMARY KEY CHECK (id = 1),"
" node_id INTEGER NOT NULL UNIQUE,"
" name TEXT NOT NULL,"
" x25519_pubkey BLOB NOT NULL,"
" x25519_privkey BLOB,"
" ed25519_pubkey BLOB,"
" created_at INTEGER DEFAULT (unixepoch()),"
" updated_at INTEGER DEFAULT (unixepoch())"
");"
"CREATE TABLE IF NOT EXISTS accounts ("
" node_id INTEGER PRIMARY KEY REFERENCES nodes(node_id),"
" display_name TEXT NOT NULL,"
" avatar_color TEXT DEFAULT '#4A90E2',"
" avatar_letter TEXT NOT NULL,"
" is_contact INTEGER DEFAULT 1,"
" created_at INTEGER DEFAULT (unixepoch())"
");"
"CREATE TABLE IF NOT EXISTS ui_state ("
" key TEXT PRIMARY KEY,"
" value TEXT"
");"
);
/* seed stub accounts if empty */
{
sqlite3_stmt* cnt_st = NULL;
if (sqlite3_prepare_v2(g_cc.db, "SELECT COUNT(*) FROM accounts",
-1, &cnt_st, NULL) == SQLITE_OK) {
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);
struct { int64_t id; const char* name; const char* letter; const char* color; } accs[] = {
{1, "Alice", "A", "#4A90E2"}, {2, "Bob", "B", "#7ED321"},
{3, "Carol", "C", "#F5A623"}, {4, "Dave", "D", "#9013FE"},
{5, "Eve", "E", "#00B894"},
};
for (int i = 0; i < 5; i++) {
sqlite3_stmt* ins = NULL;
sqlite3_prepare_v2(g_cc.db,
"INSERT OR IGNORE INTO accounts(node_id,display_name,avatar_color,avatar_letter,created_at)"
" VALUES(?,?,?,?,?)", -1, &ins, NULL);
if (ins) {
sqlite3_bind_int64(ins, 1, accs[i].id);
sqlite3_bind_text(ins, 2, accs[i].name, -1, SQLITE_STATIC);
sqlite3_bind_text(ins, 3, accs[i].color, -1, SQLITE_STATIC);
sqlite3_bind_text(ins, 4, accs[i].letter, -1, SQLITE_STATIC);
sqlite3_bind_int64(ins, 5, now_sec);
sqlite3_step(ins); sqlite3_finalize(ins);
}
}
}
sqlite3_finalize(cnt_st);
}
}
/* ensure infrastructure for all existing channels (survives restart) */
{
sqlite3_stmt* stmt = NULL;
if (sqlite3_prepare_v2(g_cc.db, "SELECT channel_id FROM channels ORDER BY created_at ASC",
-1, &stmt, NULL) == SQLITE_OK) {
int loaded = 0;
while (sqlite3_step(stmt) == SQLITE_ROW) {
const char* ch = (const char*)sqlite3_column_text(stmt, 0);
if (ch && ch[0]) { chat_core_ensure_channel_ready(ch); loaded++; }
}
sqlite3_finalize(stmt);
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: loaded %d channels from DB", CC_ID, loaded);
}
}
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: initialized, db=%s node_id=0x%016llx",
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); }
/* notify GUI: DB is ready (tables created, data seeded) */
gui_bridge_post(GUI_EVT_DB_READY, NULL, 0);
return 0;
}
sqlite3* chat_core_get_db(void) {
return g_cc.db;
}
struct UTUN_INSTANCE* chat_core_get_inst(void) {
return g_cc.inst;
}
int chat_core_is_initialized(void) {
return g_cc.initialized;
}
void chat_core_destroy(struct UTUN_INSTANCE* inst) {
(void)inst;
if (!g_cc.initialized) return;
g_cc.initialized = 0;
if (g_cc.db && !g_cc.shared_db) { sqlite3_close(g_cc.db); }
g_cc.db = NULL; g_cc.inst = NULL;
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: destroyed", CC_ID);
}
void chat_core_set_my_node_id(uint64_t node_id) {
g_cc.my_node_id = node_id;
}
void chat_core_update_my_name(const char* name) {
if (!g_cc.initialized || !name) return;
topo_node_sqlite_node_update_verified(g_cc.db, g_cc.my_node_id,
name, g_cc.inst->my_keys.public_key, g_cc.inst->my_ed25519_pubkey,
(uint64_t)ntp_time_get_seconds(g_cc.inst), 0);
snprintf(g_cc.inst->name, sizeof(g_cc.inst->name), "%s", name);
/* recompute join_sigs for all channels where I'm a member */
uint64_t myid = g_cc.my_node_id;
sqlite3_stmt* cs = NULL;
sqlite3_prepare_v2(g_cc.db, "SELECT channel_id, x25519_pubkey, ed25519_pubkey FROM channels", -1, &cs, NULL);
if (cs) {
while (sqlite3_step(cs) == SQLITE_ROW) {
const char* ch = (const char*)sqlite3_column_text(cs, 0);
const uint8_t* ch_x25519 = sqlite3_column_blob(cs, 1);
const uint8_t* ch_ed25519 = sqlite3_column_blob(cs, 2);
if (!ch || !ch_x25519 || !ch_ed25519) continue;
char tbl[80]; peers_table_name(ch, tbl, sizeof(tbl));
char buf[256]; snprintf(buf, sizeof(buf), "SELECT 1 FROM \"%s\" WHERE node_id=?", tbl);
sqlite3_stmt* ps = NULL;
if (sqlite3_prepare_v2(g_cc.db, buf, -1, &ps, NULL) == SQLITE_OK) {
sqlite3_bind_int64(ps, 1, (sqlite3_int64)myid);
if (sqlite3_step(ps) == SQLITE_ROW) {
uint8_t join_msg[256]; size_t mlen = 0;
memcpy(join_msg + mlen, ch_x25519, 32); mlen += 32;
memcpy(join_msg + mlen, ch_ed25519, 32); mlen += 32;
memcpy(join_msg + mlen, &myid, 8); mlen += 8;
memcpy(join_msg + mlen, g_cc.inst->my_keys.public_key, 32); mlen += 32;
uint64_t join_ts = (uint64_t)ntp_time_get_seconds(g_cc.inst);
memcpy(join_msg + mlen, &join_ts, 8); mlen += 8;
uint8_t new_sig[64]; memset(new_sig, 0, 64);
EVP_PKEY* pkey = EVP_PKEY_new_raw_private_key(EVP_PKEY_ED25519, NULL,
g_cc.inst->my_ed25519_privkey, 32);
if (pkey) {
EVP_MD_CTX* mdctx = EVP_MD_CTX_new();
if (mdctx) {
if (EVP_DigestSignInit(mdctx, NULL, NULL, NULL, pkey) == 1)
EVP_DigestSign(mdctx, new_sig, &(size_t){64}, join_msg, mlen);
EVP_MD_CTX_free(mdctx);
}
EVP_PKEY_free(pkey);
}
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,
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);
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);
if (r2 != 0) DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%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);
evt[0] = cl; memcpy(evt + 1, ch, cl);
gui_bridge_post(GUI_EVT_MEMBERS_CHANGED, evt, 1 + cl);
}
sqlite3_finalize(ps);
}
}
sqlite3_finalize(cs);
}
}
void chat_core_sync_my_addresses(void) {
if (!g_cc.initialized || !g_cc.db || !g_cc.inst) return;
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] sync_my_addresses START my_node=0x%016llx", CC_ID, (unsigned long long)g_cc.my_node_id);
sqlite3_stmt* del = NULL;
sqlite3_prepare_v2(g_cc.db, "DELETE FROM node_addresses WHERE node_id=? AND addr_type=0", -1, &del, NULL);
if (del) {
sqlite3_bind_int64(del, 1, (sqlite3_int64)g_cc.my_node_id);
sqlite3_step(del);
int deleted = sqlite3_changes(g_cc.db);
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] DELETE addr_type=0 for my_node=0x%016llx: %d rows deleted", CC_ID, (unsigned long long)g_cc.my_node_id, deleted);
sqlite3_finalize(del);
}
sqlite3_stmt* ins = NULL;
sqlite3_prepare_v2(g_cc.db,
"INSERT OR REPLACE INTO node_addresses(node_id,family,protocol,address,port,addr_type,socket_id)"
" VALUES(?,?,1,?,?,0,?)", -1, &ins, NULL);
if (!ins) {
DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] FAIL prepare INSERT stmt my_node=0x%016llx", CC_ID, (unsigned long long)g_cc.my_node_id);
return;
}
struct ETCP_SOCKET* sock = g_cc.inst->etcp_sockets;
int sock_count = 0, written = 0;
while (sock) {
sock_count++;
struct sockaddr_storage* sa = sock->interface_addr.ss_family ? &sock->interface_addr : NULL;
if (!sa) sa = sock->local_addr.ss_family ? &sock->local_addr : NULL;
if (!sa) { DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] socket %d — no interface/local addr", CC_ID, sock_count); sock = sock->next; continue; }
if (sa->ss_family == AF_INET) {
struct sockaddr_in* sin = (struct sockaddr_in*)sa;
sqlite3_bind_int64(ins, 1, (sqlite3_int64)g_cc.my_node_id);
sqlite3_bind_int(ins, 2, 4);
sqlite3_bind_blob(ins, 3, &sin->sin_addr, 4, SQLITE_STATIC);
sqlite3_bind_int(ins, 4, (int)ntohs(sin->sin_port));
sqlite3_bind_int(ins, 5, (int)sock->sock_id);
sqlite3_step(ins); sqlite3_reset(ins);
uint8_t* ip = (uint8_t*)&sin->sin_addr;
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] INSERT my_node=0x%016llx sock=%d %d.%d.%d.%d:%d",
CC_ID, (unsigned long long)g_cc.my_node_id, sock->sock_id, ip[0], ip[1], ip[2], ip[3], (int)ntohs(sin->sin_port));
written++;
} else if (sa->ss_family == AF_INET6) {
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)sa;
sqlite3_bind_int64(ins, 1, (sqlite3_int64)g_cc.my_node_id);
sqlite3_bind_int(ins, 2, 6);
sqlite3_bind_blob(ins, 3, &sin6->sin6_addr, 16, SQLITE_STATIC);
sqlite3_bind_int(ins, 4, (int)ntohs(sin6->sin6_port));
sqlite3_bind_int(ins, 5, (int)sock->sock_id);
sqlite3_step(ins); sqlite3_reset(ins);
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] INSERT v6 my_node=0x%016llx sock=%d port=%d",
CC_ID, (unsigned long long)g_cc.my_node_id, sock->sock_id, (int)ntohs(sin6->sin6_port));
written++;
}
sock = sock->next;
}
sqlite3_finalize(ins);
if (sock_count == 0)
DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] etcp_sockets is NULL — NO sockets, addresses NOT written!", CC_ID);
else if (written == 0)
DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] %d sockets found but 0 addresses written (all have no interface_addr nor local_addr)", CC_ID, sock_count);
else
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] DONE: %d of %d sockets written, my_node=0x%016llx",
CC_ID, written, sock_count, (unsigned long long)g_cc.my_node_id);
/* verify */
sqlite3_stmt* chk = NULL;
sqlite3_prepare_v2(g_cc.db, "SELECT COUNT(*) FROM node_addresses WHERE node_id=? AND addr_type=0", -1, &chk, NULL);
if (chk) {
sqlite3_bind_int64(chk, 1, (sqlite3_int64)g_cc.my_node_id);
int cnt = sqlite3_step(chk) == SQLITE_ROW ? sqlite3_column_int(chk, 0) : -1;
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] VERIFY: total addr_type=0 rows for my_node=0x%016llx = %d",
CC_ID, (unsigned long long)g_cc.my_node_id, cnt);
sqlite3_finalize(chk);
}
}
void chat_core_update_my_name_trampoline(void* arg) {
chat_core_update_my_name((const char*)arg);
u_free(arg);
}
void chat_core_save_ui_state(const char* key, const char* value) {
if (!g_cc.initialized || !key || !value) return;
sqlite3_stmt* st = NULL;
sqlite3_prepare_v2(g_cc.db,
"INSERT OR REPLACE INTO ui_state(key, value) VALUES(?, ?)", -1, &st, NULL);
if (st) {
sqlite3_bind_text(st, 1, key, -1, SQLITE_STATIC);
sqlite3_bind_text(st, 2, value, -1, SQLITE_STATIC);
sqlite3_step(st); sqlite3_finalize(st);
}
}
struct save_ui_state_arg { char data[256]; };
void chat_core_save_ui_state_trampoline(void* arg) {
struct save_ui_state_arg* a = (struct save_ui_state_arg*)arg;
const char* key = a->data;
const char* value = key + strlen(key) + 1;
chat_core_save_ui_state(key, value);
u_free(arg);
}
static void chat_core_collect_status(void);
void chat_core_collect_status_trampoline(void* arg) {
(void)arg;
chat_core_collect_status();
}
/* ─── отправка сообщения (GUI → uasync) ─── */
static int msg_get_prev_chain_hash(const char* tbl, uint8_t out[32]) {
char sql[128]; snprintf(sql, sizeof(sql),
"SELECT chain_hash FROM \"%s\" ORDER BY timestamp DESC LIMIT 1", tbl);
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(g_cc.db, sql, -1, &st, NULL) != SQLITE_OK) { memset(out,0,32); return -1; }
if (sqlite3_step(st) == SQLITE_ROW) {
const void* blob = sqlite3_column_blob(st, 0); int blen = sqlite3_column_bytes(st, 0);
if (blob && blen == 32) memcpy(out, blob, 32); else memset(out, 0, 32);
} else { memset(out, 0, 32); }
sqlite3_finalize(st); return 0;
}
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 (!req) return;
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: SUBMIT ch=%s ct=%s len=%u",
CC_ID, req->channel_id, req->content_type, req->data_len);
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; }
/* Build JSON: {"n":<node_id>,"ch":"<ch_id>","ct":"<>","d":"<>"} */
char json[4096];
snprintf(json, sizeof(json),
"{\"n\":%llu,\"ch\":\"%s\",\"ct\":\"%s\",\"d\":\"%.*s\"}",
(unsigned long long)g_cc.my_node_id, req->channel_id,
req->content_type, (int)req->data_len, (const char*)req->data);
/* Generate Ed25519 signature: ts || json */
uint64_t ts = db_sync_next_timestamp(si);
uint8_t sig_msg[8192]; size_t soff = 0;
memcpy(sig_msg + soff, &ts, 8); soff += 8;
size_t jl = strlen(json); memcpy(sig_msg + soff, json, jl); soff += jl;
uint8_t sig[64];
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);
return;
}
int ret = db_sync_insert_signed(si, json, strlen(json), sig, 64, ts);
if (ret != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: db_sync_insert_signed failed ret=%d", CC_ID, ret);
return;
}
/* notify GUI */
{ uint8_t evt[65]; uint8_t cl=(uint8_t)strlen(req->channel_id); evt[0]=cl; memcpy(evt+1,req->channel_id,cl);
gui_bridge_post(GUI_EVT_MSG_RECEIVED, evt, 1+cl); }
}
void chat_core_submit_trampoline(void* arg) { chat_core_submit_message((struct chat_msg_submit*)arg); u_free(arg); }
/* ─── DB-операции для chat_sync ─── */
uint32_t chat_core_count(const char* ch_id) {
if (!g_cc.initialized) return 0;
char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl));
char sql[128]; snprintf(sql, sizeof(sql),
"SELECT COUNT(*) FROM \"%s\"", tbl);
sqlite3_stmt* stmt = NULL;
if (sqlite3_prepare_v2(g_cc.db, sql, -1, &stmt, NULL) != SQLITE_OK) return 0;
uint32_t cnt = 0;
if (sqlite3_step(stmt) == SQLITE_ROW)
cnt = (uint32_t)sqlite3_column_int64(stmt, 0);
sqlite3_finalize(stmt);
return cnt;
}
int chat_core_chain_hash_at(const char* ch_id, uint32_t pos, uint8_t* hash_out) {
if (!g_cc.initialized || !hash_out) return -1;
char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl));
char sql[200]; snprintf(sql, sizeof(sql),
"SELECT chain_hash FROM \"%s\" ORDER BY timestamp, node_id ASC LIMIT 1 OFFSET ?",
tbl);
sqlite3_stmt* stmt = NULL;
if (sqlite3_prepare_v2(g_cc.db, sql, -1, &stmt, NULL) != SQLITE_OK) {
memset(hash_out, 0, 32); return -1;
}
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)pos);
if (sqlite3_step(stmt) == SQLITE_ROW) {
const void* blob = sqlite3_column_blob(stmt, 0);
int len = sqlite3_column_bytes(stmt, 0);
if (blob && len == 32) memcpy(hash_out, blob, 32);
else memset(hash_out, 0, 32);
} else {
memset(hash_out, 0, 32);
}
sqlite3_finalize(stmt);
return 0;
}
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;
sqlite3_stmt* stmt = NULL;
if (sqlite3_prepare_v2(g_cc.db,
"SELECT channel_id FROM channels ORDER BY created_at ASC",
-1, &stmt, NULL) != SQLITE_OK) return -1;
uint8_t* out = buf;
uint8_t* start = out;
out += 2; /* placeholder for count */
uint16_t cnt = 0;
while (sqlite3_step(stmt) == SQLITE_ROW) {
const char* ch_id = (const char*)sqlite3_column_text(stmt, 0);
int ch_len = sqlite3_column_bytes(stmt, 0);
if (!ch_id || ch_len <= 0 || ch_len > 63) continue;
size_t need = (size_t)(out - buf) + 1 + (size_t)ch_len;
if (need > buf_size) break;
*out++ = (uint8_t)ch_len;
memcpy(out, ch_id, (size_t)ch_len); out += ch_len;
cnt++;
}
sqlite3_finalize(stmt);
memcpy(start, &cnt, 2);
*out_len = (size_t)(out - buf);
return 0;
}
int chat_core_list_peers(const char* ch_id, uint8_t* buf, size_t buf_size,
size_t* out_len) {
if (!g_cc.initialized || !buf || !out_len) return -1;
char tbl[80]; peers_table_name(ch_id, tbl, sizeof(tbl));
char sql[128]; snprintf(sql, sizeof(sql),
"SELECT node_id FROM \"%s\"", tbl);
sqlite3_stmt* stmt = NULL;
if (sqlite3_prepare_v2(g_cc.db, sql, -1, &stmt, NULL) != SQLITE_OK) {
uint16_t z = 0; memcpy(buf, &z, 2); *out_len = 2; return -1;
}
uint16_t cnt = 0;
uint8_t* out = buf + 2;
while (sqlite3_step(stmt) == SQLITE_ROW) {
if ((size_t)(out - buf) + 8 > buf_size) break;
uint64_t nid = (uint64_t)sqlite3_column_int64(stmt, 0);
memcpy(out, &nid, 8); out += 8; cnt++;
}
sqlite3_finalize(stmt);
memcpy(buf, &cnt, 2);
*out_len = (size_t)(out - buf);
return 0;
}
int chat_core_load_nodeinfo(uint64_t node_id, uint8_t* buf, size_t buf_size,
size_t* out_len) {
if (!g_cc.initialized || !buf || !out_len) return -1;
sqlite3_stmt* stmt = NULL;
if (sqlite3_prepare_v2(g_cc.db,
"SELECT x25519_pubkey, ed25519_pubkey FROM nodes WHERE node_id=?",
-1, &stmt, NULL) != SQLITE_OK) return -1;
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)node_id);
if (sqlite3_step(stmt) != SQLITE_ROW) { sqlite3_finalize(stmt); return -1; }
const uint8_t* x25519 = (const uint8_t*)sqlite3_column_blob(stmt, 0);
int x25519_len = sqlite3_column_bytes(stmt, 0);
const uint8_t* ed25519 = (const uint8_t*)sqlite3_column_blob(stmt, 1);
int ed25519_len = sqlite3_column_bytes(stmt, 1);
uint8_t* out = buf;
if (x25519 && x25519_len == 32) memcpy(out, x25519, 32);
else memset(out, 0, 32);
out += 32;
if (ed25519 && ed25519_len == 32) memcpy(out, ed25519, 32);
else memset(out, 0, 32);
out += 32;
sqlite3_finalize(stmt);
/* адреса */
if (sqlite3_prepare_v2(g_cc.db,
"SELECT family, protocol, address, port, rtt FROM node_addresses WHERE node_id=?",
-1, &stmt, NULL) != SQLITE_OK) {
*out++ = 0; *out_len = (size_t)(out - buf); return 0;
}
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)node_id);
uint8_t* cnt_pos = out; *out++ = 0;
uint8_t addr_cnt = 0;
while (sqlite3_step(stmt) == SQLITE_ROW && addr_cnt < 255) {
int family = sqlite3_column_int(stmt, 0);
int proto = sqlite3_column_int(stmt, 1);
const uint8_t* addr = (const uint8_t*)sqlite3_column_blob(stmt, 2);
int addr_len = sqlite3_column_bytes(stmt, 2);
uint16_t port = (uint16_t)sqlite3_column_int(stmt, 3);
int16_t rtt = (int16_t)sqlite3_column_int(stmt, 4);
if ((size_t)(out - buf) + 7 + (size_t)addr_len > buf_size) break;
*out++ = (uint8_t)family;
*out++ = (uint8_t)proto;
*out++ = (uint8_t)addr_len;
if (addr_len > 0) { memcpy(out, addr, (size_t)addr_len); out += addr_len; }
memcpy(out, &port, 2); out += 2;
memcpy(out, &rtt, 2); out += 2;
addr_cnt++;
}
*cnt_pos = addr_cnt;
sqlite3_finalize(stmt);
*out_len = (size_t)(out - buf);
return 0;
}
/* ─── подготовка инфраструктуры канала (db_sync instance) ─── */
void chat_core_ensure_channel_ready(const char* ch_id) {
if (!g_cc.initialized || !ch_id || !ch_id[0]) return;
if (si_find(ch_id)) return;
char tbl_msg[80]; msg_table_name(ch_id, tbl_msg, sizeof(tbl_msg));
uint64_t ch_hash = 0;
{ const uint8_t* chd = (const uint8_t*)ch_id; size_t chl = strlen(ch_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, tbl_msg, ch_hash);
/* create TOPO_GROUP_TYPE_CHAT for auto-connect */
uint64_t gid = strtoull(ch_id, NULL, 10);
if (gid != 0 && g_cc.inst->topo_groups && !topo_groups_find(g_cc.inst->topo_groups, gid))
topo_groups_create_group(g_cc.inst->topo_groups, gid, TOPO_GROUP_TYPE_CHAT, ch_id);
if (si) {
si_register(si, 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",
CC_ID, ch_id, tbl_msg, (unsigned long long)ch_hash);
} else {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: db_sync_instance_add failed for ch=%s",
CC_ID, ch_id);
}
}
/* ─── создание канала ─── */
void chat_core_create_channel(struct chat_channel_create* req) {
if (!g_cc.initialized) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: create_channel NOT INITIALIZED ch=%s",
CC_ID, req ? req->channel_id : "(null)");
return;
}
if (!req) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: create_channel req=NULL", CC_ID); return; }
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: create_channel BEGIN ch=%s name=%s",
CC_ID, req->channel_id, req->name);
/* записываем канал в БД ПЕРВЫМ — создаст полную схему peers_<ch_id> (11 колонок) */
int rc = topo_node_sqlite_channel_put(g_cc.db,
req->channel_id, req->name, req->owner_node_id,
req->x25519_pubkey, req->x25519_privkey,
req->ed25519_pubkey, req->ed25519_privkey,
req->signature);
if (rc != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: channel_put FAILED ch=%s rc=%d",
CC_ID, req->channel_id, rc);
} else {
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%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);
/* после полной схемы — готовим таблицы и db_sync */
chat_core_ensure_channel_ready(req->channel_id);
uint8_t ch_id_len = (uint8_t)strlen(req->channel_id);
uint8_t data[65];
data[0] = ch_id_len;
memcpy(data + 1, req->channel_id, ch_id_len);
/* add self as first member BEFORE posting channel to GUI */
{
uint64_t myid = g_cc.inst->node_id;
uint8_t join_msg[256]; size_t mlen = 0;
memcpy(join_msg + mlen, req->x25519_pubkey, 32); mlen += 32;
memcpy(join_msg + mlen, req->ed25519_pubkey, 32); mlen += 32;
memcpy(join_msg + mlen, &myid, 8); mlen += 8;
memcpy(join_msg + mlen, g_cc.inst->my_keys.public_key, 32); mlen += 32;
uint64_t join_ts = (uint64_t)ntp_time_get_seconds(g_cc.inst);
memcpy(join_msg + mlen, &join_ts, 8); mlen += 8;
uint8_t join_sig[64];
EVP_PKEY* pkey = EVP_PKEY_new_raw_private_key(EVP_PKEY_ED25519, NULL,
g_cc.inst->my_ed25519_privkey, 32);
if (pkey) {
EVP_MD_CTX* mdctx = EVP_MD_CTX_new();
if (mdctx) {
if (EVP_DigestSignInit(mdctx, NULL, NULL, NULL, pkey) == 1)
EVP_DigestSign(mdctx, join_sig, &(size_t){64}, join_msg, mlen);
else
memset(join_sig, 0, 64);
EVP_MD_CTX_free(mdctx);
}
EVP_PKEY_free(pkey);
} else {
memset(join_sig, 0, 64);
}
uint8_t my_addrs[256]; int my_addr_cnt = 0;
{ struct ETCP_SOCKET* s = g_cc.inst->etcp_sockets;
while (s && my_addr_cnt < 16) {
struct sockaddr_storage* sa = s->interface_addr.ss_family ? &s->interface_addr : NULL;
if (!sa) sa = s->local_addr.ss_family ? &s->local_addr : NULL;
if (sa && sa->ss_family == AF_INET) {
struct sockaddr_in* sin = (struct sockaddr_in*)sa;
my_addrs[my_addr_cnt * 8] = 4;
my_addrs[my_addr_cnt * 8 + 1] = s->sock_id;
memcpy(my_addrs + my_addr_cnt * 8 + 2, &sin->sin_addr, 4);
uint16_t port = ntohs(sin->sin_port);
my_addrs[my_addr_cnt * 8 + 6] = (uint8_t)(port >> 8);
my_addrs[my_addr_cnt * 8 + 7] = (uint8_t)(port & 0xFF);
my_addr_cnt++;
}
s = s->next;
}
}
char juser3[256]; snprintf(juser3, sizeof(juser3), "{\"name\":\"%s\"}",
g_cc.inst->name[0] ? g_cc.inst->name : "");
int mrc = member_sync_put(g_cc.inst, req->channel_id, myid,
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);
if (mrc != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: member_sync_put(self) FAILED ch=%s rc=%d",
CC_ID, req->channel_id, mrc);
}
}
/* notify GUI — members already in DB, channel will show with self as participant */
gui_bridge_post(GUI_EVT_CHANNEL_UPDATED, data, 1 + ch_id_len);
}
}
void chat_core_create_channel_trampoline(void* arg) {
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "chat_core: TRAMPOLINE invoked arg=%p", arg);
struct chat_channel_create* req = (struct chat_channel_create*)arg;
chat_core_create_channel(req);
u_free(req);
}
/* ─── db_sync helpers ─── */
static void on_msg_inserted(struct DB_SYNC_INSTANCE* si, uint64_t record_ts, const char* data, size_t len, uint64_t author, void* arg) {
(void)si; (void)record_ts; (void)data; (void)len; (void)author;
const char* ch_id = (const char*)arg;
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);
}
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++;
uint32_t msg_cnt = db_sync_count(si);
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: si_register ch=%s mc=%u", CC_ID, ch_id, msg_cnt);
}
/* ─── сбор статуса (NTP + connections) и отправка в GUI ─── */
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 "?"; }
}
static void get_node_name(uint64_t node_id, char* out, size_t sz) {
out[0] = '\0';
if (!g_cc.db || node_id == 0) return;
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(g_cc.db, "SELECT name FROM nodes WHERE node_id=?", -1, &st, NULL) != SQLITE_OK) return;
sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id);
if (sqlite3_step(st) == SQLITE_ROW) {
const char* n = (const char*)sqlite3_column_text(st, 0);
if (n) snprintf(out, sz, "%s", n);
}
sqlite3_finalize(st);
}
static void chat_core_collect_status(void) {
char buf[8192]; int off = 0;
if (!g_cc.initialized || !g_cc.inst) {
off = snprintf(buf, sizeof(buf), "uTun not initialized\n");
gui_bridge_post(GUI_EVT_STATUS_REFRESH, (const uint8_t*)buf, off);
return;
}
/* ─── NTP ─── */
struct NTP_TIME* ntp = &g_cc.inst->ntp;
off += snprintf(buf + off, sizeof(buf) - off, "=== NTP ===\n");
off += snprintf(buf + off, sizeof(buf) - off, "Enabled: %s\n", ntp->enabled ? "yes" : "no");
off += snprintf(buf + off, sizeof(buf) - off, "Synced: %s\n", ntp->synced ? "yes" : "no");
off += snprintf(buf + off, sizeof(buf) - off, "Offset: %.1f s\n", ntp->offset_us / 1000000.0);
if (ntp->server_count > 0 && ntp->servers) {
off += snprintf(buf + off, sizeof(buf) - off, "Servers: %d\n", ntp->server_count);
for (int i = 0; i < ntp->server_count && ntp->servers[i]; i++)
off += snprintf(buf + off, sizeof(buf) - off, " [%d] %s%s\n", i, ntp->servers[i],
i == ntp->server_current ? " (current)" : "");
}
if (ntp->synced && ntp->last_sync_tb > 0) {
uint64_t now_tb = get_time_tb();
uint64_t ago_sec = (now_tb - ntp->last_sync_tb) / 10000;
off += snprintf(buf + off, sizeof(buf) - off, "Last sync: %llu sec ago\n", (unsigned long long)ago_sec);
}
/* ─── Connections ─── */
int conn_count = 0;
struct ll_entry* entry = g_cc.inst->connections->head;
while (entry) { conn_count++; entry = entry->next; }
off += snprintf(buf + off, sizeof(buf) - off, "=== Connections (%d) ===\n", conn_count);
uint64_t my_id = g_cc.inst->node_id;
entry = g_cc.inst->connections->head;
while (entry) {
struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data;
if (!ce || !ce->conn) { entry = entry->next; continue; }
uint64_t pid = ce->peer_node_id;
struct ETCP_CONN* conn = ce->conn;
int link_count = 0;
struct ETCP_LINK* tl = conn->links;
while (tl) { link_count++; tl = tl->next; }
const char* init_str = conn->initialized ? "" : "(!init)";
char myhex[5], peerhex[5];
snprintf(myhex, sizeof(myhex), "%04llX", (unsigned long long)(my_id & 0xFFFF));
snprintf(peerhex, sizeof(peerhex), "%04llX", (unsigned long long)(pid & 0xFFFF));
char peername[64];
get_node_name(pid, peername, sizeof(peername));
if (peername[0])
off += snprintf(buf + off, sizeof(buf) - off, "[%s]→[%s] %s ETCP:%s%s(%dL)\n",
myhex, peerhex, peername,
conn->links_up ? "UP" : "DOWN", init_str, link_count);
else
off += snprintf(buf + off, sizeof(buf) - off, "[%s]→[%s] ETCP:%s%s(%dL)\n",
myhex, peerhex,
conn->links_up ? "UP" : "DOWN", init_str, link_count);
int link_idx = 0;
struct ETCP_LINK* link = conn->links;
while (link) {
char rtt_str[32]; rtt_str[0] = '\0';
if (link->rtt_last > 0) snprintf(rtt_str, sizeof(rtt_str), " / rtt=%ums", link->rtt_last);
off += snprintf(buf + off, sizeof(buf) - off, " LINK#%d: %s /NAT=%s%s\n",
link_idx,
link->link_status ? "UP" : "DOWN",
nat_type_str(link->nat_type),
rtt_str);
link = link->next; link_idx++;
}
entry = entry->next;
}
gui_bridge_post(GUI_EVT_STATUS_REFRESH, (const uint8_t*)buf, off);
}