Browse Source

refactor: move chat_sync and member_sync to direct ETCP P2P, bypass etcp_router

Replace etcp_router_bind/etcp_route_send/etcp_router_unbind with
direct etcp_bind/etcp_send/etcp_unbind + find-conn-by-node_id helpers.
Same approach as db_sync (0x20). Service IDs unchanged (0x30, 0x31).
topo_upd
Evgeny 3 months ago
parent
commit
acaa3a767c
  1. 1
      tools/chatgui/transport/chat_core.c
  2. 25
      tools/chatgui/transport/chat_sync.c
  3. 2
      tools/chatgui/transport/chat_sync.h
  4. 18
      tools/chatgui/transport/merkle_sync.c
  5. 2
      tools/chatgui/transport/merkle_sync.h

1
tools/chatgui/transport/chat_core.c

@ -13,7 +13,6 @@
#include "member_sync.h"
#include "../../../src/utun_instance.h"
#include "../../../src/etcp_router.h"
#include "../../../src/etcp_api.h"
#include "../../../src/topo_group.h"
#include "../../../src/topo_node.h"

25
tools/chatgui/transport/chat_sync.c

@ -5,7 +5,6 @@
#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"
@ -304,6 +303,9 @@ static void cs_cancel_proto_timers(struct chat_sync* cs) {
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,
@ -324,7 +326,9 @@ static int cs_send(struct chat_sync* cs, const char* ch_id, uint64_t dst,
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);
struct ETCP_CONN* conn = cs_find_conn_for_node(cs->inst, dst);
if (!conn) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: no conn for node %016llx", CS_ID, (unsigned long long)dst); u_free(buf); queue_entry_free(entry); return -1; }
return etcp_send(conn, entry);
}
/* ── Channel cache ── */
@ -650,6 +654,17 @@ static void cs_join_timeout_cb(void* arg) {
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) {
struct ETCP_CONN* c = inst->connections;
while (c) {
if (c->peer_node_id == node_id && c->initialized && c->links_up) return c;
c = c->next;
}
return NULL;
}
/* ── helper: check if peer has active ETCP link ── */
static int cs_is_peer_online(struct UTUN_INSTANCE* inst, uint64_t peer_id) {
@ -851,7 +866,7 @@ int chat_sync_init(struct UTUN_INSTANCE* inst,
cs->initialized = 1;
g_cs = cs;
etcp_router_bind(inst, ETCP_RT_ID_CHAT_SYNC, chat_sync_recv_cb);
etcp_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;
@ -883,7 +898,7 @@ void chat_sync_destroy(struct UTUN_INSTANCE* inst) {
member_sync_destroy(inst);
cs->initialized = 0; g_cs = NULL;
etcp_router_unbind(inst, ETCP_RT_ID_CHAT_SYNC);
etcp_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);
@ -1027,7 +1042,7 @@ static void cs_propagate(struct chat_sync* cs, const char* ch_id, uint64_t exclu
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;
if (!cs_find_conn_for_node(cs->inst, pid)) continue;
cs_send(cs, ch_id, pid, payload, len);
}
}

2
tools/chatgui/transport/chat_sync.h

@ -11,7 +11,7 @@ extern "C" {
struct UTUN_INSTANCE;
struct UASYNC;
/* etcp_router service ID */
/* etcp service ID (direct P2P) */
#define ETCP_RT_ID_CHAT_SYNC 0x30
#define ETCP_RT_ID_MEMBER_SYNC 0x31

18
tools/chatgui/transport/merkle_sync.c

@ -1,7 +1,6 @@
#include "merkle_sync.h"
#include "../../../src/utun_instance.h"
#include "../../../src/etcp_router.h"
#include "../../../src/etcp_api.h"
#include "../../../src/etcp.h"
#include "../../../src/topo_group.h"
@ -174,6 +173,15 @@ static int _get_level_hashes(struct merkle_sync* ms, const char* ns,
/* ── Send helpers ── */
static struct ETCP_CONN* ms_find_conn_for_node(struct UTUN_INSTANCE* inst, uint64_t node_id) {
struct ETCP_CONN* c = inst->connections;
while (c) {
if (c->peer_node_id == node_id && c->initialized && c->links_up) return c;
c = c->next;
}
return NULL;
}
static int _send_msg(struct merkle_sync* ms, uint64_t peer,
const uint8_t* payload, size_t len) {
if (!ms->inst || len < 1) return -1;
@ -184,7 +192,9 @@ static int _send_msg(struct merkle_sync* ms, uint64_t peer,
struct ll_entry* entry = queue_entry_new(0);
if (!entry) { u_free(buf); return -1; }
entry->dgram = buf; entry->len = 1 + len;
int r = etcp_route_send(ms->inst, peer, entry, 0);
struct ETCP_CONN* conn = ms_find_conn_for_node(ms->inst, peer);
if (!conn) { DEBUG_ERROR(DEBUG_CATEGORY_CONNECTIVITY, "%s: no conn for node %016llx", MS_ID, (unsigned long long)peer); u_free(buf); queue_entry_free(entry); return -1; }
int r = etcp_send(conn, entry);
if (r != 0) { u_free(buf); queue_entry_free(entry); }
return r;
}
@ -555,7 +565,7 @@ int merkle_sync_init(struct UTUN_INSTANCE* inst, uint8_t svc_id,
_ensure_table(ms);
etcp_router_bind(inst, svc_id, _recv_cb);
etcp_bind(inst, svc_id, _recv_cb);
ms->bg_timer = uasync_set_timeout(inst->ua, (uint32_t)(MS_BG_INTERVAL_MS * 10),
ms, _bg_timer_cb, "ms_bg");
@ -569,7 +579,7 @@ void merkle_sync_destroy(struct UTUN_INSTANCE* inst) {
if (!ms || !inst) return;
ms->initialized = 0; g_merkle = NULL;
etcp_router_bind(inst, ms->svc_id, NULL);
etcp_unbind(inst, ms->svc_id);
if (ms->bg_timer) { uasync_cancel_timeout(inst->ua, ms->bg_timer); ms->bg_timer = NULL; }

2
tools/chatgui/transport/merkle_sync.h

@ -165,7 +165,7 @@ typedef void (*merkle_sync_done_cb)(uint64_t peer, const char* ns, int result, v
* ETCP-сервис svc_id (напр. 0x31), запускает фоновую проверку (bg_timer).
* Вызывается один раз, обычно из member_sync_init().
*
* inst — экземпляр uTun (нужен для etcp_route_send, uasync, sqlite3)
* inst — экземпляр uTun (нужен для etcp_send, uasync, sqlite3)
* svc_id — ID сервиса на ETCP-роутере
* ops — коллбэки модели данных (update_bucket_hash, get_items, apply_items)
* data_ctx — прозрачный контекст, передаваемый в коллбэки первым аргументом

Loading…
Cancel
Save