|
|
|
|
@ -36,17 +36,45 @@ struct chat_sync {
|
|
|
|
|
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_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; |
|
|
|
|
@ -60,6 +88,8 @@ static int cs_send(struct chat_sync* cs, const char* ch_id, uint64_t dst,
|
|
|
|
|
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_DEBUG, "%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); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
@ -282,6 +312,9 @@ static void chat_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
|
|
|
|
|
const uint8_t* pl = d + 3 + ch_len; |
|
|
|
|
size_t plen = dlen - 3 - ch_len; |
|
|
|
|
|
|
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%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; |
|
|
|
|
@ -295,27 +328,56 @@ static void chat_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
|
|
|
|
|
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_PEER_REMOVE: cs_handle_peer_remove(g_cs, peer, ch_id, pl, plen); break; |
|
|
|
|
default: break; |
|
|
|
|
default: DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%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_DEBUG, "%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_DEBUG, "%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 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; |
|
|
|
|
|
|
|
|
|
/* pending invite: send CHANNEL_INFO_REQ */ |
|
|
|
|
if (g_cs->pending_invite_node_id == peer && g_cs->pending_invite_ch_id != 0) { |
|
|
|
|
char ch_id_str[64]; |
|
|
|
|
snprintf(ch_id_str, sizeof(ch_id_str), "%llu", (unsigned long long)g_cs->pending_invite_ch_id); |
|
|
|
|
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->pending_invite_node_id = 0; |
|
|
|
|
g_cs->pending_invite_ch_id = 0; |
|
|
|
|
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; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
@ -339,6 +401,12 @@ 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_DEBUG, "%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++) |
|
|
|
|
@ -451,6 +519,8 @@ int chat_sync_init(struct UTUN_INSTANCE* inst,
|
|
|
|
|
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; |
|
|
|
|
|
|
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: initialized", CS_ID); |
|
|
|
|
return 0; |
|
|
|
|
@ -464,6 +534,7 @@ void chat_sync_destroy(struct UTUN_INSTANCE* inst) {
|
|
|
|
|
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) { |
|
|
|
|
@ -634,6 +705,7 @@ static void cs_handle_channel_info_req(struct chat_sync* cs, uint64_t peer,
|
|
|
|
|
|
|
|
|
|
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; |
|
|
|
|
@ -717,6 +789,10 @@ static void cs_handle_channel_info_resp(struct chat_sync* cs, uint64_t peer,
|
|
|
|
|
} |
|
|
|
|
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 ─── */ |
|
|
|
|
@ -812,6 +888,7 @@ static void cs_handle_channel_join(struct chat_sync* cs, uint64_t peer,
|
|
|
|
|
|
|
|
|
|
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; |
|
|
|
|
|