Browse Source

chatgui: обновление статуса/адресов узлов через чат-протокол

- merkle_sync: synced-флаг в сессии, MSG_ITEM_UPDATE(0x04), push_update() — рассылка лёгких дельт (online-статус) всем synced-соседям
- member_sync: apply_update в ops, node_updated_cb для GUI-нотификации внешних обновлений
- chat_sync: cs_on_peer_status_changed() — единая функция: БД + GUI + push_update + throttled member_sync_start (1/сек)
- cs_on_conn_down: исправлено — теперь сбрасывает nodes.online=0
- flush при destroy
topo_upd
Evgeny 3 months ago
parent
commit
0b385dca50
  1. 98
      tools/chatgui/transport/chat_sync.c
  2. 19
      tools/chatgui/transport/member_sync.c
  3. 8
      tools/chatgui/transport/member_sync.h
  4. 55
      tools/chatgui/transport/merkle_sync.c
  5. 26
      tools/chatgui/transport/merkle_sync.h

98
tools/chatgui/transport/chat_sync.c

@ -3,6 +3,7 @@
#include "gui_bridge.h"
#include "topo_node_sqlite.h"
#include "member_sync.h"
#include "merkle_sync.h"
#include "../../../src/utun_instance.h"
#include "../../../src/etcp_api.h"
@ -269,7 +270,9 @@ struct chat_sync {
void* refresh_timer;
void* info_req_timer;
void* join_timer;
void* sync_timer;
uint8_t initialized;
uint8_t sync_scheduled;
uint64_t pending_invite_ch_id;
uint64_t pending_invite_node_id;
};
@ -277,6 +280,7 @@ struct chat_sync {
#define CS_ID "chat_sync"
#define CS_INFO_REQ_TIMEOUT_MS 5000
#define CS_JOIN_TIMEOUT_MS 8000
#define CS_SYNC_INTERVAL_MS 1000
static const char* cs_msg_name(uint8_t type) {
switch (type) {
@ -483,6 +487,87 @@ static void cs_post_channel_online(struct chat_sync* cs, const char* ch_id) {
gui_bridge_post(GUI_EVT_CHANNEL_PEERS_ONLINE, evt, 1 + (int)cl + 2);
}
/* ── Throttled sync (не чаще CS_SYNC_INTERVAL_MS) ── */
static void cs_flush_sync(struct chat_sync* cs) {
uint64_t myid = cs->inst->node_id;
for (int i = 0; i < cs->channel_count; i++) {
struct channel_cache* ch = &cs->channels[i];
for (int j = 0; j < ch->peer_count; j++) {
uint64_t pid = ch->peer_ids[j];
if (pid == myid || !cs_is_peer_online(cs->inst, pid)) continue;
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: flush_sync start for peer=%016llx ch=%s",
CS_ID, (unsigned long long)pid, ch->channel_id);
member_sync_start(cs->inst, pid, ch->channel_id, NULL, NULL);
}
}
}
static void cs_sync_timer_cb(void* arg) {
struct chat_sync* cs = (struct chat_sync*)arg;
if (!cs || !cs->initialized) return;
cs->sync_timer = NULL;
cs->sync_scheduled = 0;
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: sync timer fire", CS_ID);
cs_flush_sync(cs);
}
static void cs_schedule_sync(struct chat_sync* cs) {
if (cs->sync_scheduled) return;
cs->sync_scheduled = 1;
cs->sync_timer = uasync_set_timeout(cs->inst->ua,
CS_SYNC_INTERVAL_MS * 10, cs, cs_sync_timer_cb, "cs_sync");
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: sync scheduled in %ums", CS_ID, CS_SYNC_INTERVAL_MS);
}
/* ── Unified peer status change (локальная БД + GUI + push_update + schedule_sync) ── */
static void cs_on_remote_status_changed(uint64_t peer, int online) {
if (!g_cs) return;
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: remote status changed peer=%016llx online=%d",
CS_ID, (unsigned long long)peer, online);
for (int i = 0; i < g_cs->channel_count; i++) {
for (int j = 0; j < g_cs->channels[i].peer_count; j++) {
if (g_cs->channels[i].peer_ids[j] == peer) {
size_t cl = strlen(g_cs->channels[i].channel_id);
if (cl > 63) cl = 63;
uint8_t evt[65]; evt[0] = (uint8_t)cl;
memcpy(evt + 1, g_cs->channels[i].channel_id, cl);
gui_bridge_post(GUI_EVT_MEMBERS_CHANGED, evt, 1 + (int)cl);
cs_post_channel_online(g_cs, g_cs->channels[i].channel_id);
break;
}
}
}
}
static void cs_on_peer_status_changed(uint64_t peer, int online) {
/* 1. Локальная БД */
member_sync_set_online(g_cs->inst, peer, online);
/* 2. GUI + online-бар */
for (int i = 0; i < g_cs->channel_count; i++) {
for (int j = 0; j < g_cs->channels[i].peer_count; j++) {
if (g_cs->channels[i].peer_ids[j] == peer) {
size_t cl = strlen(g_cs->channels[i].channel_id);
if (cl > 63) cl = 63;
uint8_t evt[65]; evt[0] = (uint8_t)cl;
memcpy(evt + 1, g_cs->channels[i].channel_id, cl);
gui_bridge_post(GUI_EVT_MEMBERS_CHANGED, evt, 1 + (int)cl);
cs_post_channel_online(g_cs, g_cs->channels[i].channel_id);
/* 3. Рассылка дельты всем synced-соседям */
uint8_t st = (uint8_t)online;
merkle_sync_push_update(g_cs->inst, g_cs->channels[i].channel_id, peer, 0x01, &st, 1);
break;
}
}
}
/* 4. Запланировать member_sync_start (throttled) */
cs_schedule_sync(g_cs);
}
static void _on_member_sync_done(uint64_t peer, const char* ns, int result, void* arg) {
struct channel_cache* ch = (struct channel_cache*)arg;
if (result == MT_OK && ch) ch->synced = CS_SYNC_DONE;
@ -517,11 +602,7 @@ static void cs_on_conn_up(struct ETCP_CONN* conn, void* arg) {
return;
}
member_sync_set_online(g_cs->inst, peer, 1);
for (int i = 0; i < g_cs->channel_count; i++) {
for (int j = 0; j < g_cs->channels[i].peer_count; j++)
if (g_cs->channels[i].peer_ids[j] == peer) { cs_post_channel_online(g_cs, g_cs->channels[i].channel_id); break; }
}
cs_on_peer_status_changed(peer, 1);
}
static void cs_on_conn_down(struct ETCP_CONN* conn, void* arg) {
@ -529,6 +610,9 @@ static void cs_on_conn_down(struct ETCP_CONN* conn, void* arg) {
if (!conn || !g_cs) return;
uint64_t peer = conn->peer_node_id;
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: conn_down peer=%016llx", CS_ID, (unsigned long long)peer);
cs_on_peer_status_changed(peer, 0);
uint16_t rtt = conn->rtt_avg_100;
if (rtt > 0 && peer != 0 && g_cs->inst->topo_sqlite_db) {
sqlite3* db = g_cs->inst->topo_sqlite_db;
@ -552,7 +636,6 @@ static void cs_on_conn_down(struct ETCP_CONN* conn, void* arg) {
uint8_t rem[9]; rem[0] = CS_MSG_PEER_REMOVE;
memcpy(rem + 1, &peer, 8);
cs_propagate(g_cs, g_cs->channels[i].channel_id, peer, rem, 9);
cs_post_channel_online(g_cs, g_cs->channels[i].channel_id);
}
}
}
@ -653,6 +736,7 @@ int chat_sync_init(struct UTUN_INSTANCE* inst,
cs->join_timer = NULL;
member_sync_init(inst);
member_sync_set_node_updated_cb(cs_on_remote_status_changed);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: initialized", CS_ID);
return 0;
@ -663,6 +747,8 @@ void chat_sync_destroy(struct UTUN_INSTANCE* inst) {
if (!cs || !inst) return;
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: destroy", CS_ID);
chat_sync_auto_connect_stop();
if (cs->sync_timer) { uasync_cancel_timeout(inst->ua, cs->sync_timer); cs->sync_timer = NULL; }
if (cs->sync_scheduled) { cs->sync_scheduled = 0; cs_flush_sync(cs); }
member_sync_destroy(inst);
cs->initialized = 0; g_cs = NULL;

19
tools/chatgui/transport/member_sync.c

@ -280,10 +280,25 @@ static int _member_apply_items(void* ctx, const char* ns,
return 0;
}
static member_sync_node_updated_fn g_node_updated_cb = NULL;
static int _ms_apply_update(void* ctx, const char* ns, uint64_t key,
uint8_t type, const uint8_t* data, size_t len) {
(void)ns;
struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)ctx;
if (type == 0x01 && len >= 1) {
int online = data[0];
member_sync_set_online(inst, key, online);
if (g_node_updated_cb) g_node_updated_cb(key, online);
}
return 0;
}
static const struct merkle_sync_data_ops g_member_ops = {
.update_bucket_hash = _member_update_bucket_hash,
.get_items = _member_get_items,
.apply_items = _member_apply_items,
.apply_update = _ms_apply_update,
};
/* ── Public API ── */
@ -387,6 +402,10 @@ int member_sync_count(struct UTUN_INSTANCE* inst, const char* ch_id) {
return c;
}
void member_sync_set_node_updated_cb(member_sync_node_updated_fn cb) {
g_node_updated_cb = cb;
}
void member_sync_set_online(struct UTUN_INSTANCE* inst, uint64_t node_id, int online) {
if (!inst) return;
DEBUG_TRACE(DEBUG_CATEGORY_DB_SYNC, "%s: set_online nid=%016llx online=%d", MS_ID, (unsigned long long)node_id, online);

8
tools/chatgui/transport/member_sync.h

@ -119,6 +119,14 @@ void member_sync_set_online(struct UTUN_INSTANCE* inst, uint64_t node_id, int on
const uint8_t* member_sync_get_hash(struct UTUN_INSTANCE* inst,
const char* ch_id, uint8_t level, uint64_t prefix64);
/*
* Коллбэк: узел изменился по инициативе удалённого пира (через MSG_ITEM_UPDATE).
* Вызывается из uasync-потока. Потребитель (chat_sync) может из него
* постить GUI-события.
*/
typedef void (*member_sync_node_updated_fn)(uint64_t node_id, int online);
void member_sync_set_node_updated_cb(member_sync_node_updated_fn cb);
#ifdef __cplusplus
}
#endif

55
tools/chatgui/transport/merkle_sync.c

@ -24,6 +24,7 @@ struct ms_session {
char ns[64];
uint64_t peer;
uint8_t active;
uint8_t synced;
void* timer;
uint8_t retries;
merkle_sync_done_cb done_cb;
@ -46,6 +47,8 @@ static sqlite3* _db(struct UTUN_INSTANCE* inst) {
return inst ? inst->topo_sqlite_db : NULL;
}
static int _send_msg(struct merkle_sync* ms, uint64_t peer, const uint8_t* payload, size_t len);
/* ── Prefix arithmetic ── */
uint64_t merkle_sync_level_prefix(uint64_t key, uint8_t level) {
@ -328,6 +331,8 @@ static struct ms_session* _session_find(struct merkle_sync* ms, uint64_t peer, c
static void _session_done(struct ms_session* s, int result) {
s->active = 0;
if (s->timer) { uasync_cancel_timeout(g_merkle->inst->ua, s->timer); s->timer = NULL; }
if (result == MT_OK) { s->synced = 1; DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: session SYNCED peer=%016llx ns=%s", MS_ID, (unsigned long long)s->peer, s->ns); }
else { s->synced = 0; DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: session FAILED peer=%016llx ns=%s result=%d", MS_ID, (unsigned long long)s->peer, s->ns, result); }
if (s->done_cb) s->done_cb(s->peer, s->ns, result, s->cb_arg);
}
@ -514,11 +519,60 @@ static void _recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
case 0x01: _handle_hashes(g_merkle, peer, ns, pl, plen); break;
case 0x02: _handle_request(g_merkle, peer, ns, pl, plen); break;
case 0x03: _handle_batch(g_merkle, peer, ns, pl, plen); break;
case 0x04:
if (plen >= 9 && g_merkle->ops->apply_update) {
uint64_t key; memcpy(&key, pl, 8);
uint8_t utype = pl[8];
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: recv ITEM_UPDATE key=%016llx type=%02x from=%016llx ns=%s",
MS_ID, (unsigned long long)key, utype, (unsigned long long)peer, ns);
g_merkle->ops->apply_update(g_merkle->data_ctx, ns, key, utype, pl + 9, plen - 9);
}
break;
}
u_free(entry->dgram); queue_entry_free(entry);
}
/* ── Push lightweight update to synced peers ── */
void merkle_sync_push_update(struct UTUN_INSTANCE* inst, const char* ns,
uint64_t key, uint8_t type, const uint8_t* data, size_t len) {
struct merkle_sync* ms = g_merkle;
if (!ms || !ms->initialized || !ns || !data) return;
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: push_update ns=%s key=%016llx type=%02x len=%zu",
MS_ID, ns, (unsigned long long)key, type, len);
uint8_t ch_len = (uint8_t)strlen(ns);
if (ch_len > 63) return;
size_t pkt = 1 + 1 + ch_len + 1 + 8 + 1 + len;
uint8_t* buf = u_malloc(pkt);
if (!buf) return;
uint8_t* p = buf;
*p++ = ms->svc_id;
*p++ = ch_len; memcpy(p, ns, ch_len); p += ch_len;
*p++ = 0x04; /* MSG_ITEM_UPDATE */
memcpy(p, &key, 8); p += 8;
*p++ = type;
memcpy(p, data, len);
for (struct ms_session* s = ms->sessions; s; s = s->next) {
if (!s->synced || strcmp(s->ns, ns) != 0) continue;
struct ll_entry* entry = queue_entry_new(0);
if (!entry) continue;
uint8_t* dcopy = u_malloc(pkt);
if (!dcopy) { queue_entry_free(entry); continue; }
memcpy(dcopy, buf, pkt);
entry->dgram = dcopy; entry->len = pkt;
struct ETCP_CONN* conn = ms_find_conn_for_node(inst, s->peer);
if (!conn) { u_free(dcopy); queue_entry_free(entry); continue; }
int r = etcp_send(conn, entry);
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: push_update sent type=%02x to=%016llx rc=%d",
MS_ID, type, (unsigned long long)s->peer, r);
if (r != 0) { u_free(dcopy); queue_entry_free(entry); }
}
u_free(buf);
}
/* ── Background consistency check ── */
int merkle_sync_bg_check(struct UTUN_INSTANCE* inst, const char* ns) {
@ -633,6 +687,7 @@ int merkle_sync_start(struct UTUN_INSTANCE* inst, uint64_t peer,
} else {
/* replace callback if session already exists */
if (s->timer) { uasync_cancel_timeout(inst->ua, s->timer); s->timer = NULL; }
s->synced = 0;
}
s->active = 1; s->retries = 0;
s->done_cb = done_cb; s->cb_arg = arg;

26
tools/chatgui/transport/merkle_sync.h

@ -140,6 +140,20 @@ struct merkle_sync_data_ops {
*/
int (*apply_items)(void* ctx, const char* ns,
const uint8_t* data, size_t len);
/*
* Применить лёгкое обновление (не влияющее на Merkle-хеш).
* Вызывается при приёме MSG_ITEM_UPDATE(0x04) — online-статус и т.п.
*
* ns — namespace (channel_id)
* key — node_id изменившегося узла
* type — тип обновления (0x01 = online-статус)
* data,len — данные обновления (type=0x01: [online:1])
*
* Возвращает 0 при успехе, <0 при ошибке.
*/
int (*apply_update)(void* ctx, const char* ns, uint64_t key,
uint8_t type, const uint8_t* data, size_t len);
};
/*
@ -211,6 +225,18 @@ void merkle_sync_cancel(struct UTUN_INSTANCE* inst, uint64_t peer, const char* n
*/
void merkle_sync_recompute_path(struct UTUN_INSTANCE* inst, const char* ns, uint64_t key);
/*
* Отправить лёгкое обновление всем synced-пирам в namespace ns.
* Не влияет на Merkle-дерево — используется для online-статуса и т.п.
*
* ns — namespace (channel_id)
* key — node_id изменившегося узла
* type — тип обновления (0x01 = online-статус)
* data,len — данные обновления
*/
void merkle_sync_push_update(struct UTUN_INSTANCE* inst, const char* ns,
uint64_t key, uint8_t type, const uint8_t* data, size_t len);
/* ── Background consistency check ── */
/*

Loading…
Cancel
Save