Browse Source

chat: троттлинг по фазе спячки (member_sync send-only, db_sync sound-aware) + настройка звука

proxy
evgeny 2 weeks ago
parent
commit
b39b85db13
  1. 29
      src/chat/chat_core.c
  2. 5
      src/chat/chat_core.h
  3. 35
      src/chat/chat_sync.c
  4. 4
      src/chat/chat_sync.h
  5. 28
      src/chat/db_sync.c
  6. 3
      src/chat/db_sync.h
  7. 10
      src/chat/merkle_sync.c

29
src/chat/chat_core.c

@ -15,6 +15,7 @@
#include "chat_setting.h" #include "chat_setting.h"
#include "chat_member.h" #include "chat_member.h"
#include "member_sync.h" #include "member_sync.h"
#include "chat_sync.h"
#include "../utun_instance.h" #include "../utun_instance.h"
#include "../ntp_time.h" #include "../ntp_time.h"
@ -259,6 +260,34 @@ void chat_core_set_setting_trampoline(void* arg) {
u_free(arg); u_free(arg);
} }
/* ─── per-channel «play sound on new message» ─── */
int chat_core_get_sound_on_message(struct UTUN_INSTANCE* inst, const char* ch_id) {
struct chat_core_ctx* cc = CC(inst);
if (!cc || !cc->initialized || !ch_id || !ch_id[0]) return 1;
char key[128]; snprintf(key, sizeof(key), "ch_%s_sound_on_message", ch_id);
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(cc->db, "SELECT value FROM ui_state WHERE key=?", -1, &st, NULL) != SQLITE_OK) return 1;
sqlite3_bind_text(st, 1, key, -1, SQLITE_STATIC);
int on = 1;
if (sqlite3_step(st) == SQLITE_ROW) {
const char* v = (const char*)sqlite3_column_text(st, 0);
if (v) on = atoi(v) != 0;
}
sqlite3_finalize(st);
return on;
}
void chat_core_set_sound_on_message(struct UTUN_INSTANCE* inst, const char* ch_id, int on) {
struct chat_core_ctx* cc = CC(inst);
if (!cc || !cc->initialized || !ch_id || !ch_id[0]) return;
char key[128]; snprintf(key, sizeof(key), "ch_%s_sound_on_message", ch_id);
char val[2]; snprintf(val, sizeof(val), "%d", on ? 1 : 0);
chat_core_save_ui_state(inst, key, val);
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "%s: sound_on_message ch=%s = %d", CC_ID, ch_id, on ? 1 : 0);
chat_sync_broadcast_sound_setting(inst, ch_id, on ? 1 : 0);
}
/* ─── headless API ─── */ /* ─── headless API ─── */
static void json_escape(const char* src, char* dst, size_t dst_sz) { static void json_escape(const char* src, char* dst, size_t dst_sz) {

5
src/chat/chat_core.h

@ -120,6 +120,11 @@ void chat_core_update_my_member(struct UTUN_INSTANCE* inst);
/* Сохранение key-value в ui_state (вызывается из uasync-потока) */ /* Сохранение key-value в ui_state (вызывается из uasync-потока) */
void chat_core_save_ui_state(struct UTUN_INSTANCE* inst, const char* key, const char* value); void chat_core_save_ui_state(struct UTUN_INSTANCE* inst, const char* key, const char* value);
/* Per-channel «play sound on new message» (ui_state: ch_<id>_sound_on_message, default 1=on).
* set — сохраняет и рассылает пирам через chat_sync (троттлинг db_sync на сервере). */
int chat_core_get_sound_on_message(struct UTUN_INSTANCE* inst, const char* ch_id);
void chat_core_set_sound_on_message(struct UTUN_INSTANCE* inst, const char* ch_id, int on);
struct save_ui_state_arg { struct UTUN_INSTANCE* inst; char data[256]; }; struct save_ui_state_arg { struct UTUN_INSTANCE* inst; char data[256]; };
void chat_core_save_ui_state_trampoline(void* arg); void chat_core_save_ui_state_trampoline(void* arg);

35
src/chat/chat_sync.c

@ -99,6 +99,7 @@ static const char* cs_msg_name(uint8_t type) {
case CS_MSG_JOIN_INFO_RESP: return "JOIN_INFO_RESP"; case CS_MSG_JOIN_INFO_RESP: return "JOIN_INFO_RESP";
case CS_MSG_JOIN_REQUEST: return "JOIN_REQUEST"; case CS_MSG_JOIN_REQUEST: return "JOIN_REQUEST";
case CS_MSG_JOIN_READY: return "JOIN_READY"; case CS_MSG_JOIN_READY: return "JOIN_READY";
case CS_MSG_SOUND_SETTING: return "SOUND_SETTING";
default: return "???"; default: return "???";
} }
} }
@ -239,6 +240,10 @@ static void chat_sync_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
break; break;
} }
case CS_MSG_CHANNEL_INVITE: cs_handle_channel_invite(cs, peer, ch_id, pl, plen); break; case CS_MSG_CHANNEL_INVITE: cs_handle_channel_invite(cs, peer, ch_id, pl, plen); break;
case CS_MSG_SOUND_SETTING:
if (plen >= 1)
db_sync_set_peer_sound(cs->inst, group_id, peer, pl[0] ? 1 : 0);
break;
default: DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: UNKNOWN msg type=%02x from=%016llx", CS_ID, type, (unsigned long long)peer); break; default: DEBUG_WARN(DEBUG_CATEGORY_MEMBER_SYNC, "%s: UNKNOWN msg type=%02x from=%016llx", CS_ID, type, (unsigned long long)peer); break;
} }
u_free(entry->dgram); queue_entry_free(entry); u_free(entry->dgram); queue_entry_free(entry);
@ -408,6 +413,35 @@ static void cs_on_peer_sleep(struct UTUN_INSTANCE* inst, uint64_t peer_node_id,
} }
} }
/* ── Per-chat «play sound on message» (троттлинг db_sync на сервере) ── */
static void cs_send_sound_setting(struct chat_sync* cs, const char* ch_id, uint64_t dst, int on) {
uint8_t p[2] = { CS_MSG_SOUND_SETTING, (uint8_t)(on ? 1 : 0) };
cs_send(cs, ch_id, dst, p, sizeof(p));
}
void chat_sync_broadcast_sound_setting(struct UTUN_INSTANCE* inst, const char* ch_id, int on) {
struct chat_sync* cs = cs_of(inst);
if (!cs || !cs->initialized || !ch_id || !ch_id[0]) return;
for (int i = 0; i < cs->channel_count; i++) {
if (strcmp(cs->channels[i].channel_id, ch_id) != 0) continue;
for (int j = 0; j < cs->channels[i].peer_count; j++) {
uint64_t pid = cs->channels[i].peer_ids[j];
if (pid == inst->node_id || !cs_is_peer_online(inst, pid)) continue;
cs_send_sound_setting(cs, ch_id, pid, on);
}
break;
}
}
/* На поднятии соединения: отправить пиру свои per-chat настройки звука (начальная синхронизация). */
static void cs_send_all_sound_settings(struct chat_sync* cs, uint64_t peer) {
for (int i = 0; i < cs->channel_count; i++) {
int on = chat_core_get_sound_on_message(cs->inst, cs->channels[i].channel_id);
cs_send_sound_setting(cs, cs->channels[i].channel_id, peer, on);
}
}
/* ── Peer status change (GUI + online-бар + schedule_sync; онлайн в БД не хранится) ── */ /* ── Peer status change (GUI + online-бар + schedule_sync; онлайн в БД не хранится) ── */
static void cs_on_peer_status_changed(struct chat_sync* cs, uint64_t peer, int online) { static void cs_on_peer_status_changed(struct chat_sync* cs, uint64_t peer, int online) {
@ -536,6 +570,7 @@ static void cs_on_conn_up(struct ETCP_CONN* conn, int event, void* arg) { (void)
} }
cs_on_peer_status_changed(cs, peer, 1); cs_on_peer_status_changed(cs, peer, 1);
cs_send_all_sound_settings(cs, peer);
} }
static void cs_on_conn_down(struct ETCP_CONN* conn, int event, void* arg) { (void)event; static void cs_on_conn_down(struct ETCP_CONN* conn, int event, void* arg) { (void)event;

4
src/chat/chat_sync.h

@ -51,6 +51,7 @@ struct UASYNC;
#define CS_MSG_JOIN_INFO_RESP 0x11 /* connection → joiner: channel info + connection's member record */ #define CS_MSG_JOIN_INFO_RESP 0x11 /* connection → joiner: channel info + connection's member record */
#define CS_MSG_JOIN_REQUEST 0x12 /* joiner → connection: {join_key:8, member...} */ #define CS_MSG_JOIN_REQUEST 0x12 /* joiner → connection: {join_key:8, member...} */
#define CS_MSG_JOIN_READY 0x13 /* connection → joiner: {group_id:8} */ #define CS_MSG_JOIN_READY 0x13 /* connection → joiner: {group_id:8} */
#define CS_MSG_SOUND_SETTING 0x14 /* peer → server: {sound_on:1} — per-chat play sound on message */
/* Protocol constants */ /* Protocol constants */
#define CS_SEND_DATA_MAX 32 #define CS_SEND_DATA_MAX 32
@ -114,6 +115,9 @@ void chat_sync_retry_channels_on_socket_change(struct UTUN_INSTANCE* inst);
* чтобы cs_on_conn_status не запускал member_sync для удалённого ns до следующего refresh. */ * чтобы cs_on_conn_status не запускал member_sync для удалённого ns до следующего refresh. */
void chat_sync_on_channel_deleted(struct UTUN_INSTANCE* inst, const char* ch_id); void chat_sync_on_channel_deleted(struct UTUN_INSTANCE* inst, const char* ch_id);
/* Разослать всем подключённым пирам канала смену per-chat настройки «play sound on message». */
void chat_sync_broadcast_sound_setting(struct UTUN_INSTANCE* inst, const char* ch_id, int on);
#ifdef __cplusplus #ifdef __cplusplus
} }
#endif #endif

