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.
1090 lines
44 KiB
1090 lines
44 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_router.h" |
|
#include "../../../src/etcp_api.h" |
|
#include "../../../src/etcp.h" |
|
#include "../../../src/conn_mgr.h" |
|
#include "../../../src/topo_group.h" |
|
#include "../../../src/secure_channel.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 <openssl/evp.h> |
|
|
|
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* ttl_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; } |
|
} |
|
|
|
/* ── 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_CONNECTIVITY, "%s: SEND %s to=%016llx ch=%s len=%zu", |
|
CS_ID, cs_msg_name(payload[0]), (unsigned long long)dst, ch_id, len); |
|
return etcp_route_send(cs->inst, dst, entry, 0); |
|
} |
|
|
|
/* ── 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; |
|
} |
|
|
|
/* ── INIT_SYNC handler ── */ |
|
|
|
static void cs_handle_init_sync(struct chat_sync* cs, uint64_t peer, |
|
const char* ch_id, const uint8_t* pl, size_t len) { |
|
if (len < 36) return; |
|
uint32_t peer_count; memcpy(&peer_count, pl, 4); |
|
uint8_t peer_last_ch[32]; memcpy(peer_last_ch, pl + 4, 32); |
|
|
|
uint32_t my_count = chat_core_count(ch_id); |
|
uint32_t tp = peer_count < my_count ? peer_count : my_count; |
|
if (tp > 0) tp--; |
|
uint8_t my_ch[32]; |
|
chat_core_chain_hash_at(ch_id, tp, my_ch); |
|
|
|
if (memcmp(my_ch, peer_last_ch, 32) == 0) { |
|
/* synced */ |
|
uint8_t resp[64]; resp[0] = CS_MSG_INIT_RESP; |
|
memcpy(resp + 1, &my_count, 4); resp[5] = 0; |
|
cs_send(cs, ch_id, peer, resp, 6); |
|
struct channel_cache* ch = cs_find(cs, ch_id); |
|
if (ch) ch->synced = CS_SYNC_DONE; |
|
return; |
|
} |
|
|
|
/* divergence — send INIT_RESP with 0 sparse for now, peer will handle */ |
|
uint8_t resp[64]; resp[0] = CS_MSG_INIT_RESP; |
|
memcpy(resp + 1, &my_count, 4); |
|
resp[5] = 0; /* sparse_count=0 */ |
|
cs_send(cs, ch_id, peer, resp, 6); |
|
} |
|
|
|
/* ── INIT_RESP handler ── */ |
|
|
|
static void cs_handle_init_resp(struct chat_sync* cs, uint64_t peer, |
|
const char* ch_id, const uint8_t* pl, size_t len) { |
|
if (len < 37) return; |
|
uint32_t peer_count = *(const uint32_t*)pl; |
|
uint32_t test_pos = *(const uint32_t*)(pl + 4); |
|
uint8_t peer_ch[32]; memcpy(peer_ch, pl + 8, 32); |
|
uint8_t sparse_count = pl[40]; (void)sparse_count; |
|
|
|
uint8_t my_ch[32]; |
|
chat_core_chain_hash_at(ch_id, test_pos, my_ch); |
|
|
|
if (memcmp(my_ch, peer_ch, 32) == 0) { |
|
struct channel_cache* ch = cs_find(cs, ch_id); |
|
if (ch && peer_count > test_pos + 1) { |
|
uint32_t from = test_pos + 1; |
|
uint8_t snd[7]; snd[0] = CS_MSG_SEND_DATA; |
|
memcpy(snd + 1, &from, 4); uint16_t z = 0; memcpy(snd + 5, &z, 2); |
|
cs_send(cs, ch_id, peer, snd, 7); |
|
} else { |
|
if (ch) ch->synced = CS_SYNC_DONE; |
|
} |
|
return; |
|
} |
|
|
|
/* mismatch — request from start */ |
|
uint32_t from = 0; |
|
uint8_t snd[7]; snd[0] = CS_MSG_SEND_DATA; |
|
memcpy(snd + 1, &from, 4); uint16_t z = 0; memcpy(snd + 5, &z, 2); |
|
cs_send(cs, ch_id, peer, snd, 7); |
|
} |
|
|
|
/* ── SEND_DATA handler ── */ |
|
|
|
static void cs_handle_send_data(struct chat_sync* cs, uint64_t peer, |
|
const char* ch_id, const uint8_t* pl, size_t len) { |
|
if (len < 6) return; |
|
uint32_t from = *(const uint32_t*)pl; |
|
uint16_t count = *(const uint16_t*)(pl + 4); |
|
|
|
if (count == 0) { |
|
/* peer requests OUR data starting from 'from' */ |
|
uint32_t cur = chat_core_cursor_open(ch_id); |
|
if (!cur) return; |
|
|
|
uint16_t sent = 0; |
|
uint8_t rec_buf[4096]; size_t rec_len; |
|
|
|
while (chat_core_cursor_next(cur, rec_buf, sizeof(rec_buf), &rec_len) == 0 |
|
&& rec_len > 0) { |
|
uint8_t snd[8192]; uint32_t off = 0; |
|
snd[off++] = CS_MSG_SEND_DATA; |
|
uint32_t pos = from + sent; |
|
memcpy(snd + off, &pos, 4); off += 4; |
|
uint16_t c = 1; memcpy(snd + off, &c, 2); off += 2; |
|
if (off + rec_len <= (uint32_t)sizeof(snd)) |
|
{ memcpy(snd + off, rec_buf, rec_len); off += (uint32_t)rec_len; } |
|
cs_send(cs, ch_id, peer, snd, off); |
|
sent++; |
|
} |
|
chat_core_cursor_close(cur); |
|
return; |
|
} |
|
|
|
/* peer sent US data — insert each record */ |
|
const uint8_t* ptr = pl + 6; |
|
size_t remain = len - 6; |
|
|
|
for (uint16_t i = 0; i < count && remain > 0; i++) { |
|
if (remain < 29) break; |
|
const uint8_t* rec_start = ptr; |
|
ptr += 8 + 8 + 8; remain -= 24; /* ts, dh, nid */ |
|
if (remain < 1) break; |
|
uint8_t ct_len = *ptr; ptr++; remain--; |
|
if (remain < ct_len) break; |
|
ptr += ct_len; remain -= ct_len; |
|
if (remain < 4) break; |
|
uint32_t dlen; memcpy(&dlen, ptr, 4); ptr += 4; remain -= 4; |
|
if (remain < dlen) break; |
|
ptr += dlen; remain -= dlen; |
|
|
|
size_t reclen = (size_t)(ptr - rec_start); |
|
chat_core_insert_record(ch_id, rec_start, reclen); |
|
} |
|
|
|
uint32_t next = from + count; |
|
uint8_t snd[7]; snd[0] = CS_MSG_SEND_DATA; |
|
memcpy(snd + 1, &next, 4); uint16_t z = 0; memcpy(snd + 5, &z, 2); |
|
cs_send(cs, ch_id, peer, snd, 7); |
|
|
|
struct channel_cache* ch = cs_find(cs, ch_id); |
|
if (ch) ch->synced = CS_SYNC_DONE; |
|
} |
|
|
|
/* ── PUSH handler ── */ |
|
|
|
static void cs_handle_push(struct chat_sync* cs, uint64_t peer, |
|
const char* ch_id, const uint8_t* pl, size_t len) { |
|
if (len < 24 + 1 + 4) return; |
|
uint64_t ts, dh; |
|
memcpy(&ts, pl, 8); memcpy(&dh, pl + 8, 8); |
|
|
|
const uint8_t* rec = pl + 24; |
|
int reclen = (int)(len - 24); |
|
int rc = chat_core_insert_record(ch_id, rec, (size_t)reclen); |
|
|
|
if (rc == 0) { |
|
uint8_t ack[17]; ack[0] = CS_MSG_ACK_PUSH; |
|
memcpy(ack + 1, &dh, 8); memcpy(ack + 9, &ts, 8); |
|
cs_send(cs, ch_id, peer, ack, 17); |
|
|
|
uint8_t evt[65]; evt[0] = (uint8_t)strlen(ch_id); |
|
memcpy(evt + 1, ch_id, evt[0]); |
|
gui_bridge_post(GUI_EVT_MSG_RECEIVED, evt, 1 + evt[0]); |
|
|
|
struct channel_cache* ch = cs_find(cs, ch_id); |
|
if (ch) ch->msg_count++; |
|
} |
|
} |
|
|
|
/* ── ACK_PUSH handler ── */ |
|
|
|
static void cs_handle_ack_push(struct chat_sync* cs, uint64_t peer, |
|
const char* ch_id, const uint8_t* pl, size_t len) { |
|
(void)peer; |
|
if (len < 16) return; |
|
uint64_t ts, dh; memcpy(&dh, pl, 8); memcpy(&ts, pl + 8, 8); |
|
chat_core_mark_sent(ch_id, ts, dh, cs->inst->node_id); |
|
} |
|
|
|
/* ── SYNC_DONE handler ── */ |
|
|
|
static void cs_handle_sync_done(struct chat_sync* cs, uint64_t peer, |
|
const char* ch_id, const uint8_t* pl, size_t len) { |
|
(void)peer; |
|
if (len < 4) return; |
|
uint32_t pc; memcpy(&pc, pl, 4); |
|
struct channel_cache* ch = cs_find(cs, ch_id); |
|
if (ch) { ch->msg_count = pc; ch->synced = CS_SYNC_DONE; } |
|
} |
|
|
|
/* ── 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]; |
|
const uint8_t* pl = d + 3 + ch_len; |
|
size_t plen = dlen - 3 - ch_len; |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%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_INIT_SYNC: cs_handle_init_sync(g_cs, peer, ch_id, pl, plen); break; |
|
case CS_MSG_INIT_RESP: cs_handle_init_resp(g_cs, peer, ch_id, pl, plen); break; |
|
case CS_MSG_SEND_DATA: cs_handle_send_data(g_cs, peer, ch_id, pl, plen); break; |
|
case CS_MSG_PUSH: cs_handle_push(g_cs, peer, ch_id, pl, plen); break; |
|
case CS_MSG_ACK_PUSH: cs_handle_ack_push(g_cs, peer, ch_id, pl, plen); break; |
|
case CS_MSG_SYNC_DONE: cs_handle_sync_done(g_cs, peer, ch_id, pl, plen); break; |
|
case CS_MSG_CHANNEL_INFO_REQ: cs_handle_channel_info_req(g_cs, peer, ch_id, pl, plen); break; |
|
case CS_MSG_CHANNEL_INFO_RESP:cs_handle_channel_info_resp(g_cs, peer, ch_id, pl, plen); break; |
|
case CS_MSG_CHANNEL_JOIN: cs_handle_channel_join(g_cs, peer, ch_id, pl, plen); break; |
|
case CS_MSG_WELCOME: cs_handle_welcome(g_cs, peer, ch_id, pl, plen); break; |
|
case CS_MSG_PEER_UPSERT: cs_handle_peer_upsert(g_cs, peer, ch_id, pl, plen); break; |
|
case CS_MSG_ERROR: { |
|
DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%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_CONNECTIVITY, "%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_CONNECTIVITY, "%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 = CONN_MGR_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_CONNECTIVITY, "%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; |
|
} |
|
|
|
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_CONNECTIVITY, "%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; |
|
|
|
if (g_cs->pending_invite_ch_id != 0 && !g_cs->info_req_timer) { |
|
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; |
|
} |
|
|
|
for (int i = 0; i < g_cs->channel_count; i++) { |
|
struct channel_cache* ch = &g_cs->channels[i]; |
|
int found = 0; |
|
for (int j = 0; j < ch->peer_count; j++) |
|
if (ch->peer_ids[j] == peer) { found = 1; break; } |
|
if (!found) continue; |
|
ch->synced = CS_SYNC_IN_PROGRESS; |
|
uint8_t msg[37]; msg[0] = CS_MSG_INIT_SYNC; |
|
chat_core_chain_hash_at(ch->channel_id, |
|
ch->msg_count > 0 ? ch->msg_count - 1 : 0, ch->last_chain_hash); |
|
memcpy(msg + 1, &ch->msg_count, 4); |
|
memcpy(msg + 5, ch->last_chain_hash, 32); |
|
cs_send(g_cs, ch->channel_id, peer, msg, 37); |
|
member_sync_start(g_cs->inst, peer, ch->channel_id, _on_member_sync_done, ch); |
|
} |
|
member_sync_set_online(g_cs->inst, 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; |
|
cs_cancel_proto_timers(g_cs); |
|
if (g_cs->pending_invite_node_id == peer) { |
|
DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%s: conn down while waiting invite resp peer=%016llx", CS_ID, (unsigned long long)peer); |
|
g_cs->pending_invite_node_id = 0; |
|
g_cs->pending_invite_ch_id = 0; |
|
} |
|
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_new_conn(struct ETCP_CONN* conn, void* arg) { |
|
(void)arg; |
|
if (!conn) return; |
|
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 ── */ |
|
|
|
static void ttl_timer_cb(void* arg) { |
|
(void)arg; |
|
if (!g_cs || !g_cs->initialized) return; |
|
uint64_t cutoff = get_time_us() - 86400000000ULL; |
|
for (int i = 0; i < g_cs->channel_count; i++) |
|
chat_core_ttl_delete(g_cs->channels[i].channel_id, |
|
g_cs->inst->node_id, cutoff); |
|
g_cs->ttl_timer = uasync_set_timeout(g_cs->inst->ua, |
|
3600u * 10000u, g_cs, ttl_timer_cb, "cs_ttl"); |
|
} |
|
|
|
/* ── 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; |
|
struct chat_sync* cs = u_calloc(1, sizeof(*cs)); |
|
if (!cs) return -1; |
|
cs->inst = inst; |
|
cs->initialized = 1; |
|
g_cs = cs; |
|
|
|
etcp_router_bind(inst, ETCP_RT_ID_CHAT_SYNC, chat_sync_recv_cb); |
|
etcp_add_new_conn_cbk(inst, cs_on_new_conn, NULL); |
|
|
|
struct ETCP_CONN* c = inst->connections; |
|
while (c) { |
|
etcp_conn_add_up_cbk(c, cs_on_conn_up, NULL); |
|
etcp_conn_add_down_cbk(c, cs_on_conn_down, NULL); |
|
c = c->next; |
|
} |
|
|
|
cs_refresh_channels(cs); |
|
|
|
cs->refresh_timer = uasync_set_timeout(inst->ua, 50000u, cs, |
|
refresh_timer_cb, "cs_refresh"); |
|
cs->ttl_timer = uasync_set_timeout(inst->ua, 3600u * 10000u, cs, |
|
ttl_timer_cb, "cs_ttl"); |
|
cs->info_req_timer = NULL; |
|
cs->join_timer = NULL; |
|
|
|
member_sync_init(inst); |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%s: initialized", CS_ID); |
|
return 0; |
|
} |
|
|
|
void chat_sync_destroy(struct UTUN_INSTANCE* inst) { |
|
struct chat_sync* cs = g_cs; |
|
if (!cs || !inst) return; |
|
member_sync_destroy(inst); |
|
cs->initialized = 0; g_cs = NULL; |
|
|
|
etcp_router_unbind(inst, ETCP_RT_ID_CHAT_SYNC); |
|
if (cs->refresh_timer) { uasync_cancel_timeout(inst->ua, cs->refresh_timer); cs->refresh_timer = NULL; } |
|
if (cs->ttl_timer) { uasync_cancel_timeout(inst->ua, cs->ttl_timer); cs->ttl_timer = NULL; } |
|
cs_cancel_proto_timers(cs); |
|
|
|
struct ETCP_CONN* c = inst->connections; |
|
while (c) { |
|
etcp_conn_remove_up_cbk(c, cs_on_conn_up, NULL); |
|
etcp_conn_remove_down_cbk(c, cs_on_conn_down, NULL); |
|
c = c->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); |
|
} |
|
|
|
int chat_sync_push(struct UTUN_INSTANCE* inst, |
|
const char* ch_id, uint64_t node_id, |
|
const char* content_type, const uint8_t* data, uint32_t data_len, |
|
uint64_t timestamp, uint64_t datahash) { |
|
if (!g_cs || !g_cs->initialized) return -1; |
|
struct channel_cache* ch = cs_find(g_cs, ch_id); |
|
|
|
uint8_t ct_len = content_type ? (uint8_t)strlen(content_type) : 0; |
|
if (ct_len > 63) ct_len = 63; |
|
size_t total = 1 + 8 + 8 + 8 + 1 + ct_len + 4 + data_len + 32; |
|
uint8_t* buf = u_malloc(total); |
|
if (!buf) return -1; |
|
|
|
uint8_t* p = buf; |
|
*p++ = CS_MSG_PUSH; |
|
memcpy(p, ×tamp, 8); p += 8; |
|
memcpy(p, &datahash, 8); p += 8; |
|
memcpy(p, &node_id, 8); p += 8; |
|
*p++ = ct_len; |
|
if (ct_len) { memcpy(p, content_type, ct_len); p += ct_len; } |
|
uint32_t dlen = data_len; |
|
memcpy(p, &dlen, 4); p += 4; |
|
if (data_len) { memcpy(p, data, data_len); p += data_len; } |
|
memset(p, 0, 32); |
|
|
|
int sent = 0; |
|
uint64_t myid = g_cs->inst->node_id; |
|
if (ch) { |
|
for (int i = 0; i < ch->peer_count; i++) { |
|
if (ch->peer_ids[i] == node_id || ch->peer_ids[i] == myid) continue; |
|
cs_send(g_cs, ch_id, ch->peer_ids[i], buf, (size_t)(p + 32 - buf)); |
|
sent++; |
|
} |
|
} |
|
u_free(buf); |
|
return sent; |
|
} |
|
|
|
void chat_sync_connect_node(struct UTUN_INSTANCE* inst, uint64_t node_id) { |
|
uint8_t buf[2048]; size_t buf_len; |
|
if (chat_core_load_nodeinfo(node_id, buf, sizeof(buf), &buf_len) == 0 && buf_len > 0) { |
|
/* nodeinfo loaded — conn_mgr will pick it up */ |
|
(void)inst; |
|
conn_mgr_connect_node(inst->conn_mgr, node_id, 30000, 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; |
|
|
|
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_CONNECTIVITY, "%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_CONNECTIVITY, "%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 (!etcp_router_conn_get(cs->inst, pid, ETCP_RT_ID_CHAT_SYNC)) 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; |
|
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_groups->topo_sqlite_db, |
|
ch_id, name, (int)sizeof(name), &is_dm, &owner, x25519, ed_pub, ch_sig) != 0) { |
|
DEBUG_WARN(DEBUG_CATEGORY_CONNECTIVITY, "%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; |
|
uint8_t my_join_sig[64] = {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; |
|
{ const char* nm = cs->inst->name[0] ? cs->inst->name : ""; |
|
size_t nl = strlen(nm); memcpy(join_msg + mlen, nm, nl); mlen += nl; join_msg[mlen++] = '\0'; } |
|
cs_ed25519_sign(cs->inst->my_ed25519_privkey, join_msg, mlen, my_join_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; |
|
memcpy(buf + boff, my_join_sig, 64); boff += 64; |
|
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) { |
|
if (cs->info_req_timer) { uasync_cancel_timeout(cs->inst->ua, cs->info_req_timer); cs->info_req_timer = NULL; } |
|
if (len < 1) return; |
|
uint8_t nl = pl[0]; if (1 + nl + 8 + 1 + 32 + 32 + 64 + 64 > 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; |
|
const uint8_t* inviter_join_sig = p; |
|
|
|
/* 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_CONNECTIVITY, "%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_groups->topo_sqlite_db, |
|
ch_id, name, (int)is_dm, owner, x25519, NULL, ed_pub, NULL, ch_sig); |
|
|
|
/* verify inviter's join_sig */ |
|
{ sqlite3* vdb = cs->inst->topo_groups->topo_sqlite_db; |
|
/* read inviter's Ed25519 pubkey */ |
|
uint8_t inv_ed[32] = {0}; |
|
sqlite3_stmt* es = NULL; |
|
sqlite3_prepare_v2(vdb, "SELECT ed25519_pubkey FROM nodes WHERE node_id=?", -1, &es, NULL); |
|
if (es) { sqlite3_bind_int64(es, 1, (sqlite3_int64)peer); |
|
if (sqlite3_step(es) == SQLITE_ROW) memcpy(inv_ed, sqlite3_column_blob(es, 0), 32); |
|
sqlite3_finalize(es); } |
|
/* verify */ |
|
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, x25519, 32); ilen += 32; |
|
{ char nnm[64] = ""; _get_node_name(vdb, peer, nnm, sizeof(nnm)); |
|
size_t nl = strlen(nnm); memcpy(ivmsg + ilen, nnm, nl); ilen += nl; ivmsg[ilen++] = '\0'; } |
|
if (cs_ed25519_verify(inv_ed, ivmsg, ilen, inviter_join_sig) != 0) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: CHANNEL_INFO_RESP invalid inviter_join_sig peer=%016llx", |
|
CS_ID, (unsigned long long)peer); |
|
} |
|
} |
|
|
|
/* save inviter as node and member */ |
|
topo_node_sqlite_member_put(cs->inst->topo_groups->topo_sqlite_db, ch_id, peer, |
|
inviter_join_sig, inviter_join_sig); |
|
|
|
/* 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]; |
|
{ |
|
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'; } |
|
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; |
|
/* 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 + 1) 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; |
|
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; |
|
{ sqlite3* tdb = cs->inst->topo_groups->topo_sqlite_db; char nnm[64] = ""; _get_node_name(tdb, node_id, nnm, sizeof(nnm)); |
|
size_t nl = strlen(nnm); memcpy(vmsg + vlen, nnm, nl); vlen += nl; vmsg[vlen++] = '\0'; } |
|
if (cs_ed25519_verify(ed_pub, vmsg, vlen, join_sig) != 0) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: JOIN invalid sig node=0x%016llx ch=%s", CS_ID, |
|
(unsigned long long)node_id, ch_id); |
|
return; |
|
} |
|
|
|
sqlite3* db = cs->inst->topo_groups->topo_sqlite_db; |
|
topo_node_sqlite_member_put(db, ch_id, node_id, join_sig, NULL); |
|
|
|
/* 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,is_nat)" |
|
" 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; |
|
ubuf[uoff++] = addr_cnt; |
|
size_t addr_data_sz = (size_t)(p - (pl + 8 + 32 + 32 + 64 + 1)); |
|
if (uoff + addr_data_sz <= sizeof(ubuf)) { |
|
memcpy(ubuf + uoff, pl + 8 + 32 + 32 + 64 + 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); |
|
} |
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_CONNECTIVITY, "%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) return; |
|
sqlite3* db = cs->inst->topo_groups->topo_sqlite_db; |
|
|
|
const uint8_t* p = pl; |
|
uint16_t pc; memcpy(&pc, p, 2); p += 2; |
|
for (uint16_t i = 0; i < pc; i++) { |
|
if ((size_t)(p - pl) + 8 + 32 + 32 + 64 + 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; |
|
uint8_t ac = *p++; |
|
|
|
topo_node_sqlite_member_put(db, ch_id, node_id, join_sig, NULL); |
|
|
|
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,is_nat)" |
|
" 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); |
|
|
|
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_CONNECTIVITY, "%s: WELCOME processed ch=%s peers=%d", |
|
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 + 1) 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; |
|
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; |
|
{ sqlite3* tdb = cs->inst->topo_groups->topo_sqlite_db; char nnm[64] = ""; _get_node_name(tdb, node_id, nnm, sizeof(nnm)); |
|
size_t nl = strlen(nnm); memcpy(vmsg + vlen, nnm, nl); vlen += nl; vmsg[vlen++] = '\0'; } |
|
if (cs_ed25519_verify(ed_pub, vmsg, vlen, join_sig) != 0) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: PEER_UPSERT invalid sig node=0x%016llx", CS_ID, |
|
(unsigned long long)node_id); |
|
return; |
|
} |
|
|
|
sqlite3* db = cs->inst->topo_groups->topo_sqlite_db; |
|
|
|
topo_node_sqlite_member_put(db, ch_id, node_id, join_sig, NULL); |
|
|
|
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,is_nat)" |
|
" 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_CONNECTIVITY, "%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) return; |
|
uint64_t node_id; memcpy(&node_id, pl, 8); |
|
|
|
sqlite3* db = cs->inst->topo_groups->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_CONNECTIVITY, "%s: PEER_REMOVE node=0x%016llx ch=%s", |
|
CS_ID, (unsigned long long)node_id, ch_id); |
|
}
|
|
|