#include "chat_sync.h" #include "chat_core.h" #include "gui_bridge.h" #include "topo_node_sqlite.h" #include "topo_group.h" #include "conn_mgr.h" #include "topo_group_connect.h" #include "../../../lib/json_flat.h" #include "member_sync.h" #include "merkle_sync.h" #include "../../../src/utun_instance.h" #include "etcp_api.h" #include "etcp.h" #include "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 #include #include static struct chat_sync* g_cs = NULL; 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; void* sync_timer; uint8_t initialized; uint8_t sync_scheduled; 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 #define CS_SYNC_INTERVAL_MS 1000 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; } int rc = etcp_send(conn, entry); DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "%s: etcp_send rc=%d id=0x%02x to=%016llx conn=%s", CS_ID, rc, payload[0], (unsigned long long)dst, conn->log_name); return rc; } /* ── 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 = CONN_MGR_ERR_NOT_FOUND; memcpy(err + 8, &r, 4); memcpy(err + 12, &g_cs->pending_invite_ch_id, 8); gui_bridge_post(GUI_EVT_CONNECT_RESULT, err, 20); 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[20]; memcpy(err, &cs->pending_invite_node_id, 8); int r = CONN_MGR_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; } 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 = CONN_MGR_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) { DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: is_online peer=0x%016llx CONNS=NULL", CS_ID, (unsigned long long)peer_id); 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; int onl = ce->conn->initialized && ce->conn->links_up; DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: is_online peer=0x%016llx found=1 init=%d links=%d online=%d", CS_ID, (unsigned long long)peer_id, ce->conn->initialized, ce->conn->links_up, onl); return onl; } DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: is_online peer=0x%016llx found=0", CS_ID, (unsigned long long)peer_id); 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++) { int p_on = cs_is_peer_online(cs->inst, ch->peer_ids[j]); if (p_on) online++; DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: PEERS_ONLINE peer[%d]=0x%016llx online=%d", CS_ID, j, (unsigned long long)ch->peer_ids[j], p_on); } } 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); DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: PEERS_ONLINE ch=%.*s total_online=%d", CS_ID, (int)cl, ch_id, online); } /* ── Throttled sync (не чаще CS_SYNC_INTERVAL_MS) ── */ static void cs_flush_sync(struct chat_sync* cs) { uint64_t myid = cs->inst->node_id; for (int i = 0; i < cs->channel_count; i++) { struct channel_cache* ch = &cs->channels[i]; for (int j = 0; j < ch->peer_count; j++) { uint64_t pid = ch->peer_ids[j]; if (pid == myid || !cs_is_peer_online(cs->inst, pid)) continue; DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: flush_sync start for peer=%016llx ch=%s", CS_ID, (unsigned long long)pid, ch->channel_id); member_sync_start(cs->inst, pid, ch->channel_id, NULL, NULL); } } } static void cs_sync_timer_cb(void* arg) { struct chat_sync* cs = (struct chat_sync*)arg; if (!cs || !cs->initialized) return; cs->sync_timer = NULL; cs->sync_scheduled = 0; DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: sync timer fire", CS_ID); cs_flush_sync(cs); } static void cs_schedule_sync(struct chat_sync* cs) { if (cs->sync_scheduled) return; cs->sync_scheduled = 1; cs->sync_timer = uasync_set_timeout(cs->inst->ua, CS_SYNC_INTERVAL_MS * 10, cs, cs_sync_timer_cb, "cs_sync"); DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: sync scheduled in %ums", CS_ID, CS_SYNC_INTERVAL_MS); } /* ── Unified peer status change (локальная БД + GUI + push_update + schedule_sync) ── */ static void cs_on_remote_status_changed(uint64_t peer, int online) { if (!g_cs) return; DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: remote status changed peer=%016llx online=%d", CS_ID, (unsigned long long)peer, online); 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) { size_t cl = strlen(g_cs->channels[i].channel_id); if (cl > 63) cl = 63; uint8_t evt[65]; evt[0] = (uint8_t)cl; memcpy(evt + 1, g_cs->channels[i].channel_id, cl); gui_bridge_post(GUI_EVT_MEMBERS_CHANGED, evt, 1 + (int)cl); cs_post_channel_online(g_cs, g_cs->channels[i].channel_id); break; } } } } static void cs_on_peer_status_changed(uint64_t peer, int online) { /* 1. Локальная БД */ member_sync_set_online(g_cs->inst, peer, online); /* 2. GUI + online-бар */ int found_in_channel = 0; 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) { found_in_channel = 1; size_t cl = strlen(g_cs->channels[i].channel_id); if (cl > 63) cl = 63; uint8_t evt[65]; evt[0] = (uint8_t)cl; memcpy(evt + 1, g_cs->channels[i].channel_id, cl); gui_bridge_post(GUI_EVT_MEMBERS_CHANGED, evt, 1 + (int)cl); cs_post_channel_online(g_cs, g_cs->channels[i].channel_id); /* 3. Рассылка дельты всем synced-соседям */ uint8_t st = (uint8_t)online; merkle_sync_push_update(g_cs->inst, g_cs->channels[i].channel_id, peer, 0x01, &st, 1); break; } } } DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: peer_status_changed peer=0x%016llx online=%d channels=%d found=%d", CS_ID, (unsigned long long)peer, online, g_cs->channel_count, found_in_channel); /* 4. Запланировать member_sync_start (throttled) */ cs_schedule_sync(g_cs); } 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); } struct invite_sync_arg { uint64_t node_id; char ch_id[64]; }; static void _on_invite_sync_done(uint64_t peer, const char* ns, int result, void* arg) { struct invite_sync_arg* sa = (struct invite_sync_arg*)arg; DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: invite_sync %s ns=%s peer=%016llx node=0x%016llx", CS_ID, result == MT_OK ? "OK" : "FAIL", ns, (unsigned long long)peer, sa ? (unsigned long long)sa->node_id : 0); if (result == MT_OK && sa && g_cs) { cs_post_channel_online(g_cs, sa->ch_id); uint8_t one = 1; merkle_sync_push_update(g_cs->inst, sa->ch_id, sa->node_id, 0x01, &one, 1); DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: invite_sync push_update node=0x%016llx online=1 ch=%s", CS_ID, (unsigned long long)sa->node_id, sa->ch_id); } u_free(sa); } 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"); DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: conn_up invite path, marking peer=%016llx online", CS_ID, (unsigned long long)peer); cs_on_peer_status_changed(peer, 1); return; } cs_on_peer_status_changed(peer, 1); } 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); cs_on_peer_status_changed(peer, 0); 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); } } } static void cs_on_conn_status(struct ETCP_CONN* conn, int status, void* arg) { switch (status) { case ETCP_CONN_STATUS_UP: cs_on_conn_up(conn, arg); break; case ETCP_CONN_STATUS_DOWN: cs_on_conn_down(conn, arg); break; default: break; } } #if 0 /* replaced by cs_on_conn_status */ 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); } #endif /* ── 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_conn_status_cbk(inst, cs_on_conn_status, NULL); 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); member_sync_set_node_updated_cb(cs_on_remote_status_changed); 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); if (cs->sync_timer) { uasync_cancel_timeout(inst->ua, cs->sync_timer); cs->sync_timer = NULL; } if (cs->sync_scheduled) { cs->sync_scheduled = 0; cs_flush_sync(cs); } member_sync_destroy(inst); cs->initialized = 0; g_cs = NULL; etcp_unbind(inst, ETCP_RT_ID_CHAT_SYNC); etcp_remove_conn_status_cbk(inst, cs_on_conn_status, NULL); if (cs->refresh_timer) { uasync_cancel_timeout(inst->ua, cs->refresh_timer); cs->refresh_timer = NULL; } cs_cancel_proto_timers(cs); 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; (void)node_id; } struct chat_invite { uint64_t channel_id; uint64_t node_id; uint8_t pubkey[32]; uint8_t* addrs_data; int addr_count; }; struct cm_invite_wrap { struct chat_invite inv; }; static void cm_invite_trampoline(void* arg) { struct cm_invite_wrap* w = (struct cm_invite_wrap*)arg; struct chat_invite* inv = &w->inv; struct UTUN_INSTANCE* inst = chat_core_get_inst(); sqlite3* db = chat_core_get_db(); uint64_t node_id = inv->node_id; uint64_t channel_id = inv->channel_id; /* save pubkey to nodes */ sqlite3_stmt* st = NULL; sqlite3_prepare_v2(db, "INSERT INTO nodes(node_id,x25519_pubkey,created_at) VALUES(?,?,?)" " ON CONFLICT(node_id) DO UPDATE SET x25519_pubkey=excluded.x25519_pubkey," " created_at=COALESCE(nodes.created_at, excluded.created_at)", -1, &st, NULL); if (st) { sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); sqlite3_bind_blob(st, 2, inv->pubkey, 32, SQLITE_STATIC); sqlite3_bind_int64(st, 3, 0); sqlite3_step(st); sqlite3_finalize(st); } /* save invite addresses (socket_id=0) */ sqlite3_prepare_v2(db, "DELETE FROM node_addresses WHERE node_id=? AND socket_id=0", -1, &st, NULL); if (st) { sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); sqlite3_step(st); sqlite3_finalize(st); } sqlite3_prepare_v2(db, "INSERT INTO node_addresses(node_id,family,protocol,address,port,addr_type,socket_id) VALUES(?,?,1,?,?,0,?)", -1, &st, NULL); if (st) { const uint8_t* ap = inv->addrs_data; for (int i = 0; i < inv->addr_count; i++) { uint8_t family = *ap++; ap++; /* skip sock_id */ if (family != 4) { ap += 18; continue; } sqlite3_bind_int64(st, 1, (sqlite3_int64)node_id); sqlite3_bind_int(st, 2, 4); sqlite3_bind_blob(st, 3, ap, 4, SQLITE_STATIC); ap += 4; uint16_t port = ((uint16_t)ap[0] << 8) | ap[1]; ap += 2; sqlite3_bind_int(st, 4, (int)port); sqlite3_bind_int(st, 5, 1); /* addr_type=1 (DIRECT) */ sqlite3_step(st); sqlite3_reset(st); } sqlite3_finalize(st); } char ch_str[32]; snprintf(ch_str, sizeof(ch_str), "%llu", (unsigned long long)channel_id); chat_core_ensure_channel_ready(ch_str); /* ensure TOPO_GROUP_TYPE_CHAT group exists */ if (inst->topo_groups) { uint64_t gid = channel_id; if (!topo_groups_find(inst->topo_groups, gid)) topo_groups_create_group(inst->topo_groups, gid, TOPO_GROUP_TYPE_CHAT, ch_str); } /* build TOPO_NODE from invite data */ struct TOPO_NODE* ni = u_calloc(1, sizeof(*ni)); if (!ni) { u_free(w); return; } ni->node_id = node_id; memcpy(ni->public_key, inv->pubkey, 32); const uint8_t* ap = inv->addrs_data; for (int i = 0; i < inv->addr_count; i++) { uint8_t family = *ap++; ap++; /* skip sock_id */ if (family != 4) { ap += 18; continue; } struct TOPO_ADDR4* a4 = u_calloc(1, sizeof(*a4)); if (!a4) break; memcpy(a4->addr, ap, 4); ap += 4; a4->port = ((uint16_t)ap[0] << 8) | ap[1]; ap += 2; a4->type = TOPO_ADDR_NAT; a4->protocol = TOPO_PROTO_UDP; a4->next = ni->v4_addrs; ni->v4_addrs = a4; } if (inst->topo_groups && ni->v4_addrs) { struct TOPO_GROUP* g = topo_groups_find(inst->topo_groups, channel_id); if (g && g->conn_mgr) conn_mgr_connect_from_invite(g->conn_mgr, ni, channel_id, 0, NULL, NULL); } /* free caller-owned ni */ while (ni->v4_addrs) { struct TOPO_ADDR4* n = ni->v4_addrs->next; u_free(ni->v4_addrs); ni->v4_addrs = n; } u_free(ni); u_free(w); } 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 * 8; 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); struct cm_invite_wrap { struct chat_invite inv; }* w = u_malloc(sizeof(struct cm_invite_wrap)); w->inv = *inv; u_free(inv); gui_bridge_post_uasync_fn( (void(*)(void*))cm_invite_trampoline, w); } /* ─── 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]; 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), &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; memcpy(join_msg + mlen, x25519, 32); mlen += 32; memcpy(join_msg + mlen, ed_pub, 32); mlen += 32; 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; { char juser[256]; snprintf(juser, sizeof(juser), "{\"name\":\"%s\"}", cs->inst->name[0] ? cs->inst->name : ""); uint8_t umsg[256]; size_t ulen = 0; memcpy(umsg + ulen, my_join_sig, 64); ulen += 64; my_update_ts = (uint64_t)ntp_time_get_seconds(cs->inst); memcpy(umsg + ulen, &my_update_ts, 8); ulen += 8; size_t jl = strlen(juser); memcpy(umsg + ulen, juser, jl); umsg[ulen + jl] = '\0'; ulen += jl + 1; 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; 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; char juser[256]; snprintf(juser, sizeof(juser), "{\"name\":\"%s\"}", cs->inst->name[0] ? cs->inst->name : ""); uint8_t il = (uint8_t)strlen(juser); buf[boff++] = il; memcpy(buf + boff, juser, 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 + 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; 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_userinfo_len = *p++; char inv_userinfo[128] = ""; if (p + inv_userinfo_len <= pl + len) { memcpy(inv_userinfo, p, inv_userinfo_len); inv_userinfo[inv_userinfo_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, 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; if (inviter_join_sig) memcpy(ivmsg + ilen, inviter_join_sig, 64); else memset(ivmsg + ilen, 0, 64); ilen += 64; memcpy(ivmsg + ilen, &inviter_update_ts, 8); ilen += 8; { const char* nm = inv_userinfo[0] ? inv_userinfo : ""; size_t nl2 = strlen(nm); memcpy(ivmsg + ilen, nm, nl2 + 1); ilen += nl2 + 1; } 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); } sqlite3* vdb = cs->inst->topo_sqlite_db; /* save inviter's address from ETCP connection */ { DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] INFO_RESP save inviter peer=%016llx", CS_ID, (unsigned long long)peer); struct ETCP_LINK* lk = inv_conn->links; int lk_count = 0, lk_written = 0; while (lk) { lk_count++; 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,socket_id) 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_bind_int(as, 4, (int)lk->conn->sock_id); sqlite3_step(as); sqlite3_finalize(as); lk_written++; DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] INFO_RESP INSERT peer=%016llx sock=%d %d.%d.%d.%d:%d", CS_ID, (unsigned long long)peer, (int)lk->conn->sock_id, ip[0], ip[1], ip[2], ip[3], port); } } else if (lk->initialized) { DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] INFO_RESP SKIP non-IPv4 link for peer=%016llx", CS_ID, (unsigned long long)peer); } lk = lk->next; } DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] INFO_RESP DONE peer=%016llx: %d links, %d written", CS_ID, (unsigned long long)peer, lk_count, lk_written); } 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_userinfo, 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); char inv_name[128]; json_flat_get(inv_userinfo, "name", inv_name, sizeof(inv_name)); topo_node_sqlite_node_update_verified(vdb, peer, inv_name, inv_x25519, inv_ed, inviter_join_ts, ntp_time_get_seconds(cs->inst)); } /* 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; memcpy(msg + mlen, x25519, 32); mlen += 32; memcpy(msg + mlen, ed_pub, 32); mlen += 32; memcpy(msg + mlen, &myid, 8); mlen += 8; memcpy(msg + mlen, my_x25519, 32); mlen += 32; 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; char juser2[256]; snprintf(juser2, sizeof(juser2), "{\"name\":\"%s\"}", cs->inst->name[0] ? cs->inst->name : ""); uint8_t nml = (uint8_t)strlen(juser2); jbuf[joff++] = nml; memcpy(jbuf + joff, juser2, 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 + 8 > sizeof(jbuf)) break; jbuf[joff++] = 4; jbuf[joff++] = sock->sock_id; 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 + 20 > sizeof(jbuf)) break; jbuf[joff++] = 6; jbuf[joff++] = sock->sock_id; 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"); } /* ─── 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 userinfo_len = *p++; char joiner_userinfo[128] = ""; if (p + userinfo_len <= pl + len) { memcpy(joiner_userinfo, p, userinfo_len); joiner_userinfo[userinfo_len] = '\0'; p += userinfo_len; } uint8_t addr_cnt = *p++; /* verify join_sig */ uint8_t ch_x25519[32], ch_ed25519[32]; if (topo_node_sqlite_channel_get(cs->inst->topo_sqlite_db, ch_id, NULL, 0, NULL, ch_x25519, ch_ed25519, NULL) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: JOIN channel not in DB ch=%s", CS_ID, ch_id); return; } uint8_t vmsg[256]; size_t vlen = 0; memcpy(vmsg + vlen, ch_x25519, 32); vlen += 32; memcpy(vmsg + vlen, ch_ed25519, 32); vlen += 32; memcpy(vmsg + vlen, &node_id, 8); vlen += 8; memcpy(vmsg + vlen, x25519, 32); vlen += 32; 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 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, (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_userinfo, 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); { char j_name[128]; json_flat_get(joiner_userinfo, "name", j_name, sizeof(j_name)); topo_node_sqlite_node_update_verified(db, node_id, j_name, x25519, ed_pub, join_ts, ntp_time_get_seconds(cs->inst)); } /* save/update node addresses */ DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] JOIN save addrs node=0x%016llx addr_cnt=%d", CS_ID, (unsigned long long)node_id, addr_cnt); for (uint8_t i = 0; i < addr_cnt && p + 2 <= pl + len; i++) { uint8_t fm = *p++; uint8_t sid = *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,socket_id)" " 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_bind_int(stmt, 5, (int)sid); sqlite3_step(stmt); sqlite3_finalize(stmt); } if (fm == 4) DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] JOIN INSERT node=0x%016llx sock=%d %d.%d.%d.%d:%d", CS_ID, (unsigned long long)node_id, sid, ip[0], ip[1], ip[2], ip[3], port); else DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] JOIN INSERT v6 node=0x%016llx sock=%d port=%d", CS_ID, (unsigned long long)node_id, sid, port); } /* 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++] = userinfo_len; memcpy(ubuf + uoff, joiner_userinfo, userinfo_len); uoff += userinfo_len; ubuf[uoff++] = addr_cnt; size_t addr_data_sz = (size_t)(p - (pl + 8 + 32 + 32 + 64 + 8 + 1 + userinfo_len + 1)); if (uoff + addr_data_sz <= sizeof(ubuf)) { memcpy(ubuf + uoff, pl + 8 + 32 + 32 + 64 + 8 + 1 + userinfo_len + 1, addr_data_sz); uoff += addr_data_sz; } cs_propagate(cs, ch_id, peer, ubuf, uoff); } /* add to channel cache — always refresh to include new peer */ cs_refresh_channels(cs); { uint8_t evt[65]; uint8_t cl=(uint8_t)strlen(ch_id); evt[0]=cl; memcpy(evt+1,ch_id,cl); gui_bridge_post(GUI_EVT_MEMBERS_CHANGED, evt, 1+cl); } DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: JOIN calling cs_post_channel_online ch=%s", CS_ID, ch_id); cs_post_channel_online(cs, ch_id); /* initiate member sync with the joiner (server side) */ struct invite_sync_arg* sa = u_malloc(sizeof(struct invite_sync_arg)); if (sa) { sa->node_id = node_id; strncpy(sa->ch_id, ch_id, sizeof(sa->ch_id) - 1); sa->ch_id[sizeof(sa->ch_id) - 1] = '\0'; member_sync_start(cs->inst, node_id, ch_id, _on_invite_sync_done, sa); } else 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 done_cb=%s", CS_ID, (unsigned long long)node_id, ch_id, sa ? "yes" : "no"); 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 + 1 + 64 + 8 + 1 > len) break; /* min: id+x25519+ed+flags+update_sig(64+8)+nl+ac */ 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; uint8_t flags = *p++; const uint8_t* join_sig = NULL; uint64_t join_ts = 0; if (flags & PEERS_FLAG_HAS_JOIN) { if (p + 64 + 8 > pl + len) break; join_sig = p; p += 64; memcpy(&join_ts, p, 8); p += 8; } const uint8_t* update_sig = p; p += 64; uint64_t update_ts; memcpy(&update_ts, p, 8); p += 8; (void)update_sig; (void)update_ts; uint8_t nl = *p++; char peer_userinfo[256] = ""; if (nl && p + nl <= pl + len) { memcpy(peer_userinfo, p, nl); peer_userinfo[nl] = '\0'; p += nl; } uint8_t ac = *p++; /* verify join_sig if present */ if (join_sig) { uint8_t ch_x2[32], ch_ed2[32]; if (topo_node_sqlite_channel_get(db, ch_id, NULL, 0, NULL, ch_x2, ch_ed2, NULL) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: WELCOME channel not in DB ch=%s", CS_ID, ch_id); continue; } uint8_t vmsg[256]; size_t vlen = 0; memcpy(vmsg + vlen, ch_x2, 32); vlen += 32; memcpy(vmsg + vlen, ch_ed2, 32); vlen += 32; memcpy(vmsg + vlen, &node_id, 8); vlen += 8; memcpy(vmsg + vlen, x25519, 32); vlen += 32; 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 — skipping", CS_ID, (unsigned long long)node_id); EVP_MD_CTX_free(ver); EVP_PKEY_free(pkey); continue; } 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_userinfo, 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); { char pn[256]; json_flat_get(peer_userinfo, "name", pn, sizeof(pn)); topo_node_sqlite_node_update_verified(db, node_id, pn, x25519, ed_pub, join_ts, ntp_time_get_seconds(cs->inst)); } DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] WELCOME save addrs node=0x%016llx ac=%d", CS_ID, (unsigned long long)node_id, ac); for (uint8_t j = 0; j < ac; j++) { if (p + 2 > pl + len) break; uint8_t fm = *p++; uint8_t sid = *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,socket_id)" " 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_bind_int(stmt, 5, (int)sid); sqlite3_step(stmt); sqlite3_finalize(stmt); } if (fm == 4) DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] WELCOME INSERT node=0x%016llx sock=%d %d.%d.%d.%d:%d", CS_ID, (unsigned long long)node_id, sid, ip[0], ip[1], ip[2], ip[3], port); else DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] WELCOME INSERT v6 node=0x%016llx sock=%d port=%d", CS_ID, (unsigned long long)node_id, sid, port); } member_sync_cancel(g_cs->inst, node_id, ch_id); } 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 calling cs_post_channel_online ch=%s", CS_ID, ch_id); cs_post_channel_online(cs, ch_id); 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); if (cs->pending_invite_ch_id != 0) { uint64_t ch_id_num = strtoull(ch_id, NULL, 10); uint8_t cevt[20]; memcpy(cevt, &peer, 8); int r = 0; memcpy(cevt + 8, &r, 4); memcpy(cevt + 12, &ch_id_num, 8); gui_bridge_post(GUI_EVT_CONNECT_RESULT, cevt, 20); DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: WELCOME invite success ch=%s peer=%016llx — posted GUI_EVT_CONNECT_RESULT", CS_ID, ch_id, (unsigned long long)peer); cs->pending_invite_node_id = 0; cs->pending_invite_ch_id = 0; if (cs->join_timer) { uasync_cancel_timeout(cs->inst->ua, cs->join_timer); cs->join_timer = NULL; } } } /* ─── 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 userinfo_len = *p++; char peer_userinfo[128] = ""; if (p + userinfo_len <= pl + len) { memcpy(peer_userinfo, p, userinfo_len); peer_userinfo[userinfo_len] = '\0'; p += userinfo_len; } uint8_t ac = *p++; /* verify join_sig */ uint8_t ch_x3[32], ch_ed3[32]; if (topo_node_sqlite_channel_get(cs->inst->topo_sqlite_db, ch_id, NULL, 0, NULL, ch_x3, ch_ed3, NULL) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: PEER_UPSERT channel not in DB ch=%s", CS_ID, ch_id); return; } uint8_t vmsg[256]; size_t vlen = 0; memcpy(vmsg + vlen, ch_x3, 32); vlen += 32; memcpy(vmsg + vlen, ch_ed3, 32); vlen += 32; memcpy(vmsg + vlen, &node_id, 8); vlen += 8; memcpy(vmsg + vlen, x25519, 32); vlen += 32; 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_userinfo, 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); { char pn2[128]; json_flat_get(peer_userinfo, "name", pn2, sizeof(pn2)); topo_node_sqlite_node_update_verified(db, node_id, pn2, x25519, ed_pub, join_ts, ntp_time_get_seconds(cs->inst)); } DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] PEER_UPSERT save addrs node=0x%016llx ac=%d", CS_ID, (unsigned long long)node_id, ac); for (uint8_t i = 0; i < ac; i++) { if (p + 2 > pl + len) break; uint8_t fm = *p++; uint8_t sid = *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,socket_id)" " 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_bind_int(stmt, 5, (int)sid); sqlite3_step(stmt); sqlite3_finalize(stmt); } if (fm == 4) DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] PEER_UPSERT INSERT node=0x%016llx sock=%d %d.%d.%d.%d:%d", CS_ID, (unsigned long long)node_id, sid, ip[0], ip[1], ip[2], ip[3], port); else DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: [ADDR_SYNC] PEER_UPSERT INSERT v6 node=0x%016llx sock=%d port=%d", CS_ID, (unsigned long long)node_id, sid, port); } 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); }