#include "chat_sync.h" #include "chat_core.h" #include "gui_bridge.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 "../../../lib/u_async.h" #include "../../../lib/ll_queue.h" #include "../../../lib/debug_config.h" #include "../../../lib/mem.h" #include "../../../lib/platform_compat.h" #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* ttl_timer; uint8_t initialized; }; #define CS_ID "chat_sync" /* ── Send ── */ static int cs_send(struct chat_sync* cs, const char* ch_id, uint64_t dst, const uint8_t* payload, size_t len) { 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; 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; } } /* ── 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; 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; default: break; } u_free(entry->dgram); queue_entry_free(entry); } /* ── Connection callbacks ── */ 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; 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); } } 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; 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; } } 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"); DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: initialized", CS_ID); return 0; } void chat_sync_destroy(struct UTUN_INSTANCE* inst) { struct chat_sync* cs = g_cs; if (!cs || !inst) return; 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; } 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 node_id, const uint8_t* pubkey_bin, const uint8_t* addrs_data, int addr_count, const uint8_t* channel_id, int ch_id_len) { 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->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; inv->ch_id = u_malloc((size_t)ch_id_len + 1); if (!inv->ch_id) { u_free(inv->addrs_data); u_free(inv); return; } memcpy(inv->ch_id, channel_id, (size_t)ch_id_len); inv->ch_id_len = ch_id_len; gui_bridge_post_uasync_fn( (void(*)(void*))chat_core_connect_from_invite, inv); }