Browse Source

chat: member_sync per-peer throttling on peer sleep (cancel_peer + recv/start gating)

proxy
evgeny 1 week ago
parent
commit
0a1ba79d6f
  1. 16
      src/chat/chat_sync.c
  2. 34
      src/chat/merkle_sync.c
  3. 6
      src/chat/merkle_sync.h

16
src/chat/chat_sync.c

@ -394,6 +394,20 @@ static void cs_schedule_sync(struct chat_sync* cs) {
DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "%s: sync scheduled in %ums", CS_ID, CS_SYNC_INTERVAL_MS);
}
/* ── Peer sleep (keepalive-протокол): спячка → отменить сессии; пробуждение → пересинхронизация ── */
static void cs_on_peer_sleep(struct UTUN_INSTANCE* inst, uint64_t peer_node_id, int sleeping, void* arg) {
struct chat_sync* cs = (struct chat_sync*)arg;
if (!cs || !cs->initialized) return;
if (sleeping) {
DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: peer sleep peer=0x%016llx — cancel member_sync sessions", CS_ID, (unsigned long long)peer_node_id);
merkle_sync_cancel_peer(inst, peer_node_id);
} else {
DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: peer awake peer=0x%016llx — schedule sync", CS_ID, (unsigned long long)peer_node_id);
cs_schedule_sync(cs);
}
}
/* ── Peer status change (GUI + online-бар + schedule_sync; онлайн в БД не хранится) ── */
static void cs_on_peer_status_changed(struct chat_sync* cs, uint64_t peer, int online) {
@ -643,6 +657,7 @@ int chat_sync_init(struct UTUN_INSTANCE* inst) {
member_sync_init(inst);
member_sync_add_apply_cbk(inst, cs_on_member_applied, cs);
utun_add_peer_sleep_cbk(inst, cs_on_peer_sleep, cs);
chat_admin_init(inst);
chat_join_init(inst);
@ -658,6 +673,7 @@ void chat_sync_destroy(struct UTUN_INSTANCE* inst) {
if (cs->invite_send_timer) { uasync_cancel_timeout(inst->ua, cs->invite_send_timer); cs->invite_send_timer = NULL; }
if (cs->sync_scheduled) { cs->sync_scheduled = 0; cs_flush_sync(cs); }
member_sync_remove_apply_cbk(inst, cs_on_member_applied, cs);
utun_remove_peer_sleep_cbk(inst, cs_on_peer_sleep, cs);
cs_pending_free_all(cs);
member_sync_destroy(inst);
chat_admin_destroy(inst);

34
src/chat/merkle_sync.c

@ -692,6 +692,12 @@ static void _recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
struct merkle_sync* ms = inst ? inst->msync : NULL;
if (!ms || !ms->initialized) { u_free(entry->dgram); queue_entry_free(entry); return; }
uint64_t peer = conn ? conn->peer_node_id : 0;
/* пир в спячке — не обрабатываем его синхронизацию (server-side throttling) */
if (conn && conn->peer_sleeping) {
DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "%s: recv DROP — peer=0x%016llx sleeping", MS_ID, (unsigned long long)peer);
u_free(entry->dgram); queue_entry_free(entry); return;
}
const uint8_t* d = entry->dgram;
size_t dlen = entry->len;
@ -882,6 +888,15 @@ int merkle_sync_start(struct UTUN_INSTANCE* inst, uint64_t peer,
if (!inst || !ns || !ms || !ms->initialized) return -1;
_ensure_table(ms);
/* пир в спячке — не инициируем синхронизацию (server-side throttling) */
{
struct ETCP_CONN* pc = ms_find_conn_for_node(inst, peer);
if (pc && pc->peer_sleeping) {
DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "%s: start SKIP — peer=0x%016llx sleeping", MS_ID, (unsigned long long)peer);
return 0;
}
}
struct ms_session* s = _session_find(ms, peer, ns);
DEBUG_TRACE(DEBUG_CATEGORY_MEMBER_SYNC, "%s: start peer=%016llx ns=%s new=%d", MS_ID, (unsigned long long)peer, ns, s ? 0 : 1);
if (!s) {
@ -939,3 +954,22 @@ void merkle_sync_cancel_ns(struct UTUN_INSTANCE* inst, const char* ns) {
}
if (removed) DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: cancel_ns ns=%s removed=%d", MS_ID, ns, removed);
}
void merkle_sync_cancel_peer(struct UTUN_INSTANCE* inst, uint64_t peer) {
struct merkle_sync* ms = inst ? inst->msync : NULL;
if (!ms) return;
struct ms_session** p = &ms->sessions;
int removed = 0;
while (*p) {
struct ms_session* s = *p;
if (s->peer == peer) {
*p = s->next;
ms_session_outq_free(s);
u_free(s);
removed++;
continue;
}
p = &(*p)->next;
}
if (removed) DEBUG_INFO(DEBUG_CATEGORY_MEMBER_SYNC, "%s: cancel_peer peer=%016llx removed=%d", MS_ID, (unsigned long long)peer, removed);
}

6
src/chat/merkle_sync.h

@ -230,6 +230,12 @@ void merkle_sync_cancel(struct UTUN_INSTANCE* inst, uint64_t peer, const char* n
*/
void merkle_sync_cancel_ns(struct UTUN_INSTANCE* inst, const char* ns);
/*
* Отменить и удалить ВСЕ сессии с конкретным пиром peer (независимо от ns).
* Используется при уходе пира в спячку (server-side throttling). done_cb НЕ вызывается.
*/
void merkle_sync_cancel_peer(struct UTUN_INSTANCE* inst, uint64_t peer);
/* ── Recompute tree path after data change ── */
/*

Loading…
Cancel
Save