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.
 
 
 
 
 
 

1341 lines
62 KiB

#include "chat_sync.h"
#include "chat_core.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/etcp.h"
#include "../../../src/secure_channel.h"
#include "../../../src/ntp_time.h"
#include "../../../lib/u_async.h"
#include "../../../lib/ll_queue.h"
#include "../../../lib/debug_config.h"
#include "../../../lib/mem.h"
#include "../../../lib/platform_compat.h"
#include <string.h>
#include <sqlite3.h>
#include <openssl/evp.h>
static struct chat_sync* g_cs = NULL;
/* ═══════════════════════════════════════════════════════════════════════
* Auto-connect: cursor-based peer iteration with GC every 1s
* ══════════════════════════════════════════════════════════════════════ */
#define AC_MAX_FLIGHTS 10
#define AC_RETRY_MS 1000
#define AC_GC_TIMEOUT_MS 3000
#define AC_ID "auto_connect"
struct ac_flight {
uint64_t node_id;
void* ca_state; /* opaque, owned by chat_core_connect_auto */
uint64_t created_tb; /* get_time_tb() when launched */
};
struct auto_connect {
struct UTUN_INSTANCE* inst;
void* retry_timer;
uint8_t active;
char** channel_ids; /* loaded once, refreshed on cursor wrap */
int channel_count;
int ch_cursor; /* current channel index */
int peer_cursor; /* current peer index within channel */
struct ac_flight flights[AC_MAX_FLIGHTS];
};
static struct auto_connect* g_ac = NULL;
static void ac_retry_timer_cb(void* arg);
static void ac_result_cb(int result, uint64_t node_id, void* arg);
/* ── load channel IDs into auto_connect ── */
static int ac_load_channels(struct auto_connect* ac) {
if (ac->channel_ids) {
for (int i = 0; i < ac->channel_count; i++) u_free(ac->channel_ids[i]);
u_free(ac->channel_ids);
ac->channel_ids = NULL;
}
ac->channel_count = 0;
uint8_t buf[4096]; size_t buf_len;
if (chat_core_list_channels(buf, sizeof(buf), &buf_len) != 0 || buf_len < 2) return 0;
uint16_t cnt; memcpy(&cnt, buf, 2);
const uint8_t* p = buf + 2; size_t rem = buf_len - 2;
ac->channel_ids = u_calloc(cnt, sizeof(char*));
if (!ac->channel_ids) return 0;
for (uint16_t i = 0; i < cnt && rem >= 1; i++) {
uint8_t id_len = *p++; rem--;
if (rem < id_len) break;
ac->channel_ids[i] = u_malloc(id_len + 1);
if (ac->channel_ids[i]) { memcpy(ac->channel_ids[i], p, id_len); ac->channel_ids[i][id_len] = '\0'; ac->channel_count++; }
p += id_len; rem -= id_len;
}
return ac->channel_count;
}
/* ── GC: close expired flights (no link_status after AC_GC_TIMEOUT_MS) ── */
static void ac_gc(struct auto_connect* ac) {
uint64_t now = get_time_tb();
uint64_t deadline = (uint64_t)AC_GC_TIMEOUT_MS * 10;
for (int i = 0; i < AC_MAX_FLIGHTS; i++) {
if (!ac->flights[i].ca_state) continue;
if (now - ac->flights[i].created_tb < deadline) continue;
uint64_t nid = ac->flights[i].node_id;
/* check if link already UP — if so, just free slot (already connected) */
int found_up = 0;
{
struct ll_entry* entry = ac->inst->connections->head;
while (entry) {
struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data;
if (ce->conn->peer_node_id == nid) {
struct ETCP_LINK* l = ce->conn->links;
while (l) { if (l->initialized && l->link_status) { found_up = 1; break; } l = l->next; }
if (found_up) break;
}
entry = entry->next;
}
}
if (found_up) {
chat_core_connect_auto_cancel(ac->flights[i].ca_state);
memset(&ac->flights[i], 0, sizeof(ac->flights[i]));
continue;
}
/* expired and not up — cancel */
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: GC closing flight node=0x%016llx (expired)", AC_ID, (unsigned long long)nid);
chat_core_connect_auto_cancel(ac->flights[i].ca_state);
memset(&ac->flights[i], 0, sizeof(ac->flights[i]));
}
}
/* ── count occupied flight slots ── */
static int ac_flight_count(struct auto_connect* ac) {
int n = 0;
for (int i = 0; i < AC_MAX_FLIGHTS; i++) if (ac->flights[i].ca_state) n++;
return n;
}
/* ── find first free flight slot ── */
static int ac_find_free_slot(struct auto_connect* ac) {
for (int i = 0; i < AC_MAX_FLIGHTS; i++) if (!ac->flights[i].ca_state) return i;
return -1;
}
/* ── check if node already has active ETCP connection ── */
static int ac_node_has_conn(struct UTUN_INSTANCE* inst, uint64_t nid) {
if (!inst->connections) return 0;
struct ll_entry* e = queue_find_data_by_index(inst->connections, (const uint8_t*)&nid);
return e != NULL;
}
/* ── fill up to AC_MAX_FLIGHTS by advancing cursor ── */
static void ac_fill(struct auto_connect* ac) {
if (ac->channel_count == 0) return;
/* try all channels up to 2 full cycles without finding a candidate → give up */
int max_iter = ac->channel_count * 2 + 10;
while (ac_flight_count(ac) < AC_MAX_FLIGHTS && max_iter-- > 0) {
int slot = ac_find_free_slot(ac);
if (slot < 0) break;
char* ch_id = ac->channel_ids[ac->ch_cursor];
uint8_t pb[2048]; size_t plen;
if (chat_core_list_peers(ch_id, pb, sizeof(pb), &plen) != 0 || plen < 2) {
/* channel has no peers or error — skip to next */
ac->ch_cursor++; ac->peer_cursor = 0;
if (ac->ch_cursor >= ac->channel_count) { ac->ch_cursor = 0; ac_load_channels(ac); }
continue;
}
uint16_t pc; memcpy(&pc, pb, 2);
if (ac->peer_cursor >= (int)pc) {
ac->ch_cursor++; ac->peer_cursor = 0;
if (ac->ch_cursor >= ac->channel_count) { ac->ch_cursor = 0; ac_load_channels(ac); }
continue;
}
/* get next peer */
uint64_t nid; memcpy(&nid, pb + 2 + ac->peer_cursor * 8, 8);
ac->peer_cursor++;
if (nid == 0 || nid == ac->inst->node_id) continue;
/* already connected? */
if (ac_node_has_conn(ac->inst, nid)) continue;
/* already in a flight slot? */
int dup = 0;
for (int i = 0; i < AC_MAX_FLIGHTS; i++)
if (ac->flights[i].ca_state && ac->flights[i].node_id == nid) { dup = 1; break; }
if (dup) continue;
/* launch */
struct ac_flight* f = &ac->flights[slot];
f->node_id = nid;
f->created_tb = get_time_tb();
chat_core_connect_auto(nid, ac_result_cb, f, &f->ca_state);
if (!f->ca_state) {
f->node_id = 0; f->created_tb = 0;
continue;
}
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: flight[%d] launched node=0x%016llx", AC_ID, slot, (unsigned long long)nid);
uint8_t evt[7]; evt[0] = 0; uint16_t t = 0, n = (uint16_t)ac->channel_count, s = 0;
memcpy(evt + 1, &t, 2); memcpy(evt + 3, &n, 2); memcpy(evt + 5, &s, 2);
gui_bridge_post(GUI_EVT_AUTO_CONNECT_STATUS, evt, 7);
}
}
/* ── result callback (fired by chat_core_connect_auto) ── */
static void ac_result_cb(int result, uint64_t node_id, void* arg) {
struct ac_flight* f = (struct ac_flight*)arg;
if (!g_ac || !g_ac->active) return;
const char* rs = result == CC_OK ? "OK" : "FAIL";
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: result node=0x%016llx %s", AC_ID, (unsigned long long)node_id, rs);
if (result == CC_OK) f->ca_state = NULL; /* already freed by ca_ready_cb, prevent double-free in GC/stop */
}
/* ── retry timer callback (every AC_RETRY_MS) ── */
static void ac_retry_timer_cb(void* arg) {
struct auto_connect* ac = (struct auto_connect*)arg;
ac->retry_timer = NULL;
if (!ac->active) return;
if (ac->channel_count == 0) ac_load_channels(ac);
ac_gc(ac);
ac_fill(ac);
ac->retry_timer = uasync_set_timeout(ac->inst->ua, AC_RETRY_MS * 10, ac, ac_retry_timer_cb, "ac_retry");
}
/* ═══════════════════════════════════════════════════════════════════════
* Public auto_connect API
* ══════════════════════════════════════════════════════════════════════ */
void chat_sync_auto_connect_start(struct UTUN_INSTANCE* inst) {
if (!inst || !inst->ua) return;
chat_sync_auto_connect_stop();
struct auto_connect* ac = u_calloc(1, sizeof(*ac));
if (!ac) return;
ac->inst = inst;
ac->active = 1;
g_ac = ac;
ac_load_channels(ac);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: start, %d channels loaded", AC_ID, ac->channel_count);
ac_gc(ac);
ac_fill(ac);
ac->retry_timer = uasync_set_timeout(inst->ua, AC_RETRY_MS * 10, ac, ac_retry_timer_cb, "ac_retry");
}
void chat_sync_auto_connect_stop(void) {
struct auto_connect* ac = g_ac;
if (!ac) return;
ac->active = 0; g_ac = NULL;
if (ac->retry_timer) { uasync_cancel_timeout(ac->inst->ua, ac->retry_timer); ac->retry_timer = NULL; }
for (int i = 0; i < AC_MAX_FLIGHTS; i++)
if (ac->flights[i].ca_state) { chat_core_connect_auto_cancel(ac->flights[i].ca_state); memset(&ac->flights[i], 0, sizeof(ac->flights[i])); }
if (ac->channel_ids) { for (int i = 0; i < ac->channel_count; i++) u_free(ac->channel_ids[i]); u_free(ac->channel_ids); }
uint8_t evt[7]; evt[0] = 2; uint16_t z = 0;
memcpy(evt + 1, &z, 2); memcpy(evt + 3, &z, 2); memcpy(evt + 5, &z, 2);
gui_bridge_post(GUI_EVT_AUTO_CONNECT_STATUS, evt, 7);
u_free(ac);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: stopped", AC_ID);
}
void chat_sync_auto_connect_switch_group(struct UTUN_INSTANCE* inst, uint64_t new_group_id) {
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: switch group 0x%016llx", AC_ID, (unsigned long long)new_group_id);
chat_sync_auto_connect_stop();
chat_sync_auto_connect_start(inst);
}
struct channel_cache {
char channel_id[64];
uint64_t* peer_ids;
int peer_count;
uint32_t msg_count;
uint8_t last_chain_hash[32];
uint8_t synced;
};
struct chat_sync {
struct UTUN_INSTANCE* inst;
struct channel_cache* channels;
int channel_count;
void* refresh_timer;
void* info_req_timer;
void* join_timer;
uint8_t initialized;
uint64_t pending_invite_ch_id;
uint64_t pending_invite_node_id;
};
#define CS_ID "chat_sync"
#define CS_INFO_REQ_TIMEOUT_MS 5000
#define CS_JOIN_TIMEOUT_MS 8000
static const char* cs_msg_name(uint8_t type) {
switch (type) {
case CS_MSG_INIT_SYNC: return "INIT_SYNC";
case CS_MSG_INIT_RESP: return "INIT_RESP";
case CS_MSG_SEND_DATA: return "SEND_DATA";
case CS_MSG_PUSH: return "PUSH";
case CS_MSG_ACK_PUSH: return "ACK_PUSH";
case CS_MSG_SYNC_DONE: return "SYNC_DONE";
case CS_MSG_CHANNEL_INFO_REQ: return "CHANNEL_INFO_REQ";
case CS_MSG_CHANNEL_INFO_RESP: return "CHANNEL_INFO_RESP";
case CS_MSG_CHANNEL_JOIN: return "CHANNEL_JOIN";
case CS_MSG_WELCOME: return "WELCOME";
case CS_MSG_PEER_UPSERT: return "PEER_UPSERT";
case CS_MSG_ERROR: return "ERROR";
case CS_MSG_PEER_REMOVE: return "PEER_REMOVE";
default: return "???";
}
}
static void cs_cancel_proto_timers(struct chat_sync* cs) {
if (cs->info_req_timer) { uasync_cancel_timeout(cs->inst->ua, cs->info_req_timer); cs->info_req_timer = NULL; }
if (cs->join_timer) { uasync_cancel_timeout(cs->inst->ua, cs->join_timer); cs->join_timer = NULL; }
}
/* ── Forward decl ── */
static struct ETCP_CONN* cs_find_conn_for_node(struct UTUN_INSTANCE* inst, uint64_t node_id);
/* ── Send ── */
static int cs_send(struct chat_sync* cs, const char* ch_id, uint64_t dst,
const uint8_t* payload, size_t len) {
if (len < 1) return -1;
uint8_t ch_len = (uint8_t)strlen(ch_id);
if (ch_len > 63) return -1;
size_t total = 1 + 1 + ch_len + len;
uint8_t* buf = u_malloc(total);
if (!buf) return -1;
uint8_t* p = buf;
*p++ = ETCP_RT_ID_CHAT_SYNC;
*p++ = ch_len;
memcpy(p, ch_id, ch_len); p += ch_len;
memcpy(p, payload, len);
struct ll_entry* entry = queue_entry_new(0);
if (!entry) { u_free(buf); return -1; }
entry->dgram = buf; entry->len = total;
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: SEND %s to=%016llx ch=%s len=%zu",
CS_ID, cs_msg_name(payload[0]), (unsigned long long)dst, ch_id, len);
struct ETCP_CONN* conn = cs_find_conn_for_node(cs->inst, dst);
if (!conn) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: no conn for node %016llx", CS_ID, (unsigned long long)dst); u_free(buf); queue_entry_free(entry); return -1; }
return etcp_send(conn, entry);
}
/* ── Channel cache ── */
static struct channel_cache* cs_find(struct chat_sync* cs, const char* ch_id) {
for (int i = 0; i < cs->channel_count; i++)
if (strcmp(cs->channels[i].channel_id, ch_id) == 0) return &cs->channels[i];
return NULL;
}
/* ── Forward declarations for new message handlers ── */
static void cs_propagate(struct chat_sync* cs, const char* ch_id, uint64_t exclude_id,
const uint8_t* payload, size_t len);
static void cs_handle_channel_info_req(struct chat_sync* cs, uint64_t peer,
const char* ch_id, const uint8_t* pl, size_t len);
static void cs_handle_channel_info_resp(struct chat_sync* cs, uint64_t peer,
const char* ch_id, const uint8_t* pl, size_t len);
static void cs_handle_channel_join(struct chat_sync* cs, uint64_t peer,
const char* ch_id, const uint8_t* pl, size_t len);
static void cs_handle_welcome(struct chat_sync* cs, uint64_t peer,
const char* ch_id, const uint8_t* pl, size_t len);
static void cs_handle_peer_upsert(struct chat_sync* cs, uint64_t peer,
const char* ch_id, const uint8_t* pl, size_t len);
static void cs_handle_peer_remove(struct chat_sync* cs, uint64_t peer,
const char* ch_id, const uint8_t* pl, size_t len);
/* ── Recv dispatcher ── */
static void chat_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
if (!entry || entry->len < 4) {
if (entry) { if (entry->dgram) u_free(entry->dgram); queue_entry_free(entry); }
return;
}
if (!g_cs || !g_cs->initialized) { u_free(entry->dgram); queue_entry_free(entry); return; }
uint64_t peer = conn ? conn->peer_node_id : 0;
const uint8_t* d = entry->dgram;
size_t dlen = entry->len;
uint8_t ch_len = d[1];
if (dlen < (size_t)(2 + ch_len + 1)) { u_free(entry->dgram); queue_entry_free(entry); return; }
char ch_id[64]; memcpy(ch_id, d + 2, ch_len); ch_id[ch_len] = '\0';
uint8_t type = d[2 + ch_len];
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: recv %s(%02x) from=%016llx ch=%s len=%zu",
CS_ID, cs_msg_name(type), type, (unsigned long long)peer, ch_id, dlen);
const uint8_t* pl = d + 3 + ch_len;
size_t plen = dlen - 3 - ch_len;
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: RECV %s from=%016llx ch=%s len=%zu",
CS_ID, cs_msg_name(type), (unsigned long long)peer, ch_id, plen);
switch (type) {
case CS_MSG_CHANNEL_INFO_REQ: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: → handle INFO_REQ", CS_ID); cs_handle_channel_info_req(g_cs, peer, ch_id, pl, plen); break;
case CS_MSG_CHANNEL_INFO_RESP:DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: → handle INFO_RESP", CS_ID); cs_handle_channel_info_resp(g_cs, peer, ch_id, pl, plen); break;
case CS_MSG_CHANNEL_JOIN: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: → handle JOIN", CS_ID); cs_handle_channel_join(g_cs, peer, ch_id, pl, plen); break;
case CS_MSG_WELCOME: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: → handle WELCOME", CS_ID); cs_handle_welcome(g_cs, peer, ch_id, pl, plen); break;
case CS_MSG_PEER_UPSERT: DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: → handle PEER_UPSERT", CS_ID); cs_handle_peer_upsert(g_cs, peer, ch_id, pl, plen); break;
case CS_MSG_ERROR: {
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; }
uint8_t err[20]; memcpy(err, &peer, 8); int r = CC_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);
g_cs->pending_invite_node_id = 0; g_cs->pending_invite_ch_id = 0;
break;
}
case CS_MSG_PEER_REMOVE: cs_handle_peer_remove(g_cs, peer, ch_id, pl, plen); break;
default: DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: UNKNOWN msg type=%02x from=%016llx", CS_ID, type, (unsigned long long)peer); break;
}
u_free(entry->dgram); queue_entry_free(entry);
}
/* ── Connection callbacks ── */
static void cs_info_req_timeout_cb(void* arg) {
struct chat_sync* cs = (struct chat_sync*)arg;
if (!cs || !cs->initialized) return;
cs->info_req_timer = NULL;
DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: CHANNEL_INFO_REQ timeout peer=%016llx ch=%llu",
CS_ID, (unsigned long long)cs->pending_invite_node_id,
(unsigned long long)cs->pending_invite_ch_id);
uint8_t err[12]; memcpy(err, &cs->pending_invite_node_id, 8);
int r = CC_ERR_TIMEOUT; memcpy(err + 8, &r, 4);
gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 12);
cs->pending_invite_node_id = 0;
cs->pending_invite_ch_id = 0;
}
static void cs_join_timeout_cb(void* arg) {
struct chat_sync* cs = (struct chat_sync*)arg;
if (!cs || !cs->initialized) return;
cs->join_timer = NULL;
uint64_t peer = cs->pending_invite_node_id;
DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: CHANNEL_JOIN timeout peer=%016llx ch=%llu",
CS_ID, (unsigned long long)peer, (unsigned long long)cs->pending_invite_ch_id);
uint8_t err[20]; memcpy(err, &peer, 8);
int r = CC_ERR_TIMEOUT; memcpy(err + 8, &r, 4);
memcpy(err + 12, &cs->pending_invite_ch_id, 8);
gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20);
cs->pending_invite_node_id = 0;
cs->pending_invite_ch_id = 0;
}
/* ── helper: find active ETCP_CONN for node ── */
static struct ETCP_CONN* cs_find_conn_for_node(struct UTUN_INSTANCE* inst, uint64_t node_id) {
if (!inst->connections) { DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: find_conn connections=NULL", CS_ID); return NULL; }
struct ll_entry* e = queue_find_data_by_index(inst->connections, (const uint8_t*)&node_id);
if (e) {
struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data;
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: find_conn node=%016llx found=%p initialized=%d links_up=%d state=%d",
CS_ID, (unsigned long long)node_id, (void*)ce->conn, ce->conn->initialized, ce->conn->links_up, ce->conn->state);
if (ce->conn->initialized && ce->conn->links_up) return ce->conn;
return NULL;
}
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: find_conn node=%016llx NOT FOUND in queue (head=%p count=%d)",
CS_ID, (unsigned long long)node_id, (void*)inst->connections->head, queue_entry_count(inst->connections));
return NULL;
}
/* ── helper: check if peer has active ETCP link ── */
static int cs_is_peer_online(struct UTUN_INSTANCE* inst, uint64_t peer_id) {
if (!inst->connections) return 0;
struct ll_entry* e = queue_find_data_by_index(inst->connections, (const uint8_t*)&peer_id);
if (e) {
struct conn_queue_entry* ce = (struct conn_queue_entry*)e->data;
struct ETCP_LINK* l = ce->conn->links;
while (l) { if (l->initialized && l->link_status) return 1; l = l->next; }
}
return 0;
}
/* ── post per-channel online peers count to GUI ── */
static void cs_post_channel_online(struct chat_sync* cs, const char* ch_id) {
int online = 0;
struct channel_cache* ch = NULL;
for (int i = 0; i < cs->channel_count; i++)
if (strcmp(cs->channels[i].channel_id, ch_id) == 0) { ch = &cs->channels[i]; break; }
if (ch) {
for (int j = 0; j < ch->peer_count; j++)
if (cs_is_peer_online(cs->inst, ch->peer_ids[j])) online++;
}
size_t cl = strlen(ch_id);
if (cl > 63) cl = 63;
uint8_t evt[66]; evt[0] = (uint8_t)cl;
memcpy(evt + 1, ch_id, cl);
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);
}
static void _on_member_sync_done(uint64_t peer, const char* ns, int result, void* arg) {
struct channel_cache* ch = (struct channel_cache*)arg;
if (result == MT_OK && ch) ch->synced = CS_SYNC_DONE;
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: member_sync %s ns=%s peer=%016llx",
CS_ID, result == MT_OK ? "OK" : "FAIL", ns, (unsigned long long)peer);
}
static void cs_on_conn_up(struct ETCP_CONN* conn, void* arg) {
(void)arg;
if (!conn || !g_cs) return;
uint64_t peer = conn->peer_node_id;
if (peer == 0 || peer == g_cs->inst->node_id) return;
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: conn_up peer=%016llx init=%d links=%d", CS_ID, (unsigned long long)peer, conn->initialized, conn->links_up);
if (!conn->initialized) { DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: conn_up SKIP — not initialized peer=%016llx", CS_ID, (unsigned long long)peer); return; }
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: conn_up peer=%016llx pending_invite=%llu links_up=%d", CS_ID, (unsigned long long)peer, (unsigned long long)g_cs->pending_invite_ch_id, conn->links_up);
if (g_cs->pending_invite_ch_id != 0 && !g_cs->info_req_timer) {
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: conn_up invite path peer=%016llx ch=%llu",
CS_ID, (unsigned long long)peer, (unsigned long long)g_cs->pending_invite_ch_id);
if (g_cs->pending_invite_node_id != 0 && g_cs->pending_invite_node_id != peer)
DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: invite node_id MISMATCH: invite=0x%016llx ETCP_peer=0x%016llx — invite is STALE!",
CS_ID, (unsigned long long)g_cs->pending_invite_node_id, (unsigned long long)peer);
g_cs->pending_invite_node_id = peer;
char ch_id_str[64];
snprintf(ch_id_str, sizeof(ch_id_str), "%llu",
(unsigned long long)g_cs->pending_invite_ch_id);
uint8_t req[1] = { CS_MSG_CHANNEL_INFO_REQ };
cs_send(g_cs, ch_id_str, peer, req, 1);
g_cs->info_req_timer = uasync_set_timeout(g_cs->inst->ua,
CS_INFO_REQ_TIMEOUT_MS * 10, g_cs, cs_info_req_timeout_cb, "cs_info_req");
return;
}
member_sync_set_online(g_cs->inst, peer, 1);
for (int i = 0; i < g_cs->channel_count; i++) {
for (int j = 0; j < g_cs->channels[i].peer_count; j++)
if (g_cs->channels[i].peer_ids[j] == peer) { cs_post_channel_online(g_cs, g_cs->channels[i].channel_id); break; }
}
}
static void cs_on_conn_down(struct ETCP_CONN* conn, void* arg) {
(void)arg;
if (!conn || !g_cs) return;
uint64_t peer = conn->peer_node_id;
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: conn_down peer=%016llx", CS_ID, (unsigned long long)peer);
uint16_t rtt = conn->rtt_avg_100;
if (rtt > 0 && peer != 0 && g_cs->inst->topo_sqlite_db) {
sqlite3* db = g_cs->inst->topo_sqlite_db;
sqlite3_stmt* st = NULL;
sqlite3_prepare_v2(db, "UPDATE node_addresses SET rtt=? WHERE node_id=?", -1, &st, NULL);
if (st) { sqlite3_bind_int(st, 1, (int)rtt); sqlite3_bind_int64(st, 2, (sqlite3_int64)peer); sqlite3_step(st); sqlite3_finalize(st); }
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: conn_down saved rtt=%u for node %016llx", CS_ID, rtt, (unsigned long long)peer);
}
cs_cancel_proto_timers(g_cs);
if (g_cs->pending_invite_node_id == peer) {
DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: conn down while waiting invite resp peer=%016llx, timers cancelled, state kept for retry on reconnect",
CS_ID, (unsigned long long)peer);
/* invite state survives connection flaps: only timeout or explicit protocol end clears it */
}
for (int i = 0; i < g_cs->channel_count; i++) {
int found = 0;
for (int j = 0; j < g_cs->channels[i].peer_count; j++)
if (g_cs->channels[i].peer_ids[j] == peer) { found = 1; break; }
if (found) {
g_cs->channels[i].synced = CS_SYNC_NONE;
uint8_t rem[9]; rem[0] = CS_MSG_PEER_REMOVE;
memcpy(rem + 1, &peer, 8);
cs_propagate(g_cs, g_cs->channels[i].channel_id, peer, rem, 9);
cs_post_channel_online(g_cs, g_cs->channels[i].channel_id);
}
}
}
static void cs_on_new_conn(struct ETCP_CONN* conn, void* arg) {
(void)arg;
if (!conn) return;
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: new_conn conn=%p", CS_ID, (void*)conn);
etcp_conn_add_up_cbk(conn, cs_on_conn_up, NULL);
etcp_conn_add_down_cbk(conn, cs_on_conn_down, NULL);
}
/* ── Periodic refresh from DB ── */
static void cs_refresh_channels(struct chat_sync* cs) {
uint8_t buf[4096]; size_t buf_len;
if (chat_core_list_channels(buf, sizeof(buf), &buf_len) != 0) return;
if (buf_len < 2) return;
uint16_t cnt; memcpy(&cnt, buf, 2);
const uint8_t* p = buf + 2; size_t rem = buf_len - 2;
for (int i = 0; i < cs->channel_count; i++)
if (cs->channels[i].peer_ids) u_free(cs->channels[i].peer_ids);
if (cs->channels) u_free(cs->channels);
cs->channels = u_calloc(cnt, sizeof(struct channel_cache));
cs->channel_count = cnt;
if (!cs->channels) { cs->channel_count = 0; return; }
for (int i = 0; i < (int)cnt && rem >= 1; i++) {
uint8_t id_len = *p++; rem--;
if (rem < id_len) break;
memcpy(cs->channels[i].channel_id, p, id_len);
cs->channels[i].channel_id[id_len] = '\0';
p += id_len; rem -= id_len;
cs->channels[i].msg_count = chat_core_count(cs->channels[i].channel_id);
if (cs->channels[i].msg_count > 0)
chat_core_chain_hash_at(cs->channels[i].channel_id,
cs->channels[i].msg_count - 1, cs->channels[i].last_chain_hash);
else
memset(cs->channels[i].last_chain_hash, 0, 32);
uint8_t peer_buf[2048]; size_t peer_len;
if (chat_core_list_peers(cs->channels[i].channel_id, peer_buf,
sizeof(peer_buf), &peer_len) == 0 && peer_len >= 2) {
uint16_t pc; memcpy(&pc, peer_buf, 2);
cs->channels[i].peer_ids = u_calloc(pc, sizeof(uint64_t));
cs->channels[i].peer_count = pc;
for (uint16_t j = 0; j < pc; j++)
memcpy(&cs->channels[i].peer_ids[j], peer_buf + 2 + j * 8, 8);
}
}
}
static void refresh_timer_cb(void* arg) {
(void)arg;
if (!g_cs || !g_cs->initialized) return;
cs_refresh_channels(g_cs);
g_cs->refresh_timer = uasync_set_timeout(g_cs->inst->ua,
30u * 10000u, g_cs, refresh_timer_cb, "cs_refresh");
}
/* ── TTL cleanup ── */
/* ── Public API ── */
int chat_sync_init(struct UTUN_INSTANCE* inst,
void (*gui_cb)(void*, int, const uint8_t*, int)) {
(void)gui_cb;
if (!inst) return -1;
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: init", CS_ID);
struct chat_sync* cs = u_calloc(1, sizeof(*cs));
if (!cs) return -1;
cs->inst = inst;
cs->initialized = 1;
g_cs = cs;
etcp_bind(inst, ETCP_RT_ID_CHAT_SYNC, chat_sync_recv_cb);
etcp_add_new_conn_cbk(inst, cs_on_new_conn, NULL);
{
struct ll_entry* entry = inst->connections->head;
while (entry) {
struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data;
etcp_conn_add_up_cbk(ce->conn, cs_on_conn_up, NULL);
etcp_conn_add_down_cbk(ce->conn, cs_on_conn_down, NULL);
entry = entry->next;
}
}
cs_refresh_channels(cs);
cs->refresh_timer = uasync_set_timeout(inst->ua, 50000u, cs,
refresh_timer_cb, "cs_refresh");
cs->info_req_timer = NULL;
cs->join_timer = NULL;
member_sync_init(inst);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: initialized", CS_ID);
return 0;
}
void chat_sync_destroy(struct UTUN_INSTANCE* inst) {
struct chat_sync* cs = g_cs;
if (!cs || !inst) return;
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: destroy", CS_ID);
chat_sync_auto_connect_stop();
member_sync_destroy(inst);
cs->initialized = 0; g_cs = NULL;
etcp_unbind(inst, ETCP_RT_ID_CHAT_SYNC);
if (cs->refresh_timer) { uasync_cancel_timeout(inst->ua, cs->refresh_timer); cs->refresh_timer = NULL; }
cs_cancel_proto_timers(cs);
{
struct ll_entry* entry = inst->connections->head;
while (entry) {
struct conn_queue_entry* ce = (struct conn_queue_entry*)entry->data;
etcp_conn_remove_up_cbk(ce->conn, cs_on_conn_up, NULL);
etcp_conn_remove_down_cbk(ce->conn, cs_on_conn_down, NULL);
entry = entry->next;
}
}
for (int i = 0; i < cs->channel_count; i++)
if (cs->channels[i].peer_ids) u_free(cs->channels[i].peer_ids);
if (cs->channels) u_free(cs->channels);
u_free(cs);
}
void chat_sync_connect_node(struct UTUN_INSTANCE* inst, uint64_t node_id) {
(void)inst;
chat_core_connect_auto(node_id, NULL, NULL, NULL);
}
void chat_sync_connect_from_invite(uint64_t channel_id, uint64_t node_id,
const uint8_t* pubkey_bin,
const uint8_t* addrs_data, int addr_count) {
if (!g_cs || !g_cs->inst || !g_cs->inst->ua) {
int r = -7;
uint8_t err[12]; memcpy(err, &node_id, 8); memcpy(err + 8, &r, 4);
gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 12);
return;
}
struct chat_invite* inv = u_calloc(1, sizeof(struct chat_invite));
if (!inv) return;
inv->channel_id = channel_id;
inv->node_id = node_id;
memcpy(inv->pubkey, pubkey_bin, 32);
size_t addrs_sz = (size_t)addr_count * 7;
inv->addrs_data = u_malloc(addrs_sz);
if (!inv->addrs_data) { u_free(inv); return; }
memcpy(inv->addrs_data, addrs_data, addrs_sz);
inv->addr_count = addr_count;
g_cs->pending_invite_ch_id = channel_id;
g_cs->pending_invite_node_id = node_id;
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: invite ch=%llu node=%016llx pubkey=%016llx addrs=%d",
CS_ID, channel_id, node_id, *(const uint64_t*)pubkey_bin, addr_count);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: invite start ch=%llu node=0x%016llx pubkey=%016llx... addrs=%d",
CS_ID, channel_id, node_id, *(const uint64_t*)pubkey_bin, addr_count);
gui_bridge_post_uasync_fn(
(void(*)(void*))chat_core_connect_from_invite, inv);
}
/* ─── Ed25519 sign / verify helpers ─── */
static int cs_ed25519_sign(const uint8_t* privkey, const uint8_t* msg, size_t msg_len,
uint8_t* sig_out) {
EVP_PKEY* pkey = EVP_PKEY_new_raw_private_key(EVP_PKEY_ED25519, NULL, privkey, SC_PRIVKEY_SIZE);
if (!pkey) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: EVP_PKEY_new failed", CS_ID); return -1; }
EVP_MD_CTX* mdctx = EVP_MD_CTX_new();
if (!mdctx) { EVP_PKEY_free(pkey); return -1; }
int ok = (EVP_DigestSignInit(mdctx, NULL, NULL, NULL, pkey) == 1)
&& (EVP_DigestSign(mdctx, sig_out, &(size_t){64}, msg, msg_len) == 1);
EVP_MD_CTX_free(mdctx);
EVP_PKEY_free(pkey);
return ok ? 0 : -1;
}
static int cs_ed25519_verify(const uint8_t* pubkey, const uint8_t* msg, size_t msg_len,
const uint8_t* sig) {
EVP_PKEY* pkey = EVP_PKEY_new_raw_public_key(EVP_PKEY_ED25519, NULL, pubkey, SC_PUBKEY_SIZE);
if (!pkey) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: EVP_PKEY_new pub failed", CS_ID); return -1; }
EVP_MD_CTX* mdctx = EVP_MD_CTX_new();
if (!mdctx) { EVP_PKEY_free(pkey); return -1; }
int rc = EVP_DigestVerifyInit(mdctx, NULL, NULL, NULL, pkey)
&& EVP_DigestVerify(mdctx, sig, 64, msg, msg_len) == 1;
EVP_MD_CTX_free(mdctx);
EVP_PKEY_free(pkey);
return rc ? 0 : -1;
}
/* ─── Propagation helper ─── */
static void cs_propagate(struct chat_sync* cs, const char* ch_id, uint64_t exclude_id,
const uint8_t* payload, size_t len) {
struct channel_cache* ch = cs_find(cs, ch_id);
if (!ch) return;
uint64_t myid = cs->inst->node_id;
for (int i = 0; i < ch->peer_count; i++) {
uint64_t pid = ch->peer_ids[i];
if (pid == exclude_id || pid == myid) continue;
if (!cs_find_conn_for_node(cs->inst, pid)) continue;
cs_send(cs, ch_id, pid, payload, len);
}
}
static int _get_node_name(sqlite3* db, uint64_t node_id, char* out, size_t sz) {
if (!db || !out || !sz) return -1;
sqlite3_stmt* st = NULL;
sqlite3_prepare_v2(db, "SELECT name FROM nodes WHERE node_id=?", -1, &st, NULL);
if (!st) return -1;
sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id);
int rc = -1;
if (sqlite3_step(st) == SQLITE_ROW) {
const unsigned char* nm = sqlite3_column_text(st, 0);
if (nm) { snprintf(out, sz, "%s", nm); rc = 0; }
else { out[0] = '\0'; rc = 0; }
}
sqlite3_finalize(st);
return rc;
}
/* ─── CHANNEL_INFO_REQ (0x09): joiner → inviter ─── */
static void cs_handle_channel_info_req(struct chat_sync* cs, uint64_t peer,
const char* ch_id, const uint8_t* pl, size_t len) {
(void)pl; (void)len;
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: INFO_REQ ch=%s from=%016llx", CS_ID, ch_id, (unsigned long long)peer);
char name[128]; int is_dm; uint64_t owner; uint8_t x25519[32], ed_pub[32], ch_sig[64];
if (topo_node_sqlite_channel_get(cs->inst->topo_sqlite_db,
ch_id, name, (int)sizeof(name), &is_dm, &owner, x25519, ed_pub, ch_sig) != 0) {
DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: CHANNEL_INFO_REQ unknown ch=%s from=%016llx", CS_ID, ch_id, (unsigned long long)peer);
uint8_t err[1] = { CS_MSG_ERROR }; cs_send(cs, ch_id, peer, err, 1);
return;
}
uint64_t myid = cs->inst->node_id;
if (memcmp(cs->inst->my_keys.public_key, x25519, 32) != 0)
DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: CHANNEL_INFO_RESP pubkey MISMATCH: my=%016llx ch=%016llx — channel was created with DIFFERENT keys!",
CS_ID, *(const uint64_t*)cs->inst->my_keys.public_key, *(const uint64_t*)x25519);
/* load or create join_sig */
uint8_t my_join_sig[64] = {0}; uint64_t my_join_ts = 0;
{
sqlite3* vdb = cs->inst->topo_sqlite_db;
if (topo_node_sqlite_member_get_join(vdb, ch_id, myid, my_join_sig, &my_join_ts) != 0) {
uint8_t join_msg[256]; size_t mlen = 0;
mlen += snprintf((char*)join_msg + mlen, sizeof(join_msg) - mlen, "%s", ch_id) + 1;
memcpy(join_msg + mlen, &myid, 8); mlen += 8;
memcpy(join_msg + mlen, cs->inst->my_keys.public_key, 32); mlen += 32;
my_join_ts = (uint64_t)ntp_time_get_seconds(cs->inst);
memcpy(join_msg + mlen, &my_join_ts, 8); mlen += 8;
cs_ed25519_sign(cs->inst->my_ed25519_privkey, join_msg, mlen, my_join_sig);
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: CHANNEL_INFO_RESP created NEW join_sig: myid=%016llx ts=%llu",
CS_ID, (unsigned long long)myid, (unsigned long long)my_join_ts);
}
}
/* generate update_sig */
uint8_t my_update_sig[64] = {0}; uint64_t my_update_ts;
{
uint8_t umsg[256]; size_t ulen = 0;
ulen += snprintf((char*)umsg + ulen, sizeof(umsg) - ulen, "%s", ch_id) + 1;
memcpy(umsg + ulen, &myid, 8); ulen += 8;
memcpy(umsg + ulen, cs->inst->my_keys.public_key, 32); ulen += 32;
my_update_ts = (uint64_t)ntp_time_get_seconds(cs->inst);
memcpy(umsg + ulen, &my_update_ts, 8); ulen += 8;
memcpy(umsg + ulen, my_join_sig, 64); ulen += 64;
cs_ed25519_sign(cs->inst->my_ed25519_privkey, umsg, ulen, my_update_sig);
}
uint8_t buf[1024]; size_t boff = 0;
buf[boff++] = CS_MSG_CHANNEL_INFO_RESP;
uint8_t nl = (uint8_t)strlen(name);
buf[boff++] = nl; memcpy(buf + boff, name, nl); boff += nl;
memcpy(buf + boff, &owner, 8); boff += 8;
buf[boff++] = (uint8_t)is_dm;
memcpy(buf + boff, x25519, 32); boff += 32;
memcpy(buf + boff, ed_pub, 32); boff += 32;
memcpy(buf + boff, ch_sig, 64); boff += 64;
uint8_t inv_flags = (my_join_sig[0] || my_join_ts) ? PEERS_FLAG_HAS_JOIN : 0;
buf[boff++] = inv_flags;
if (inv_flags & PEERS_FLAG_HAS_JOIN) {
memcpy(buf + boff, my_join_sig, 64); boff += 64;
memcpy(buf + boff, &my_join_ts, 8); boff += 8;
}
memcpy(buf + boff, my_update_sig, 64); boff += 64;
memcpy(buf + boff, &my_update_ts, 8); boff += 8;
{ const char* iname = cs->inst->name[0] ? cs->inst->name : "";
uint8_t il = (uint8_t)strlen(iname);
buf[boff++] = il; memcpy(buf + boff, iname, il); boff += il; }
cs_send(cs, ch_id, peer, buf, boff);
}
/* ─── CHANNEL_INFO_RESP (0x0A): inviter → joiner ─── */
static void cs_handle_channel_info_resp(struct chat_sync* cs, uint64_t peer,
const char* ch_id, const uint8_t* pl, size_t len) {
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: INFO_RESP ch=%s from=%016llx len=%zu", CS_ID, ch_id, (unsigned long long)peer, len);
if (cs->info_req_timer) { uasync_cancel_timeout(cs->inst->ua, cs->info_req_timer); cs->info_req_timer = NULL; }
if (len < 1) { DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: CHANNEL_INFO_RESP too short len=%zu peer=%016llx", CS_ID, len, (unsigned long long)peer); return; }
uint8_t nl = pl[0];
if (1 + nl + 8 + 1 + 32 + 32 + 64 + 1 + 64 + 8 + 1 > len) { DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: CHANNEL_INFO_RESP truncated len=%zu", CS_ID, len); return; }
const uint8_t* p = pl + 1;
char name[128]; memcpy(name, p, nl); name[nl] = '\0'; p += nl;
uint64_t owner; memcpy(&owner, p, 8); p += 8;
uint8_t is_dm = *p++;
const uint8_t* x25519 = p; p += 32;
const uint8_t* ed_pub = p; p += 32;
const uint8_t* ch_sig = p; p += 64;
uint8_t inv_flags = *p++;
const uint8_t* inviter_join_sig = NULL; uint64_t inviter_join_ts = 0;
if (inv_flags & PEERS_FLAG_HAS_JOIN) {
if (p + 72 > pl + len) return; inviter_join_sig = p; p += 64; memcpy(&inviter_join_ts, p, 8); p += 8; }
const uint8_t* inviter_update_sig = p; p += 64; uint64_t inviter_update_ts;
memcpy(&inviter_update_ts, p, 8); p += 8;
uint8_t inv_name_len = *p++;
char inv_name[128] = "";
if (p + inv_name_len <= pl + len) { memcpy(inv_name, p, inv_name_len); inv_name[inv_name_len] = '\0'; }
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: CHANNEL_INFO_RESP from peer=%016llx owner=%016llx ch_x25519=%016llx ch_ed=%016llx name=%s",
CS_ID, (unsigned long long)peer, (unsigned long long)owner,
*(const uint64_t*)x25519, *(const uint64_t*)ed_pub, name);
/* verify channel signature */
uint8_t vmsg[1024]; size_t vlen = 0;
vlen += snprintf((char*)vmsg + vlen, sizeof(vmsg) - vlen, "%s", ch_id) + 1;
vlen += snprintf((char*)vmsg + vlen, sizeof(vmsg) - vlen, "%s", name) + 1;
memcpy(vmsg + vlen, &owner, 8); vlen += 8;
memcpy(vmsg + vlen, x25519, 32); vlen += 32;
memcpy(vmsg + vlen, ed_pub, 32); vlen += 32;
if (cs_ed25519_verify(ed_pub, vmsg, vlen, ch_sig) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: CHANNEL_INFO_RESP invalid ch_sig ch=%s", CS_ID, ch_id);
return;
}
/* save channel to local DB */
topo_node_sqlite_channel_put(cs->inst->topo_sqlite_db,
ch_id, name, (int)is_dm, owner, x25519, NULL, ed_pub, NULL, ch_sig);
chat_core_ensure_channel_ready(ch_id);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: channel ready for sync ch=%s, db_sync will pick up via periodic check or active conn",
CS_ID, ch_id);
/* verify inviter's join_sig AND save node info using ETCP-authenticated keys */
{ struct ETCP_CONN* inv_conn = cs_find_conn_for_node(cs->inst, peer);
if (!inv_conn) {
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: CHANNEL_INFO_RESP no ETCP conn for inviter peer=%016llx", CS_ID, (unsigned long long)peer);
return;
}
const uint8_t* inv_x25519 = inv_conn->crypto_ctx.peer_public_key;
const uint8_t* inv_ed = inv_conn->peer_ed25519_pubkey;
uint8_t ivmsg[256]; size_t ilen = 0;
ilen += snprintf((char*)ivmsg + ilen, sizeof(ivmsg) - ilen, "%s", ch_id) + 1;
memcpy(ivmsg + ilen, &peer, 8); ilen += 8;
memcpy(ivmsg + ilen, inv_x25519, 32); ilen += 32;
memcpy(ivmsg + ilen, &inviter_update_ts, 8); ilen += 8;
{ const char* nm = inv_name[0] ? inv_name : "";
size_t nl2 = strlen(nm); (void)nl2; }
if (inviter_join_sig) memcpy(ivmsg + ilen, inviter_join_sig, 64); else memset(ivmsg + ilen, 0, 64);
ilen += 64;
if (cs_ed25519_verify(inv_ed, ivmsg, ilen, inviter_update_sig) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: CHANNEL_INFO_RESP invalid inviter_update_sig peer=%016llx ts=%llu",
CS_ID, (unsigned long long)peer, (unsigned long long)inviter_update_ts);
}
/* save inviter node_info to local DB */
sqlite3* vdb = cs->inst->topo_sqlite_db;
if (vdb && inv_name[0]) {
sqlite3_stmt* ns = NULL;
sqlite3_prepare_v2(vdb, "INSERT OR REPLACE INTO nodes(node_id,name,x25519_pubkey,ed25519_pubkey) VALUES(?,?,?,?)", -1, &ns, NULL);
if (ns) { sqlite3_bind_int64(ns, 1, (sqlite3_int64)peer);
sqlite3_bind_text(ns, 2, inv_name, -1, SQLITE_STATIC);
sqlite3_bind_blob(ns, 3, inv_x25519, 32, SQLITE_STATIC);
sqlite3_bind_blob(ns, 4, inv_ed, 32, SQLITE_STATIC);
sqlite3_step(ns); sqlite3_finalize(ns); }
}
/* save inviter's address from ETCP connection */
{ struct ETCP_LINK* lk = inv_conn->links;
while (lk) {
if (lk->initialized && lk->remote_addr.ss_family == AF_INET) {
struct sockaddr_in* sin = (struct sockaddr_in*)&lk->remote_addr;
uint8_t ip[4]; memcpy(ip, &sin->sin_addr, 4); uint16_t port = ntohs(sin->sin_port);
char sql[256]; snprintf(sql, sizeof(sql),
"INSERT OR REPLACE INTO node_addresses(node_id,family,address,port,addr_type) VALUES(?,4,?,?,0)");
sqlite3_stmt* as = NULL;
if (sqlite3_prepare_v2(vdb, sql, -1, &as, NULL) == SQLITE_OK) {
sqlite3_bind_int64(as, 1, (sqlite3_int64)peer);
sqlite3_bind_blob(as, 2, ip, 4, SQLITE_STATIC);
sqlite3_bind_int(as, 3, port);
sqlite3_step(as); sqlite3_finalize(as);
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: saved inviter addr peer=%016llx %d.%d.%d.%d:%d", CS_ID,
(unsigned long long)peer, ip[0], ip[1], ip[2], ip[3], port);
}
}
lk = lk->next;
}
}
int rc = topo_node_sqlite_member_put(vdb, ch_id, peer, inviter_join_sig, inviter_join_ts, inviter_update_sig, inviter_update_ts, inv_x25519, inv_ed, inv_name, inviter_join_sig);
if (rc != 0) DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: member_put(inviter) FAILED ch=%s peer=%016llx rc=%d", CS_ID, ch_id, (unsigned long long)peer, rc);
topo_node_sqlite_node_update_verified(vdb, peer, inv_name, inv_x25519, inv_ed, inviter_join_ts);
}
/* generate our own join_sig */
uint64_t myid = cs->inst->node_id;
uint8_t my_x25519[32];
memcpy(my_x25519, cs->inst->my_keys.public_key, 32);
uint8_t join_sig[64]; uint64_t join_ts;
{
uint8_t msg[256]; size_t mlen = 0;
mlen += snprintf((char*)msg + mlen, sizeof(msg) - mlen, "%s", ch_id) + 1;
memcpy(msg + mlen, &myid, 8); mlen += 8;
memcpy(msg + mlen, my_x25519, 32); mlen += 32;
{ const char* nm = cs->inst->name[0] ? cs->inst->name : "";
size_t nl = strlen(nm); memcpy(msg + mlen, nm, nl); mlen += nl; msg[mlen++] = '\0'; }
join_ts = (uint64_t)ntp_time_get_seconds(cs->inst);
memcpy(msg + mlen, &join_ts, 8); mlen += 8;
cs_ed25519_sign(cs->inst->my_ed25519_privkey, msg, mlen, join_sig);
}
/* send JOIN_CHANNEL */
uint8_t jbuf[512]; size_t joff = 0;
jbuf[joff++] = CS_MSG_CHANNEL_JOIN;
memcpy(jbuf + joff, &myid, 8); joff += 8;
memcpy(jbuf + joff, my_x25519, 32); joff += 32;
memcpy(jbuf + joff, cs->inst->my_ed25519_pubkey, 32); joff += 32;
memcpy(jbuf + joff, join_sig, 64); joff += 64;
memcpy(jbuf + joff, &join_ts, 8); joff += 8;
{ const char* nm = cs->inst->name[0] ? cs->inst->name : "";
uint8_t nml = (uint8_t)strlen(nm);
jbuf[joff++] = nml; memcpy(jbuf + joff, nm, nml); joff += nml; }
/* add local addresses */
uint8_t addr_cnt = 0; size_t ac_pos = joff;
jbuf[joff++] = 0;
struct ETCP_SOCKET* sock = cs->inst->etcp_sockets;
while (sock && addr_cnt < 255) {
struct sockaddr_storage* sa;
if (sock->nat_addr.ss_family) sa = &sock->nat_addr;
else if (sock->interface_addr.ss_family) sa = &sock->interface_addr;
else { sock = sock->next; continue; }
if (sa->ss_family == AF_INET) {
struct sockaddr_in* sin = (struct sockaddr_in*)sa;
if (joff + 7 > sizeof(jbuf)) break;
jbuf[joff++] = 4;
memcpy(jbuf + joff, &sin->sin_addr, 4); joff += 4;
uint16_t p = ntohs(sin->sin_port);
jbuf[joff++] = (uint8_t)((p >> 8) & 0xFF);
jbuf[joff++] = (uint8_t)(p & 0xFF);
addr_cnt++;
} else if (sa->ss_family == AF_INET6) {
struct sockaddr_in6* sin6 = (struct sockaddr_in6*)sa;
if (joff + 19 > sizeof(jbuf)) break;
jbuf[joff++] = 6;
memcpy(jbuf + joff, &sin6->sin6_addr, 16); joff += 16;
uint16_t p = ntohs(sin6->sin6_port);
jbuf[joff++] = (uint8_t)((p >> 8) & 0xFF);
jbuf[joff++] = (uint8_t)(p & 0xFF);
addr_cnt++;
}
sock = sock->next;
}
jbuf[ac_pos] = addr_cnt;
cs_send(cs, ch_id, peer, jbuf, joff);
cs->join_timer = uasync_set_timeout(cs->inst->ua,
CS_JOIN_TIMEOUT_MS * 10, cs, cs_join_timeout_cb, "cs_join");
cs->pending_invite_node_id = 0;
cs->pending_invite_ch_id = 0;
}
/* ─── CHANNEL_JOIN (0x0B): joiner → inviter ─── */
static void cs_handle_channel_join(struct chat_sync* cs, uint64_t peer,
const char* ch_id, const uint8_t* pl, size_t len) {
if (len < 8 + 32 + 32 + 64 + 8 + 1 + 1) { DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: CHANNEL_JOIN too short len=%zu peer=%016llx", CS_ID, len, (unsigned long long)peer); return; }
const uint8_t* p = pl;
uint64_t node_id; memcpy(&node_id, p, 8); p += 8;
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: JOIN ch=%s node=%016llx from=%016llx", CS_ID, ch_id, (unsigned long long)node_id, (unsigned long long)peer);
const uint8_t* x25519 = p; p += 32;
const uint8_t* ed_pub = p; p += 32;
const uint8_t* join_sig = p; p += 64;
uint64_t join_ts; memcpy(&join_ts, p, 8); p += 8;
uint8_t name_len = *p++;
char joiner_name[128] = "";
if (p + name_len <= pl + len) { memcpy(joiner_name, p, name_len); joiner_name[name_len] = '\0'; p += name_len; }
uint8_t addr_cnt = *p++;
/* verify join_sig */
uint8_t vmsg[256]; size_t vlen = 0;
vlen += snprintf((char*)vmsg + vlen, sizeof(vmsg) - vlen, "%s", ch_id) + 1;
memcpy(vmsg + vlen, &node_id, 8); vlen += 8;
memcpy(vmsg + vlen, x25519, 32); vlen += 32;
{ const char* nm = joiner_name[0] ? joiner_name : "";
size_t nl = strlen(nm); memcpy(vmsg + vlen, nm, nl); vlen += nl; vmsg[vlen++] = '\0'; }
memcpy(vmsg + vlen, &join_ts, 8); vlen += 8;
if (cs_ed25519_verify(ed_pub, vmsg, vlen, join_sig) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: JOIN invalid sig node=0x%016llx ch=%s — signed(node=%016llx x25519=%016llx ed=%016llx name=%s ts=%llu)", CS_ID,
(unsigned long long)node_id, ch_id,
(unsigned long long)node_id, *(const uint64_t*)x25519, *(const uint64_t*)ed_pub, joiner_name, (unsigned long long)join_ts);
return;
}
sqlite3* db = cs->inst->topo_sqlite_db;
int rc = topo_node_sqlite_member_put(db, ch_id, node_id, join_sig, join_ts, NULL, 0, x25519, ed_pub, joiner_name, NULL);
if (rc != 0) DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: member_put(joiner) FAILED ch=%s node=0x%016llx rc=%d", CS_ID, ch_id, (unsigned long long)node_id, rc);
topo_node_sqlite_node_update_verified(db, node_id, joiner_name, x25519, ed_pub, join_ts);
/* save joiner node_info to local DB */
if (db && joiner_name[0]) {
sqlite3_stmt* ns = NULL;
sqlite3_prepare_v2(db, "INSERT OR REPLACE INTO nodes(node_id,name,x25519_pubkey,ed25519_pubkey) VALUES(?,?,?,?)", -1, &ns, NULL);
if (ns) { sqlite3_bind_int64(ns, 1, (sqlite3_int64)node_id);
sqlite3_bind_text(ns, 2, joiner_name, -1, SQLITE_STATIC);
sqlite3_bind_blob(ns, 3, x25519, 32, SQLITE_STATIC);
sqlite3_bind_blob(ns, 4, ed_pub, 32, SQLITE_STATIC);
sqlite3_step(ns); sqlite3_finalize(ns); }
}
/* save/update node addresses */
for (uint8_t i = 0; i < addr_cnt && p + 1 <= pl + len; i++) {
uint8_t fm = *p++;
int ip_len = (fm == 4) ? 4 : 16;
if (p + ip_len + 2 > pl + len) break;
const uint8_t* ip = p; p += ip_len;
uint16_t port = ((uint16_t)p[0] << 8) | p[1]; p += 2;
char sql[256]; snprintf(sql, sizeof(sql),
"INSERT OR REPLACE INTO node_addresses(node_id,family,address,port,addr_type)"
" VALUES(?,?,?,?,0)");
sqlite3_stmt* stmt = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &stmt, NULL) == SQLITE_OK) {
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)node_id);
sqlite3_bind_int(stmt, 2, fm);
sqlite3_bind_blob(stmt, 3, ip, ip_len, SQLITE_STATIC);
sqlite3_bind_int(stmt, 4, (int)port);
sqlite3_step(stmt); sqlite3_finalize(stmt);
}
}
/* build WELCOME with all current peers */
uint8_t wbuf[8192];
{
uint8_t peers_data[8192]; size_t peers_len = 0;
if (topo_node_sqlite_channel_peers_all(db, ch_id,
peers_data, sizeof(peers_data), &peers_len) == 0) {
size_t woff = 0;
wbuf[woff++] = CS_MSG_WELCOME;
if (woff + peers_len <= sizeof(wbuf))
{ memcpy(wbuf + woff, peers_data, peers_len); woff += peers_len; }
cs_send(cs, ch_id, peer, wbuf, woff);
}
}
/* propagate PEER_UPSERT to other channel members */
{
uint8_t ubuf[512]; size_t uoff = 0;
ubuf[uoff++] = CS_MSG_PEER_UPSERT;
memcpy(ubuf + uoff, &node_id, 8); uoff += 8;
memcpy(ubuf + uoff, x25519, 32); uoff += 32;
memcpy(ubuf + uoff, ed_pub, 32); uoff += 32;
memcpy(ubuf + uoff, join_sig, 64); uoff += 64;
memcpy(ubuf + uoff, &join_ts, 8); uoff += 8;
ubuf[uoff++] = name_len;
memcpy(ubuf + uoff, joiner_name, name_len); uoff += name_len;
ubuf[uoff++] = addr_cnt;
size_t addr_data_sz = (size_t)(p - (pl + 8 + 32 + 32 + 64 + 8 + 1 + name_len + 1));
if (uoff + addr_data_sz <= sizeof(ubuf)) {
memcpy(ubuf + uoff, pl + 8 + 32 + 32 + 64 + 8 + 1 + name_len + 1, addr_data_sz);
uoff += addr_data_sz;
}
cs_propagate(cs, ch_id, peer, ubuf, uoff);
}
/* add to channel cache */
struct channel_cache* ch = cs_find(cs, ch_id);
if (!ch) {
cs_refresh_channels(cs);
ch = cs_find(cs, ch_id);
}
{ 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); }
/* initiate member sync with the joiner (server side) */
member_sync_start(cs->inst, node_id, ch_id, NULL, NULL);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: JOIN starting member_sync with joiner node=0x%016llx ch=%s",
CS_ID, (unsigned long long)node_id, ch_id);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: JOIN accepted node=0x%016llx ch=%s addrs=%d",
CS_ID, (unsigned long long)node_id, ch_id, addr_cnt);
}
/* ─── WELCOME (0x0C): inviter → joiner with full peer list ─── */
static void cs_handle_welcome(struct chat_sync* cs, uint64_t peer,
const char* ch_id, const uint8_t* pl, size_t len) {
if (cs->join_timer) { uasync_cancel_timeout(cs->inst->ua, cs->join_timer); cs->join_timer = NULL; }
(void)peer;
if (len < 2) { DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: WELCOME too short len=%zu peer=%016llx", CS_ID, len, (unsigned long long)peer); return; }
sqlite3* db = cs->inst->topo_sqlite_db;
const uint8_t* p = pl;
uint16_t pc; memcpy(&pc, p, 2); p += 2;
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: WELCOME ch=%s from=%016llx peers=%u", CS_ID, ch_id, (unsigned long long)peer, pc);
for (uint16_t i = 0; i < pc; i++) {
if ((size_t)(p - pl) + 8 + 32 + 32 + 64 + 8 + 1 + 1 > len) break;
uint64_t node_id; memcpy(&node_id, p, 8); p += 8;
const uint8_t* x25519 = p; p += 32;
const uint8_t* ed_pub = p; p += 32;
const uint8_t* join_sig = p; p += 64;
uint64_t join_ts; memcpy(&join_ts, p, 8); p += 8;
uint8_t nl = *p++;
char peer_name[256] = "";
if (nl && p + nl <= pl + len) { memcpy(peer_name, p, nl); peer_name[nl] = '\0'; p += nl; }
uint8_t ac = *p++;
/* verify join_sig — all fields present in WELCOME wire */
uint8_t vmsg[256]; size_t vlen = 0;
vlen += snprintf((char*)vmsg + vlen, sizeof(vmsg) - vlen, "%s", ch_id) + 1;
memcpy(vmsg + vlen, &node_id, 8); vlen += 8;
memcpy(vmsg + vlen, x25519, 32); vlen += 32;
{ const char* nm = peer_name[0] ? peer_name : "";
size_t nls = strlen(nm); memcpy(vmsg + vlen, nm, nls); vlen += nls; vmsg[vlen++] = '\0'; }
memcpy(vmsg + vlen, &join_ts, 8); vlen += 8;
{ EVP_PKEY* pkey = EVP_PKEY_new_raw_public_key(EVP_PKEY_ED25519, NULL, ed_pub, 32);
if (pkey) {
EVP_MD_CTX* ver = EVP_MD_CTX_new();
if (ver) {
if (EVP_DigestVerifyInit(ver, NULL, NULL, NULL, pkey) != 1
|| EVP_DigestVerify(ver, join_sig, 64, vmsg, vlen) != 1)
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: WELCOME invalid join_sig node=0x%016llx", CS_ID, (unsigned long long)node_id);
EVP_MD_CTX_free(ver);
}
EVP_PKEY_free(pkey);
}
}
int rc = topo_node_sqlite_member_put(db, ch_id, node_id, join_sig, join_ts, NULL, 0, x25519, ed_pub, peer_name, NULL);
if (rc != 0) DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: member_put(welcome) FAILED ch=%s node=0x%016llx rc=%d", CS_ID, ch_id, (unsigned long long)node_id, rc);
topo_node_sqlite_node_update_verified(db, node_id, peer_name, x25519, ed_pub, join_ts);
for (uint8_t j = 0; j < ac; j++) {
if (p + 1 > pl + len) break;
uint8_t fm = *p++;
int ip_len = (fm == 4) ? 4 : 16;
if (p + ip_len + 2 > pl + len) break;
const uint8_t* ip = p; p += ip_len;
uint16_t port = ((uint16_t)p[0] << 8) | p[1]; p += 2;
char sql[256]; snprintf(sql, sizeof(sql),
"INSERT OR REPLACE INTO node_addresses(node_id,family,address,port,addr_type)"
" VALUES(?,?,?,?,0)");
sqlite3_stmt* stmt = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &stmt, NULL) == SQLITE_OK) {
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)node_id);
sqlite3_bind_int(stmt, 2, fm);
sqlite3_bind_blob(stmt, 3, ip, ip_len, SQLITE_STATIC);
sqlite3_bind_int(stmt, 4, (int)port);
sqlite3_step(stmt); sqlite3_finalize(stmt);
}
member_sync_cancel(g_cs->inst, peer, ch_id);
member_sync_set_online(g_cs->inst, peer, 0);
}
}
cs_refresh_channels(cs);
/* add self to local peers table (server added joiner to its DB; client must do the same) */
{
uint64_t myid = cs->inst->node_id;
const uint8_t* my_x25 = cs->inst->my_keys.public_key;
const uint8_t* my_ed = cs->inst->my_ed25519_pubkey;
const char* my_name = cs->inst->name[0] ? cs->inst->name : "";
int src = topo_node_sqlite_member_put(db, ch_id, myid, NULL, 0, NULL, 0, my_x25, my_ed, my_name, NULL);
if (src != 0)
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: member_put(self) FAILED in WELCOME ch=%s rc=%d", CS_ID, ch_id, src);
else
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: WELCOME added self to peers ch=%s node=0x%016llx", CS_ID, ch_id, (unsigned long long)myid);
}
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);
gui_bridge_post(GUI_EVT_CHANNEL_UPDATED, evt, 1 + ch_id_len);
gui_bridge_post(GUI_EVT_MEMBERS_CHANGED, evt, 1 + ch_id_len);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: WELCOME processed ch=%s peers=%d — message sync triggered via db_sync (active conn or peer_check timer every 5s)",
CS_ID, ch_id, pc);
}
/* ─── PEER_UPSERT (0x0D): propagate new peer ─── */
static void cs_handle_peer_upsert(struct chat_sync* cs, uint64_t peer,
const char* ch_id, const uint8_t* pl, size_t len) {
if (len < 8 + 32 + 32 + 64 + 8 + 1 + 1 + 1) { DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: PEER_UPSERT too short len=%zu", CS_ID, len); return; }
const uint8_t* p = pl;
uint64_t node_id; memcpy(&node_id, p, 8); p += 8;
const uint8_t* x25519 = p; p += 32;
const uint8_t* ed_pub = p; p += 32;
const uint8_t* join_sig = p; p += 64;
uint64_t join_ts; memcpy(&join_ts, p, 8); p += 8;
uint8_t name_len = *p++;
char peer_name[128] = "";
if (p + name_len <= pl + len) { memcpy(peer_name, p, name_len); peer_name[name_len] = '\0'; p += name_len; }
uint8_t ac = *p++;
/* verify join_sig */
uint8_t vmsg[256]; size_t vlen = 0;
vlen += snprintf((char*)vmsg + vlen, sizeof(vmsg) - vlen, "%s", ch_id) + 1;
memcpy(vmsg + vlen, &node_id, 8); vlen += 8;
memcpy(vmsg + vlen, x25519, 32); vlen += 32;
{ const char* nm = peer_name[0] ? peer_name : "";
size_t nl = strlen(nm); memcpy(vmsg + vlen, nm, nl); vlen += nl; vmsg[vlen++] = '\0'; }
memcpy(vmsg + vlen, &join_ts, 8); vlen += 8;
if (cs_ed25519_verify(ed_pub, vmsg, vlen, join_sig) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: PEER_UPSERT invalid sig node=0x%016llx", CS_ID,
(unsigned long long)node_id);
return;
}
sqlite3* db = cs->inst->topo_sqlite_db;
int rc = topo_node_sqlite_member_put(db, ch_id, node_id, join_sig, join_ts, NULL, 0, x25519, ed_pub, peer_name, NULL);
if (rc != 0) DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: member_put(peer_upsert) FAILED ch=%s node=0x%016llx rc=%d", CS_ID, ch_id, (unsigned long long)node_id, rc);
topo_node_sqlite_node_update_verified(db, node_id, peer_name, x25519, ed_pub, join_ts);
if (db && peer_name[0]) {
sqlite3_stmt* ns = NULL;
sqlite3_prepare_v2(db, "INSERT OR REPLACE INTO nodes(node_id,name,x25519_pubkey,ed25519_pubkey) VALUES(?,?,?,?)", -1, &ns, NULL);
if (ns) { sqlite3_bind_int64(ns, 1, (sqlite3_int64)node_id);
sqlite3_bind_text(ns, 2, peer_name, -1, SQLITE_STATIC);
sqlite3_bind_blob(ns, 3, x25519, 32, SQLITE_STATIC);
sqlite3_bind_blob(ns, 4, ed_pub, 32, SQLITE_STATIC);
sqlite3_step(ns); sqlite3_finalize(ns); }
}
for (uint8_t i = 0; i < ac; i++) {
if (p + 1 > pl + len) break;
uint8_t fm = *p++;
int ip_len = (fm == 4) ? 4 : 16;
if (p + ip_len + 2 > pl + len) break;
const uint8_t* ip = p; p += ip_len;
uint16_t port = ((uint16_t)p[0] << 8) | p[1]; p += 2;
char sql[256]; snprintf(sql, sizeof(sql),
"INSERT OR REPLACE INTO node_addresses(node_id,family,address,port,addr_type)"
" VALUES(?,?,?,?,0)");
sqlite3_stmt* stmt = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &stmt, NULL) == SQLITE_OK) {
sqlite3_bind_int64(stmt, 1, (sqlite3_int64)node_id);
sqlite3_bind_int(stmt, 2, fm);
sqlite3_bind_blob(stmt, 3, ip, ip_len, SQLITE_STATIC);
sqlite3_bind_int(stmt, 4, (int)port);
sqlite3_step(stmt); sqlite3_finalize(stmt);
}
}
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); }
/* propagate to others (except sender and the subject node) */
cs_propagate(cs, ch_id, peer, pl, len);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: PEER_UPSERT node=0x%016llx ch=%s",
CS_ID, (unsigned long long)node_id, ch_id);
}
/* ─── PEER_REMOVE (0x0E) ─── */
static void cs_handle_peer_remove(struct chat_sync* cs, uint64_t peer,
const char* ch_id, const uint8_t* pl, size_t len) {
if (len < 8) { DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: PEER_REMOVE too short len=%zu", CS_ID, len); return; }
uint64_t node_id; memcpy(&node_id, pl, 8);
sqlite3* db = cs->inst->topo_sqlite_db;
topo_node_sqlite_member_del(db, ch_id, node_id);
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); }
cs_propagate(cs, ch_id, peer, pl, len);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: PEER_REMOVE node=0x%016llx ch=%s",
CS_ID, (unsigned long long)node_id, ch_id);
}