Browse Source

transport: per-peer sleep callback (peer_sleep_cbk) fired from keepalive protocol

proxy
evgeny 2 weeks ago
parent
commit
e4e9fca243
  1. 5
      src/transport_layer/etcp.c
  2. 1
      src/transport_layer/etcp.h
  3. 12
      src/transport_layer/etcp_keepalive.c
  4. 35
      src/utun_instance.c
  5. 14
      src/utun_instance.h

5
src/transport_layer/etcp.c

@ -400,6 +400,11 @@ void etcp_connection_close(struct ETCP_CONN* etcp) {
etcp->state = 2; // deleted — blocks ref_take and etcp_send immediately
if (etcp->links_up != 0) etcp->links_up = 0;
if (etcp->peer_sleeping) {
etcp->peer_sleeping = 0;
utun_fire_peer_sleep_cbk(etcp->instance, etcp->peer_node_id, 0);
}
etcp_fire_conn_status(etcp, ETCP_CONN_STATUS_DELETE);
DEBUG_DEBUG(DEBUG_CATEGORY_CONNECTION, "[%s] close: state=2 conn=%p peer=%016llx (DELETE fired)", etcp->log_name, (void*)etcp, (unsigned long long)etcp->peer_node_id);
etcp_cbk_fire(etcp, ETCP_CBK_EVENT_DOWN);

1
src/transport_layer/etcp.h

@ -125,6 +125,7 @@ struct ETCP_CONN {
// Peer info
uint64_t peer_node_id; // Peer node ID
uint8_t peer_ed25519_pubkey[SC_PUBKEY_SIZE]; // Ed25519 pubkey пира (из INIT)
uint8_t peer_sleeping; // 1 = пир в спячке (keepalive sleep-анонс), серверная сторона
// ============ Processing incoming data to be sent by ETCP
struct ll_queue* input_queue; // Incoming packets to send (rx_pool -> ETCP_FRAGMENT)

12
src/transport_layer/etcp_keepalive.c

@ -179,6 +179,10 @@ static void keepalive_timer_cb(void* arg) {
link->ka_sleeping = 0;
link->ka_period_ms = link->keepalive_interval;
link->keepalive_timeout = (uint32_t)link->keepalive_interval * KA_TIMEOUT_MULT;
if (link->etcp && link->etcp->peer_sleeping) {
link->etcp->peer_sleeping = 0;
utun_fire_peer_sleep_cbk(link->etcp->instance, link->etcp->peer_node_id, 0);
}
DEBUG_INFO(DEBUG_CATEGORY_KEEPALIVE, "[%s] server: sleep deadline expired, back to normal (interval=%dms)",
link->etcp->log_name, link->keepalive_interval);
}
@ -351,6 +355,10 @@ int etcp_keepalive_on_recv(struct ETCP_SOCKET* e_sock, struct ETCP_DGRAM* pkt,
link->ka_sleep_time = sleep_time;
link->ka_sleep_deadline_tb = get_time_tb() + (uint64_t)sleep_time * 1000;
link->ka_period_ms = peer_period;
if (link->etcp && !link->etcp->peer_sleeping) {
link->etcp->peer_sleeping = 1;
utun_fire_peer_sleep_cbk(link->etcp->instance, link->etcp->peer_node_id, 1);
}
DEBUG_INFO(DEBUG_CATEGORY_KEEPALIVE, "[%s] server: peer sleeping (period=%dms deadline=%us), echo ack",
link->etcp->log_name, peer_period, (unsigned)(sleep_time / 10));
etcp_link_send_keepalive(link);
@ -371,6 +379,10 @@ int etcp_keepalive_on_recv(struct ETCP_SOCKET* e_sock, struct ETCP_DGRAM* pkt,
link->ka_sleeping = 0;
link->ka_period_ms = (uint16_t)link->keepalive_interval;
link->keepalive_timeout = (uint32_t)link->keepalive_interval * KA_TIMEOUT_MULT;
if (link->etcp && link->etcp->peer_sleeping) {
link->etcp->peer_sleeping = 0;
utun_fire_peer_sleep_cbk(link->etcp->instance, link->etcp->peer_node_id, 0);
}
DEBUG_INFO(DEBUG_CATEGORY_KEEPALIVE, "[%s] server: peer woke up, back to normal (interval=%dms)",
link->etcp->log_name, link->keepalive_interval);
restart_keepalive_timer_ms(link, link->ka_period_ms);

35
src/utun_instance.c

@ -645,6 +645,12 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) {
instance->activity_cbks = r->next; u_free(r);
}
// Cleanup peer-sleep callback chain
while (instance->peer_sleep_cbks) {
struct peer_sleep_cbk_entry* r = instance->peer_sleep_cbks;
instance->peer_sleep_cbks = r->next; u_free(r);
}
// Shared SQLite DB закрываем в самом конце: routing/chat/media/db_sync
// используют его во время всего teardown (см. topo_groups_destroy).
if (instance->topo_sqlite_db) {
@ -1130,6 +1136,35 @@ static void utun_fire_activity_cbk(struct UTUN_INSTANCE* instance, int active) {
}
}
void utun_add_peer_sleep_cbk(struct UTUN_INSTANCE* instance, peer_sleep_cbk_fn fn, void* arg) {
if (!instance || !fn) return;
struct peer_sleep_cbk_entry* e = u_malloc(sizeof(*e));
if (!e) return;
e->fn = fn; e->arg = arg; e->next = instance->peer_sleep_cbks;
instance->peer_sleep_cbks = e;
}
void utun_remove_peer_sleep_cbk(struct UTUN_INSTANCE* instance, peer_sleep_cbk_fn fn, void* arg) {
if (!instance || !fn) return;
struct peer_sleep_cbk_entry** p = &instance->peer_sleep_cbks;
while (*p) {
if ((*p)->fn == fn && (*p)->arg == arg) {
struct peer_sleep_cbk_entry* rm = *p;
*p = rm->next; u_free(rm); return;
}
p = &(*p)->next;
}
}
void utun_fire_peer_sleep_cbk(struct UTUN_INSTANCE* instance, uint64_t peer_node_id, int sleeping) {
struct peer_sleep_cbk_entry* e = instance ? instance->peer_sleep_cbks : NULL;
while (e) {
struct peer_sleep_cbk_entry* n = e->next;
e->fn(instance, peer_node_id, sleeping, e->arg);
e = n;
}
}
static void client_activity_timeout_cb(void* arg) {
struct UTUN_INSTANCE* instance = (struct UTUN_INSTANCE*)arg;
DEBUG_INFO(DEBUG_CATEGORY_BGP, "client_activity_timer fired: node=%016llx type=%u -> standby",

14
src/utun_instance.h

@ -99,6 +99,16 @@ struct utun_activity_cbk_entry {
struct utun_activity_cbk_entry* next;
};
// Подписка на смену спячки пира (keepalive-протокол). Сервер (или любой узел)
// получает от мобильного пира sleep-анонс и уведомляет подписчиков:
// sleeping=1 — пир ушёл в спячку (в фон), 0 — проснулся (дедлайн истёк / обычный keepalive).
typedef void (*peer_sleep_cbk_fn)(struct UTUN_INSTANCE* instance, uint64_t peer_node_id, int sleeping, void* arg);
struct peer_sleep_cbk_entry {
peer_sleep_cbk_fn fn;
void* arg;
struct peer_sleep_cbk_entry* next;
};
// uTun instance configuration
struct UTUN_INSTANCE {
// Identification
@ -221,6 +231,7 @@ struct UTUN_INSTANCE {
uint8_t client_activity; // CLIENT_ACTIVITY_STANDBY/ACTIVE
void* client_activity_timer; // uasync timer handle for inactivity timeout
struct utun_activity_cbk_entry* activity_cbks; // подписки на смену client_activity
struct peer_sleep_cbk_entry* peer_sleep_cbks; // подписки на смену спячки пира (keepalive)
uint8_t standby_enabled; // 1 = standby duty-cycle активен (chatgui-android)
// TCP proxy server (exit node)
@ -259,6 +270,9 @@ void utun_instance_set_topo_group_enabled(int enabled);
void utun_set_client_activity(struct UTUN_INSTANCE* instance, int active);
void utun_add_activity_cbk(struct UTUN_INSTANCE* instance, utun_activity_cbk_fn fn, void* arg);
void utun_remove_activity_cbk(struct UTUN_INSTANCE* instance, utun_activity_cbk_fn fn, void* arg);
void utun_add_peer_sleep_cbk(struct UTUN_INSTANCE* instance, peer_sleep_cbk_fn fn, void* arg);
void utun_remove_peer_sleep_cbk(struct UTUN_INSTANCE* instance, peer_sleep_cbk_fn fn, void* arg);
void utun_fire_peer_sleep_cbk(struct UTUN_INSTANCE* instance, uint64_t peer_node_id, int sleeping);
// Diagnostic function for memory leak analysis
void utun_instance_diagnose_leaks(struct UTUN_INSTANCE* instance, const char* phase);

Loading…
Cancel
Save