28
src/chat/db_sync.c

@ -51,6 +51,7 @@ struct SI_PEER {
uint8_t sync_state; uint8_t sync_state;
uint64_t sync_start_tb; uint64_t sync_start_tb;
uint32_t retry_count; uint32_t retry_count;
uint8_t sound_on; /* 1 = пиру важен этот чат (play sound on message), не троттлим */
}; };
struct DB_SYNC_INSTANCE { struct DB_SYNC_INSTANCE {
@ -311,6 +312,7 @@ static struct SI_PEER* si_peer_add(struct DB_SYNC_INSTANCE* si, uint64_t node_id
p = &si->peers[si->peer_count++]; p = &si->peers[si->peer_count++];
memset(p, 0, sizeof(*p)); memset(p, 0, sizeof(*p));
p->node_id = node_id; p->node_id = node_id;
p->sound_on = 1;
return p; return p;
} }
@ -1475,6 +1477,22 @@ void db_sync_peer_set_state(struct DB_SYNC_INSTANCE* si, uint64_t node_id, uint8
if (p) p->sync_state = state; if (p) p->sync_state = state;
} }
// Задать per-(chat,peer) флаг «play sound on message» (троттлинг: sound_off + спящий пир → пропуск).
void db_sync_set_peer_sound(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t peer_node_id, int sound_on)
{
if (!inst || !inst->db_sync) return;
struct DB_SYNC* db = inst->db_sync;
struct DB_SYNC_INSTANCE* si = db_instance_find(db, group_id);
if (!si) return; /* чат ещё не создан у нас — настройка не критична */
struct SI_PEER* p = si_peer_add(si, peer_node_id);
if (p && p->sound_on != (uint8_t)(sound_on ? 1 : 0)) {
p->sound_on = sound_on ? 1 : 0;
DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "sync [%s] peer=%016llx sound_on=%d (throttle=%s)",
SI_SHRT(si), (unsigned long long)peer_node_id, p->sound_on,
p->sound_on ? "off" : "on");
}
}
// Гейт авто-синхронизации: пока gated — инстанс не запускает sync автоматически // Гейт авто-синхронизации: пока gated — инстанс не запускает sync автоматически
// (ждёт member_sync/pubkeys). При снятии гейта запускает sync к пирам в sync_state==0. // (ждёт member_sync/pubkeys). При снятии гейта запускает sync к пирам в sync_state==0.
void db_sync_instance_set_gated(struct DB_SYNC_INSTANCE* si, int gated) void db_sync_instance_set_gated(struct DB_SYNC_INSTANCE* si, int gated)
@ -1552,6 +1570,15 @@ static void db_sync_resume_peer_check(struct DB_SYNC* db)
// Периодическая проверка: запускает синхронизацию для пиров в состоянии sync_state==0, // Периодическая проверка: запускает синхронизацию для пиров в состоянии sync_state==0,
// детектирует таймауты синхронизации (DB_SYNC_SYNC_TIMEOUT секунд), логирует зависшие. // детектирует таймауты синхронизации (DB_SYNC_SYNC_TIMEOUT секунд), логирует зависшие.
// Work-driven: перевзводится только пока есть пиры с sync_state != 2. // Work-driven: перевзводится только пока есть пиры с sync_state != 2.
/* Троттлинг: пир в SLEEP-фазе И чат для него «тихий» (sound_on=0) — пропускаем. */
static int db_sync_peer_throttled(struct DB_SYNC* db, struct SI_PEER* p) {
if (!db || !p) return 0;
if (p->sound_on) return 0; /* чат пиру важен (play sound) — не троттлим */
struct ETCP_CONN* conn = instance_find_conn(db->inst, p->node_id);
return conn && conn->peer_sleep_phase;
}
static void db_sync_peer_check_cb(void* arg) static void db_sync_peer_check_cb(void* arg)
{ {
struct DB_SYNC* db = (struct DB_SYNC*)arg; struct DB_SYNC* db = (struct DB_SYNC*)arg;
@ -1591,6 +1618,7 @@ static void db_sync_peer_check_cb(void* arg)
uint32_t min_pos = UINT32_MAX; uint32_t min_pos = UINT32_MAX;
for (int j = 0; j < si->peer_count; j++) { for (int j = 0; j < si->peer_count; j++) {
if (si->peers[j].sync_state == 0 && si->peers[j].synced_pos < min_pos) { if (si->peers[j].sync_state == 0 && si->peers[j].synced_pos < min_pos) {
if (db_sync_peer_throttled(db, &si->peers[j])) continue; /* спящий пир в «тихом» чате — пропускаем */
min_pos = si->peers[j].synced_pos; min_pos = si->peers[j].synced_pos;
best = &si->peers[j]; best = &si->peers[j];
} }

3
src/chat/db_sync.h

@ -138,6 +138,9 @@ int db_sync_chain_verify(struct DB_SYNC_INSTANCE* si);
// Peer state control (for testing — disable PUSH to isolate sync protocol) // Peer state control (for testing — disable PUSH to isolate sync protocol)
void db_sync_peer_set_state(struct DB_SYNC_INSTANCE* si, uint64_t node_id, uint8_t state); void db_sync_peer_set_state(struct DB_SYNC_INSTANCE* si, uint64_t node_id, uint8_t state);
// Per-(chat,peer) «play sound on message» flag: sound_off + sleeping peer → throttled.
void db_sync_set_peer_sound(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t peer_node_id, int sound_on);
// Gate auto-sync: while gated the instance does not auto-initiate sync (waits for // Gate auto-sync: while gated the instance does not auto-initiate sync (waits for
// member_sync/pubkeys). Un-gating launches sync to peers in sync_state==0. // member_sync/pubkeys). Un-gating launches sync to peers in sync_state==0.
void db_sync_instance_set_gated(struct DB_SYNC_INSTANCE* si, int gated); void db_sync_instance_set_gated(struct DB_SYNC_INSTANCE* si, int gated);

10
src/chat/merkle_sync.c

@ -692,12 +692,6 @@ static void _recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
struct merkle_sync* ms = inst ? inst->msync : NULL; struct merkle_sync* ms = inst ? inst->msync : NULL;
if (!ms || !ms->initialized) { u_free(entry->dgram); queue_entry_free(entry); return; } if (!ms || !ms->initialized) { u_free(entry->dgram); queue_entry_free(entry); return; }
uint64_t peer = conn ? conn->peer_node_id : 0; 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; const uint8_t* d = entry->dgram;
size_t dlen = entry->len; size_t dlen = entry->len;
@ -888,10 +882,10 @@ int merkle_sync_start(struct UTUN_INSTANCE* inst, uint64_t peer,
if (!inst || !ns || !ms || !ms->initialized) return -1; if (!inst || !ns || !ms || !ms->initialized) return -1;
_ensure_table(ms); _ensure_table(ms);
/* пир в спячке — не инициируем синхронизацию (server-side throttling) */ /* пир в SLEEP-фазе — не инициируем синхронизацию (server-side throttling) */
{ {
struct ETCP_CONN* pc = ms_find_conn_for_node(inst, peer); struct ETCP_CONN* pc = ms_find_conn_for_node(inst, peer);
if (pc && pc->peer_sleeping) { if (pc && pc->peer_sleep_phase) {
DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "%s: start SKIP — peer=0x%016llx sleeping", MS_ID, (unsigned long long)peer); DEBUG_DEBUG(DEBUG_CATEGORY_MEMBER_SYNC, "%s: start SKIP — peer=0x%016llx sleeping", MS_ID, (unsigned long long)peer);
return 0; return 0;
} }

Loading…
Cancel
Save