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.
1172 lines
50 KiB
1172 lines
50 KiB
/* |
|
* chat_core.c — центральный API чата в потоке uasync |
|
* |
|
* Все DB-операции — через своё sqlite3-соединение (тот же файл БД, WAL). |
|
* Сетевые вызовы — напрямую в ETCP. |
|
*/ |
|
|
|
#include "chat_core.h" |
|
#include "db_sync.h" |
|
#include "gui_bridge.h" |
|
#include "topo_node_sqlite.h" |
|
#include "member_sync.h" |
|
|
|
#include "../../../src/utun_instance.h" |
|
#include "../../../src/etcp_api.h" |
|
#include "../../../src/topo_group.h" |
|
#include "../../../src/topo_node.h" |
|
#include "../../../src/conn_mgr.h" |
|
#include "../../../src/etcp.h" |
|
#include "../../../src/etcp_connections.h" |
|
#include "../../../src/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> |
|
|
|
#define CC_ID "chat_core" |
|
|
|
/* ─── глобальное состояние ─── */ |
|
|
|
static struct chat_core_ctx { |
|
struct UTUN_INSTANCE* inst; |
|
sqlite3* 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_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) { |
|
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 void compute_datahash(const uint8_t* data, size_t len, uint64_t* dh) { |
|
uint8_t hash[32]; SHA256(data, len, hash); |
|
memcpy(dh, hash, 8); |
|
} |
|
|
|
static void compute_chain_hash(const uint8_t* prev_chain, int64_t ts, |
|
uint64_t dh, uint8_t* out) { |
|
uint8_t buf[32 + 8 + 8]; |
|
memcpy(buf, prev_chain, 32); |
|
memcpy(buf + 32, &ts, 8); |
|
memcpy(buf + 40, &dh, 8); |
|
SHA256(buf, 48, out); |
|
} |
|
|
|
static int get_last_chain_hash(const char* ch_id, uint8_t* out) { |
|
char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl)); |
|
char sql[256]; snprintf(sql, sizeof(sql), |
|
"SELECT chain_hash FROM \"%s\" ORDER BY timestamp, datahash DESC LIMIT 1", tbl); |
|
sqlite3_stmt* stmt = NULL; |
|
if (sqlite3_prepare_v2(g_cc.db, sql, -1, &stmt, NULL) != SQLITE_OK) { |
|
memset(out, 0, 32); return -1; |
|
} |
|
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(out, blob, 32); |
|
else memset(out, 0, 32); |
|
} else { |
|
memset(out, 0, 32); |
|
} |
|
sqlite3_finalize(stmt); |
|
return 0; |
|
} |
|
|
|
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)); |
|
|
|
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; |
|
|
|
topo_node_sqlite_init(g_cc.db); |
|
topo_groups_set_sqlite_db(inst->topo_groups, g_cc.db); |
|
|
|
/* записать себя в nodes + local_identity с реальным node_id и именем */ |
|
{ |
|
const char* my_name = inst->name[0] ? inst->name : "Me"; |
|
sqlite3_stmt* st = NULL; |
|
sqlite3_prepare_v2(g_cc.db, |
|
"INSERT OR REPLACE INTO nodes(node_id, name, x25519_pubkey, ed25519_pubkey, online)" |
|
" 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_step(st); sqlite3_finalize(st); |
|
} |
|
st = NULL; |
|
sqlite3_prepare_v2(g_cc.db, |
|
"INSERT OR REPLACE INTO local_identity(id,node_id,name,x25519_pubkey,ed25519_pubkey)" |
|
" 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_step(st); sqlite3_finalize(st); |
|
} |
|
} |
|
|
|
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)" |
|
" 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_step(ins); sqlite3_finalize(ins); |
|
} |
|
} |
|
} |
|
sqlite3_finalize(cnt_st); |
|
} |
|
} |
|
|
|
g_cc.initialized = 1; |
|
|
|
/* 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) { |
|
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); |
|
} |
|
sqlite3_finalize(stmt); |
|
} |
|
} |
|
|
|
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; |
|
} |
|
|
|
void chat_core_destroy(struct UTUN_INSTANCE* inst) { |
|
(void)inst; |
|
if (!g_cc.initialized) return; |
|
g_cc.initialized = 0; |
|
|
|
if (g_cc.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; |
|
sqlite3_stmt* st = NULL; |
|
sqlite3_prepare_v2(g_cc.db, "UPDATE nodes SET name=? WHERE node_id=?", -1, &st, NULL); |
|
if (st) { |
|
sqlite3_bind_text(st, 1, name, -1, SQLITE_STATIC); |
|
sqlite3_bind_int64(st, 2, (sqlite3_int64)g_cc.my_node_id); |
|
sqlite3_step(st); sqlite3_finalize(st); |
|
} |
|
/* update instance name for NODEINFO */ |
|
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 FROM channels", -1, &cs, NULL); |
|
if (cs) { |
|
while (sqlite3_step(cs) == SQLITE_ROW) { |
|
const char* ch = (const char*)sqlite3_column_text(cs, 0); |
|
if (!ch) 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; |
|
mlen += snprintf((char*)join_msg + mlen, sizeof(join_msg) - mlen, "%s", ch) + 1; |
|
memcpy(join_msg + mlen, &myid, 8); mlen += 8; |
|
memcpy(join_msg + mlen, g_cc.inst->my_keys.public_key, 32); mlen += 32; |
|
{ const char* nm = name[0] ? name : ""; |
|
size_t nl = strlen(nm); memcpy(join_msg + mlen, nm, nl); mlen += nl; join_msg[mlen++] = '\0'; } |
|
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); |
|
} |
|
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, name, NULL); |
|
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, name, NULL, 0); |
|
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); |
|
} |
|
|
|
/* broadcast updated NODEINFO to connected peers */ |
|
struct TOPO_GROUP* grp = topo_groups_get_default(g_cc.inst->topo_groups); |
|
if (grp) topo_group_update_my_nodeinfo(g_cc.inst, grp); |
|
} |
|
|
|
void chat_core_update_my_name_trampoline(void* arg) { |
|
chat_core_update_my_name((const char*)arg); |
|
u_free(arg); |
|
} |
|
|
|
/* ─── отправка сообщения (GUI → uasync) ─── */ |
|
|
|
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); |
|
/* FIXME: JSON escaping for d (quotes/backslashes in data) */ |
|
/* For now, naive format; will break on special chars. TODO: use base64. */ |
|
|
|
int ret = db_sync_insert_signed(si, json, strlen(json), NULL, 0); |
|
if (ret != 0) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: db_sync_insert_signed failed ret=%d", CC_ID, ret); |
|
/* Insert directly to msg_ table as fallback (dual-write not triggered by callback) */ |
|
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)); |
|
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); |
|
sqlite3_stmt* st=NULL; |
|
if (sqlite3_prepare_v2(g_cc.db,sql,-1,&st,NULL)==SQLITE_OK) { |
|
sqlite3_bind_int64(st,1,(sqlite3_int64)g_cc.my_node_id); |
|
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(st,4,(sqlite3_int64)req->timestamp); |
|
sqlite3_bind_int64(st,5,(sqlite3_int64)dh); |
|
static const uint8_t ch32[32]={0}; sqlite3_bind_blob(st,6,ch32,32,SQLITE_STATIC); |
|
static const uint8_t sig64[64]={0}; sqlite3_bind_blob(st,7,sig64,64,SQLITE_STATIC); |
|
sqlite3_step(st); sqlite3_finalize(st); |
|
} |
|
/* 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); |
|
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, datahash 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; |
|
} |
|
|
|
/* ─── подключение к пиру из invite-ссылки ─── */ |
|
|
|
static void connect_result_cb(int result, uint64_t node_id, void* arg) { |
|
uint64_t channel_id = arg ? *(uint64_t*)arg : 0; |
|
uint8_t data[20]; |
|
memcpy(data, &node_id, 8); |
|
memcpy(data + 8, &result, 4); |
|
memcpy(data + 12, &channel_id, 8); |
|
gui_bridge_post(GUI_EVT_CONNECT_RESULT, data, 20); |
|
} |
|
|
|
#define CC_PARALLEL_CONNECT_TIMEOUT_MS 3000 |
|
|
|
struct cc_parallel_state { |
|
struct UTUN_INSTANCE* inst; |
|
int addr_count; |
|
int pending_count; |
|
int completed; /* 0=pending, 1=success, -1=fail-delivered */ |
|
uint64_t node_id; |
|
uint64_t channel_id; |
|
struct ETCP_CONN** conns; |
|
void** timers; |
|
}; |
|
|
|
struct cc_parallel_ctx { |
|
struct cc_parallel_state* state; |
|
int addr_index; |
|
}; |
|
|
|
static void cc_parallel_cleanup(struct cc_parallel_state* st) { |
|
if (!st) return; |
|
u_free(st->conns); |
|
u_free(st->timers); |
|
u_free(st); |
|
} |
|
|
|
static void cc_parallel_ready_cb(struct ETCP_CONN* conn, void* arg) { |
|
struct cc_parallel_ctx* pctx = (struct cc_parallel_ctx*)arg; |
|
struct cc_parallel_state* st = pctx->state; |
|
if (st->completed) { u_free(pctx); return; } |
|
st->completed = 1; |
|
|
|
uint8_t data[20]; |
|
memcpy(data, &st->node_id, 8); |
|
int r = CONN_MGR_OK; |
|
memcpy(data + 8, &r, 4); |
|
memcpy(data + 12, &st->channel_id, 8); |
|
gui_bridge_post(GUI_EVT_CONNECT_RESULT, data, 20); |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: parallel connect SUCCESS idx=%d peer=0x%016llx", |
|
CC_ID, pctx->addr_index, (unsigned long long)st->node_id); |
|
|
|
for (int i = 0; i < st->addr_count; i++) { |
|
if (st->timers[i]) { uasync_cancel_timeout(st->inst->ua, st->timers[i]); st->timers[i] = NULL; } |
|
if (st->conns[i] && i != pctx->addr_index) { uasync_call_soon(st->inst->ua, st->conns[i], (timeout_callback_t)etcp_connection_close); st->conns[i] = NULL; } |
|
} |
|
u_free(pctx); |
|
cc_parallel_cleanup(st); |
|
} |
|
|
|
static void cc_parallel_timeout_cb(void* arg) { |
|
struct cc_parallel_ctx* pctx = (struct cc_parallel_ctx*)arg; |
|
struct cc_parallel_state* st = pctx->state; |
|
if (st->completed) { u_free(pctx); return; } |
|
|
|
if (st->conns[pctx->addr_index]) { etcp_connection_close(st->conns[pctx->addr_index]); st->conns[pctx->addr_index] = NULL; } |
|
st->timers[pctx->addr_index] = NULL; |
|
st->pending_count--; |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: parallel connect TIMEOUT idx=%d pending=%d/%d peer=0x%016llx", |
|
CC_ID, pctx->addr_index, st->pending_count, st->addr_count, (unsigned long long)st->node_id); |
|
|
|
if (st->pending_count <= 0 && !st->completed) { |
|
st->completed = -1; |
|
uint8_t data[20]; |
|
memcpy(data, &st->node_id, 8); |
|
int r = CONN_MGR_ERR_TIMEOUT; |
|
memcpy(data + 8, &r, 4); |
|
memcpy(data + 12, &st->channel_id, 8); |
|
gui_bridge_post(GUI_EVT_CONNECT_RESULT, data, 20); |
|
cc_parallel_cleanup(st); |
|
} |
|
u_free(pctx); |
|
} |
|
|
|
void chat_core_connect_from_invite(struct chat_invite* inv) { |
|
if (!g_cc.initialized || !g_cc.inst || !inv) return; |
|
struct TOPO_GROUP* group = topo_groups_get_default(g_cc.inst->topo_groups); |
|
if (!group || !g_cc.inst->conn_mgr) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: bgp/conn_mgr not available", CC_ID); |
|
uint8_t err[12]; int r = -7; |
|
memcpy(err, &inv->node_id, 8); memcpy(err + 8, &r, 4); |
|
gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 12); |
|
return; |
|
} |
|
|
|
uint64_t node_id = inv->node_id; |
|
|
|
if (topo_node_find_by_id(group, node_id)) { |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: node 0x%016llx already in BGP, connecting", |
|
CC_ID, (unsigned long long)node_id); |
|
conn_mgr_connect_node(g_cc.inst->conn_mgr, node_id, 30000, |
|
connect_result_cb, &inv->channel_id); |
|
return; |
|
} |
|
|
|
struct ll_entry* qe = queue_entry_new(sizeof(struct TOPO_NODEQ)); |
|
if (!qe) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: failed to alloc TOPO_NODEQ", CC_ID); |
|
uint8_t err[12]; int r = -7; |
|
memcpy(err, &node_id, 8); memcpy(err + 8, &r, 4); |
|
gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 12); |
|
return; |
|
} |
|
struct TOPO_NODEQ* nq = (struct TOPO_NODEQ*)qe; |
|
memset((uint8_t*)nq + sizeof(struct ll_entry), 0, sizeof(*nq) - sizeof(struct ll_entry)); |
|
|
|
struct TOPO_NODE* ni = u_calloc(1, sizeof(struct TOPO_NODE)); |
|
if (!ni) { queue_entry_free(qe); return; } |
|
ni->node_id = node_id; |
|
ni->group_id = group->group_id; |
|
ni->ver = 0; |
|
memcpy(ni->public_key, inv->pubkey, SC_PUBKEY_SIZE); |
|
memset(ni->ed25519_public_key, 0, SC_PUBKEY_SIZE); |
|
ni->node_name = u_strdup(""); |
|
|
|
struct TOPO_ADDR4* addrs_head = NULL; |
|
const uint8_t* src = inv->addrs_data; |
|
for (int i = 0; i < inv->addr_count; i++) { |
|
uint8_t family = *src++; |
|
if (family == 4) { |
|
struct TOPO_ADDR4* a4 = memory_pool_alloc(group->instance->topo_groups->v4_addr_pool); |
|
if (!a4) continue; |
|
memcpy(a4->addr, src, 4); src += 4; |
|
a4->port = ((uint16_t)src[0] << 8) | src[1]; src += 2; |
|
a4->type = TOPO_ADDR_REAL; a4->socket_id = 0; a4->protocol = TOPO_PROTO_UDP; |
|
a4->next = addrs_head; addrs_head = a4; |
|
} else { |
|
src += 18; /* skip v6 */ |
|
} |
|
} |
|
ni->v4_addrs = addrs_head; |
|
|
|
/* alien-узел: sock_meta с неизвестной конфигурацией цели */ |
|
struct TOPO_SOCKMETA4* sm = memory_pool_alloc(group->instance->topo_groups->v4_sock_meta_pool); |
|
if (sm) { sm->id = 0; sm->config_type = CFG_SERVER_TYPE_UNKNOWN; sm->nat_type = NAT_TYPE_UNKNOWN; sm->next = NULL; ni->v4_sock_meta = sm; } |
|
|
|
nq->node = ni; topo_node_ref(ni); |
|
nq->hash_node_id = node_id; nq->alien = 0; |
|
nq->dirty = 0; nq->last_ver = 0; |
|
nq->conn_mgr_type = CONN_TYPE_NONE; nq->conn_mgr_intermediariy_count = 0; |
|
memset(&nq->connectivity, 0, sizeof(nq->connectivity)); |
|
nq->best_socket = NULL; |
|
|
|
queue_data_put_with_index(group->nodes, &nq->ll); |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: created NODEINFO for 0x%016llx, %d addrs, parallel connect", |
|
CC_ID, (unsigned long long)node_id, inv->addr_count); |
|
|
|
/* count IPv4 addresses and collect them for parallel connect */ |
|
int v4_count = 0; |
|
for (struct TOPO_ADDR4* a = addrs_head; a; a = a->next) v4_count++; |
|
if (v4_count == 0) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: no IPv4 addresses in invite", CC_ID); |
|
uint8_t err[20]; memcpy(err, &node_id, 8); int r = CONN_MGR_ERR_NO_ADDRESSES; |
|
memcpy(err + 8, &r, 4); memset(err + 12, 0, 8); |
|
gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20); |
|
return; |
|
} |
|
|
|
struct ETCP_SOCKET* best_socket = g_cc.inst->etcp_sockets; |
|
while (best_socket && best_socket->local_addr.ss_family != AF_INET) best_socket = best_socket->next; |
|
if (!best_socket) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: no AF_INET socket", CC_ID); |
|
uint8_t err[20]; memcpy(err, &node_id, 8); int r = CONN_MGR_ERR_INTERNAL; |
|
memcpy(err + 8, &r, 4); memset(err + 12, 0, 8); |
|
gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20); |
|
return; |
|
} |
|
|
|
struct cc_parallel_state* pst = u_calloc(1, sizeof(struct cc_parallel_state)); |
|
if (!pst) { |
|
uint8_t err[20]; memcpy(err, &node_id, 8); int r = CONN_MGR_ERR_INTERNAL; |
|
memcpy(err + 8, &r, 4); memset(err + 12, 0, 8); |
|
gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20); |
|
return; |
|
} |
|
pst->inst = g_cc.inst; |
|
pst->addr_count = v4_count; |
|
pst->pending_count = v4_count; |
|
pst->completed = 0; |
|
pst->node_id = node_id; |
|
pst->channel_id = inv->channel_id; |
|
pst->conns = u_calloc(v4_count, sizeof(struct ETCP_CONN*)); |
|
pst->timers = u_calloc(v4_count, sizeof(void*)); |
|
if (!pst->conns || !pst->timers) { cc_parallel_cleanup(pst); return; } |
|
|
|
int idx = 0; |
|
for (struct TOPO_ADDR4* a = addrs_head; a && idx < v4_count; a = a->next, idx++) { |
|
struct sockaddr_in sin; |
|
memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET; |
|
memcpy(&sin.sin_addr.s_addr, a->addr, 4); sin.sin_port = htons(a->port); |
|
struct sockaddr_storage sa; |
|
memset(&sa, 0, sizeof(sa)); memcpy(&sa, &sin, sizeof(sin)); |
|
|
|
struct ETCP_CONN* conn = etcp_connection_create(g_cc.inst, NULL); |
|
if (!conn) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: etcp_connection_create failed idx=%d", CC_ID, idx); |
|
pst->conns[idx] = NULL; pst->timers[idx] = NULL; pst->pending_count--; continue; |
|
} |
|
sc_init_ctx(&conn->crypto_ctx, &g_cc.inst->my_keys); |
|
sc_set_peer_public_key(&conn->crypto_ctx, inv->pubkey, 0); |
|
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: invite connect set peer_pubkey=%016llx my_pub=%016llx peer_node=0x%016llx", |
|
CC_ID, *(uint64_t*)inv->pubkey, *(uint64_t*)conn->crypto_ctx.pk->public_key, inv->node_id); |
|
|
|
struct cc_parallel_ctx* pctx = u_calloc(1, sizeof(struct cc_parallel_ctx)); |
|
if (!pctx) { etcp_connection_close(conn); pst->conns[idx] = NULL; pst->timers[idx] = NULL; pst->pending_count--; continue; } |
|
pctx->state = pst; pctx->addr_index = idx; |
|
|
|
etcp_conn_set_ready_cbk(conn, cc_parallel_ready_cb, pctx); |
|
|
|
if (!etcp_link_new(conn, best_socket, &sa, 0)) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: etcp_link_new failed idx=%d", CC_ID, idx); |
|
u_free(pctx); etcp_connection_close(conn); |
|
pst->conns[idx] = NULL; pst->timers[idx] = NULL; pst->pending_count--; continue; |
|
} |
|
pst->conns[idx] = conn; |
|
etcp_conn_set_peer_node_id(conn, inv->node_id); |
|
pst->timers[idx] = uasync_set_timeout(g_cc.inst->ua, CC_PARALLEL_CONNECT_TIMEOUT_MS * 10, |
|
pctx, cc_parallel_timeout_cb, "cc_parallel"); |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: parallel connect attempt %d/%d to %d.%d.%d.%d:%d", |
|
CC_ID, idx + 1, v4_count, sin.sin_addr.s_addr & 0xFF, (sin.sin_addr.s_addr >> 8) & 0xFF, |
|
(sin.sin_addr.s_addr >> 16) & 0xFF, (sin.sin_addr.s_addr >> 24) & 0xFF, a->port); |
|
} |
|
|
|
if (pst->pending_count <= 0 && !pst->completed) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: all %d parallel connects failed to start", CC_ID, v4_count); |
|
uint8_t err[20]; memcpy(err, &node_id, 8); int r = CONN_MGR_ERR_UNREACHABLE; |
|
memcpy(err + 8, &r, 4); memcpy(err + 12, &inv->channel_id, 8); |
|
gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20); |
|
cc_parallel_cleanup(pst); |
|
} |
|
} |
|
|
|
/* ═══ Auto-connect: direct ETCP connection from SQLite (no BGP) ═══ */ |
|
|
|
#define CA_CONNECT_TIMEOUT_MS 3000 |
|
|
|
struct ca_state; |
|
|
|
struct ca_ctx { |
|
struct ca_state* state; |
|
int addr_index; |
|
}; |
|
|
|
struct ca_state { |
|
struct UTUN_INSTANCE* inst; |
|
int addr_count; |
|
int pending_count; |
|
int delivered; /* 0=pending, 1=result already delivered */ |
|
int cancelled; /* 1=externally cancelled, do not deliver result */ |
|
uint64_t node_id; |
|
struct ETCP_CONN** conns; |
|
void** timers; |
|
struct ca_ctx** ctxs; |
|
void (*result_cb)(int result, uint64_t node_id, void* arg); |
|
void* result_arg; |
|
uint8_t pubkey[SC_PUBKEY_SIZE]; |
|
}; |
|
|
|
static void ca_cleanup(struct ca_state* st) { |
|
if (!st) return; |
|
if (st->ctxs) { for (int i = 0; i < st->addr_count; i++) u_free(st->ctxs[i]); u_free(st->ctxs); } |
|
u_free(st->conns); |
|
u_free(st->timers); |
|
u_free(st); |
|
} |
|
|
|
static void ca_ready_cb(struct ETCP_CONN* conn, void* arg) { |
|
struct ca_ctx* ctx = (struct ca_ctx*)arg; |
|
struct ca_state* st = ctx->state; |
|
if (st->delivered || st->cancelled) return; |
|
st->delivered = 1; |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect SUCCESS idx=%d peer=0x%016llx", |
|
CC_ID, ctx->addr_index, (unsigned long long)st->node_id); |
|
for (int i = 0; i < st->addr_count; i++) { |
|
if (st->timers[i]) { uasync_cancel_timeout(st->inst->ua, st->timers[i]); st->timers[i] = NULL; } |
|
if (st->conns[i] && i != ctx->addr_index) { uasync_call_soon(st->inst->ua, st->conns[i], (timeout_callback_t)etcp_connection_close); st->conns[i] = NULL; } |
|
} |
|
etcp_conn_remove_ready_cbk(conn, ca_ready_cb, ctx); |
|
st->result_cb(CONN_MGR_OK, st->node_id, st->result_arg); |
|
ca_cleanup(st); |
|
} |
|
|
|
static void ca_timeout_cb(void* arg) { |
|
struct ca_ctx* ctx = (struct ca_ctx*)arg; |
|
struct ca_state* st = ctx->state; |
|
if (st->delivered || st->cancelled) return; |
|
st->timers[ctx->addr_index] = NULL; |
|
st->pending_count--; |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect TIMEOUT idx=%d pending=%d/%d peer=0x%016llx", |
|
CC_ID, ctx->addr_index, st->pending_count, st->addr_count, (unsigned long long)st->node_id); |
|
if (st->pending_count <= 0 && !st->delivered) { |
|
st->delivered = 1; |
|
st->result_cb(CONN_MGR_ERR_TIMEOUT, st->node_id, st->result_arg); |
|
} |
|
} |
|
|
|
void chat_core_connect_auto_cancel(void* state) { |
|
if (!state) return; |
|
struct ca_state* st = (struct ca_state*)state; |
|
st->cancelled = 1; |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect CANCELLED peer=0x%016llx", |
|
CC_ID, (unsigned long long)st->node_id); |
|
for (int i = 0; i < st->addr_count; i++) { |
|
if (st->timers[i]) { uasync_cancel_timeout(st->inst->ua, st->timers[i]); st->timers[i] = NULL; } |
|
if (st->conns[i]) { etcp_connection_close(st->conns[i]); st->conns[i] = NULL; } |
|
} |
|
ca_cleanup(st); |
|
} |
|
|
|
void chat_core_connect_auto(uint64_t node_id, |
|
void (*cb)(int result, uint64_t node_id, void* arg), |
|
void* arg, |
|
void** out_state) { |
|
if (!g_cc.initialized || !g_cc.inst || !cb) return; |
|
if (out_state) *out_state = NULL; |
|
|
|
struct ETCP_SOCKET* best_socket = g_cc.inst->etcp_sockets; |
|
while (best_socket && best_socket->local_addr.ss_family != AF_INET) best_socket = best_socket->next; |
|
if (!best_socket) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect no AF_INET socket for node 0x%016llx", |
|
CC_ID, (unsigned long long)node_id); |
|
cb(CONN_MGR_ERR_INTERNAL, node_id, arg); |
|
return; |
|
} |
|
|
|
uint8_t pubkey[SC_PUBKEY_SIZE] = {0}; |
|
sqlite3_stmt* st = NULL; |
|
if (sqlite3_prepare_v2(g_cc.db, |
|
"SELECT x25519_pubkey FROM nodes WHERE node_id=?", |
|
-1, &st, NULL) == SQLITE_OK) { |
|
sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); |
|
if (sqlite3_step(st) == SQLITE_ROW) { |
|
const void* pk = sqlite3_column_blob(st, 0); |
|
if (pk && sqlite3_column_bytes(st, 0) == SC_PUBKEY_SIZE) |
|
memcpy(pubkey, pk, SC_PUBKEY_SIZE); |
|
} |
|
sqlite3_finalize(st); |
|
} |
|
if (pubkey[0] == 0 && pubkey[1] == 0) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect no pubkey for node 0x%016llx", |
|
CC_ID, (unsigned long long)node_id); |
|
cb(CONN_MGR_ERR_NOT_FOUND, node_id, arg); |
|
return; |
|
} |
|
|
|
/* collect IPv4 addresses */ |
|
struct { uint8_t addr[4]; uint16_t port; } addrs[16]; |
|
int addr_count = 0; |
|
if (sqlite3_prepare_v2(g_cc.db, |
|
"SELECT address, port FROM node_addresses WHERE node_id=? AND family=4", |
|
-1, &st, NULL) == SQLITE_OK) { |
|
sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); |
|
while (sqlite3_step(st) == SQLITE_ROW && addr_count < 16) { |
|
const void* a = sqlite3_column_blob(st, 0); |
|
int alen = sqlite3_column_bytes(st, 0); |
|
if (a && alen == 4) { |
|
memcpy(addrs[addr_count].addr, a, 4); |
|
addrs[addr_count].port = (uint16_t)sqlite3_column_int(st, 1); |
|
addr_count++; |
|
} |
|
} |
|
sqlite3_finalize(st); |
|
} |
|
/* dedup by (addr, port) — can't have two links to same socket */ |
|
{ int uniq = 0; |
|
for (int i = 0; i < addr_count; i++) { |
|
int dup = 0; |
|
for (int j = 0; j < uniq; j++) |
|
if (memcmp(addrs[i].addr, addrs[j].addr, 4) == 0 && addrs[i].port == addrs[j].port) { dup = 1; break; } |
|
if (!dup) { if (i != uniq) addrs[uniq] = addrs[i]; uniq++; } |
|
} |
|
addr_count = uniq; } |
|
if (addr_count == 0) { |
|
DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "%s: auto_connect no IPv4 addrs yet for node 0x%016llx (BGP not synced?)", CC_ID, (unsigned long long)node_id); |
|
cb(CONN_MGR_ERR_NO_ADDRESSES, node_id, arg); |
|
return; |
|
} |
|
|
|
struct ca_state* pst = u_calloc(1, sizeof(struct ca_state)); |
|
if (!pst) { cb(CONN_MGR_ERR_INTERNAL, node_id, arg); return; } |
|
pst->inst = g_cc.inst; |
|
pst->addr_count = addr_count; |
|
pst->pending_count = addr_count; |
|
pst->node_id = node_id; |
|
pst->result_cb = cb; |
|
pst->result_arg = arg; |
|
memcpy(pst->pubkey, pubkey, SC_PUBKEY_SIZE); |
|
pst->conns = u_calloc(addr_count, sizeof(struct ETCP_CONN*)); |
|
pst->timers = u_calloc(addr_count, sizeof(void*)); |
|
pst->ctxs = u_calloc(addr_count, sizeof(struct ca_ctx*)); |
|
if (!pst->conns || !pst->timers || !pst->ctxs) { ca_cleanup(pst); cb(CONN_MGR_ERR_INTERNAL, node_id, arg); return; } |
|
|
|
for (int i = 0; i < addr_count; i++) { |
|
struct sockaddr_in sin; |
|
memset(&sin, 0, sizeof(sin)); sin.sin_family = AF_INET; |
|
memcpy(&sin.sin_addr.s_addr, addrs[i].addr, 4); |
|
sin.sin_port = htons(addrs[i].port); |
|
struct sockaddr_storage sa; |
|
memset(&sa, 0, sizeof(sa)); memcpy(&sa, &sin, sizeof(sin)); |
|
|
|
struct ETCP_CONN* conn = etcp_connection_create(g_cc.inst, NULL); |
|
if (!conn) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect etcp_connection_create failed idx=%d", CC_ID, i); |
|
pst->conns[i] = NULL; pst->timers[i] = NULL; pst->ctxs[i] = NULL; pst->pending_count--; continue; |
|
} |
|
sc_init_ctx(&conn->crypto_ctx, &g_cc.inst->my_keys); |
|
sc_set_peer_public_key(&conn->crypto_ctx, pubkey, 0); |
|
etcp_conn_set_peer_node_id(conn, node_id); |
|
|
|
struct ca_ctx* pctx = u_calloc(1, sizeof(struct ca_ctx)); |
|
if (!pctx) { etcp_connection_close(conn); pst->conns[i] = NULL; pst->timers[i] = NULL; pst->ctxs[i] = NULL; pst->pending_count--; continue; } |
|
pctx->state = pst; pctx->addr_index = i; |
|
pst->ctxs[i] = pctx; |
|
|
|
etcp_conn_set_ready_cbk(conn, ca_ready_cb, pctx); |
|
|
|
if (!etcp_link_new(conn, best_socket, &sa, 0)) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect etcp_link_new failed idx=%d", CC_ID, i); |
|
u_free(pctx); etcp_connection_close(conn); |
|
pst->conns[i] = NULL; pst->timers[i] = NULL; pst->ctxs[i] = NULL; pst->pending_count--; continue; |
|
} |
|
pst->conns[i] = conn; |
|
pst->timers[i] = uasync_set_timeout(g_cc.inst->ua, CA_CONNECT_TIMEOUT_MS * 10, |
|
pctx, ca_timeout_cb, "ca_timeout"); |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect attempt %d/%d to %d.%d.%d.%d:%d", |
|
CC_ID, i + 1, addr_count, |
|
addrs[i].addr[0], addrs[i].addr[1], addrs[i].addr[2], addrs[i].addr[3], |
|
addrs[i].port); |
|
} |
|
|
|
if (out_state) *out_state = pst; |
|
|
|
if (pst->pending_count <= 0 && !pst->delivered) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: auto_connect all %d attempts failed to start", CC_ID, addr_count); |
|
if (out_state) *out_state = NULL; |
|
cb(CONN_MGR_ERR_UNREACHABLE, node_id, arg); |
|
ca_cleanup(pst); |
|
} |
|
} |
|
|
|
/* ─── подготовка инфраструктуры канала (msg-таблица + 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)); |
|
char sql[512]; |
|
|
|
snprintf(sql, sizeof(sql), |
|
"CREATE TABLE IF NOT EXISTS \"%s\" (" |
|
" id INTEGER PRIMARY KEY AUTOINCREMENT," |
|
" node_id INTEGER NOT NULL," |
|
" content_type TEXT NOT NULL," |
|
" data BLOB NOT NULL," |
|
" timestamp INTEGER NOT NULL," |
|
" datahash INTEGER NOT NULL," |
|
" chain_hash BLOB NOT NULL," |
|
" signature BLOB," |
|
" is_outgoing INTEGER DEFAULT 0," |
|
" is_read INTEGER DEFAULT 0," |
|
" sync_flags INTEGER DEFAULT 0," |
|
" UNIQUE(timestamp, datahash))", tbl_msg); |
|
db_exec(sql); |
|
snprintf(sql, sizeof(sql), |
|
"CREATE INDEX IF NOT EXISTS \"idx_%s_ts_dh\" ON \"%s\"(timestamp, datahash)", |
|
tbl_msg, tbl_msg); |
|
db_exec(sql); |
|
|
|
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, "chats", ch_hash); |
|
if (si) { |
|
si_register(si, ch_id); |
|
db_sync_set_insert_cb(si, on_db_sync_insert, u_strdup(ch_id)); |
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: channel ready ch=%s hash=0x%016llx", |
|
CC_ID, ch_id, (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); |
|
|
|
chat_core_ensure_channel_ready(req->channel_id); |
|
|
|
/* записываем канал в БД */ |
|
int rc = topo_node_sqlite_channel_put(g_cc.db, |
|
req->channel_id, req->name, req->is_dm, 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); |
|
|
|
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; |
|
mlen += snprintf((char*)join_msg + mlen, sizeof(join_msg) - mlen, "%s", req->channel_id) + 1; |
|
memcpy(join_msg + mlen, &myid, 8); mlen += 8; |
|
memcpy(join_msg + mlen, g_cc.inst->my_keys.public_key, 32); mlen += 32; |
|
{ const char* nm = g_cc.inst->name[0] ? g_cc.inst->name : ""; |
|
size_t nl = strlen(nm); memcpy(join_msg + mlen, nm, nl); mlen += nl; join_msg[mlen++] = '\0'; } |
|
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); |
|
} |
|
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, g_cc.inst->name, NULL, 0); |
|
} |
|
|
|
/* 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 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; |
|
/* parse JSON: {"n":<uint64>,"ch":"<str>","ct":"<str>","d":"<str>"} */ |
|
const char* p = data; const char* end = data + len; uint64_t jn = 0; |
|
{ /* "n" */ |
|
const char* n = strstr(p, "\"n\":"); if (!n) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: on_db_sync_insert JSON parse fail: no '\"n\"' field ch=%s", CC_ID, ch_id); return; } |
|
n += 4; jn = strtoull(n, NULL, 0); p = n; |
|
} |
|
const char* jct = "text/plain"; size_t jct_len = 10; |
|
{ /* "ct" */ |
|
const char* ct = strstr(p, "\"ct\":\""); if (ct) { ct += 6; const char* ce = strchr(ct, '"'); if (ce) { jct = ct; jct_len = (size_t)(ce - ct); p = ce; } } |
|
} |
|
const char* jd_start = NULL; size_t jd_len = 0; |
|
{ /* "d" */ |
|
const char* d = strstr(p, "\"d\":\""); if (d) { d += 5; const char* de = d; while (de < end) { if (*de == '"' && (de == d || *(de - 1) != '\\')) break; de++; } if (de < end) { jd_start = d; jd_len = (size_t)(de - d); } } |
|
} |
|
if (!jd_start) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: on_db_sync_insert JSON parse fail: no '\"d\"' field ch=%s", CC_ID, ch_id); return; } |
|
|
|
uint64_t dh; compute_datahash((const uint8_t*)jd_start, jd_len, &dh); |
|
uint64_t ts = db_sync_get_last_timestamp(si); |
|
char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl)); |
|
uint8_t prev_ch[32]; get_last_chain_hash(ch_id, prev_ch); |
|
uint8_t chain_h[32]; compute_chain_hash(prev_ch, (int64_t)ts, dh, chain_h); |
|
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) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: on_db_sync_insert prepare failed ch=%s", CC_ID, ch_id); return; } |
|
sqlite3_bind_int64(st,1,(sqlite3_int64)jn); |
|
sqlite3_bind_text(st,2,jct,(int)jct_len,SQLITE_STATIC); |
|
sqlite3_bind_blob(st,3,jd_start,(int)jd_len,SQLITE_STATIC); |
|
sqlite3_bind_int64(st,4,(sqlite3_int64)ts); |
|
sqlite3_bind_int64(st,5,(sqlite3_int64)dh); |
|
sqlite3_bind_blob(st,6,chain_h,32,SQLITE_STATIC); |
|
{ static const uint8_t z64[64]={0}; sqlite3_bind_blob(st,7,z64,64,SQLITE_STATIC); } |
|
sqlite3_bind_int(st,8,(uint64_t)jn==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); |
|
} else if (rc == SQLITE_CONSTRAINT) { |
|
DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: on_db_sync_insert duplicate ch=%s ts=%lld dh=0x%016llx", |
|
CC_ID, ch_id, (long long)ts, (unsigned long long)dh); |
|
} else { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: on_db_sync_insert step failed rc=%d ch=%s", |
|
CC_ID, rc, ch_id); |
|
} |
|
}
|
|
|