Browse Source

Deliver DM history snapshots asynchronously as JSON events

master
evgeny 2 days ago
parent
commit
139356bd18
  1. 2
      src/chat/chat_event.h
  2. 88
      src/dm/dm_core.c
  3. 4
      src/dm/dm_core.h

2
src/chat/chat_event.h

@ -61,7 +61,7 @@ extern "C" {
#define CHAT_EVT_CALL_DECLINED 38 /* [call_id:8][reason:1] */
#define CHAT_EVT_CALL_STATS 39 /* [call_id:8][rtt:2][buffer:2][tempo:2][dropped:2][underruns:2][min:2][max:2][reserve:2] */
#define CHAT_EVT_CALL_ERROR 40 /* [call_id:8][err:1][text:var] */
#define CHAT_EVT_DM_MESSAGES 41 /* [conv_id_len:1][conv_id:var][count:2][msg:var]* — расшифрованные сообщения беседы */
#define CHAT_EVT_DM_MESSAGES 41 /* [conv_id_len:1][conv_id:var][JSON:utf8] — сообщения беседы и состояния доставки */
#define CHAT_EVT_CALL_PATH 42 /* [call_id:8][text:var] — BGP-маршрут до пира звонка */
#define CHAT_EVT_RADIO_TALK 43 /* [group_id:8][src_node_id:8][stream_id:2][on:1] — кто-то начал/закончил говорить */
#define CHAT_EVT_CALL_CONNECTION 44 /* [call_id:8][CONN_TYPE:1][next_hop_node_id:8] — текущий путь отправки звонка */

88
src/dm/dm_core.c

@ -999,85 +999,21 @@ void dm_send_trampoline(void* arg) {
DM_ID, conv);
}
/* Собрать расшифрованные сообщения беседы и отправить в GUI бинарным событием:
* [conv_id_len:1][conv_id:var][count:2]([dir:1][seq:8][ts:8][author:8][ct_len:1][ct:var][data_len:2][data:var])* */
/* Ограниченный JSON-снимок доставляется GUI асинхронно, без доступа UI к instance/SQLite. */
void dm_messages_trampoline(void* arg) {
struct dm_messages_req* req = (struct dm_messages_req*)arg;
struct dm_messages_req* req = arg;
if (!req) return;
struct UTUN_INSTANCE* inst = req->inst;
char conv_id[64]; strncpy(conv_id, req->conv_id, sizeof(conv_id) - 1); conv_id[sizeof(conv_id) - 1] = '\0';
u_free(arg);
struct dm_state* dm = dm_of(inst);
if (!dm || !dm->initialized) return;
struct dm_conv c;
if (dm_conv_load(dm, conv_id, &c) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: dm_messages_trampoline — no conversation %s", DM_ID, conv_id);
return;
}
uint64_t conv_num = (uint64_t)strtoull(conv_id, NULL, 10);
sqlite3_stmt* st = NULL;
uint8_t* msgs = NULL;
size_t msgs_len = 0, msgs_cap = 0;
int count = 0;
if (sqlite3_prepare_v2(dm->db,
"SELECT dir,seq,ts,author,ct,data FROM dm_messages WHERE conv_id=? ORDER BY ts, seq",
-1, &st, NULL) == SQLITE_OK) {
sqlite3_bind_text(st, 1, conv_id, -1, SQLITE_STATIC);
while (sqlite3_step(st) == SQLITE_ROW) {
uint8_t dir = (uint8_t)sqlite3_column_int(st, 0);
uint64_t seq = (uint64_t)sqlite3_column_int64(st, 1);
uint64_t ts = (uint64_t)sqlite3_column_int64(st, 2);
uint64_t author = (uint64_t)sqlite3_column_int64(st, 3);
const char* ct = (const char*)sqlite3_column_text(st, 4);
const uint8_t* enc = (const uint8_t*)sqlite3_column_blob(st, 5);
int enc_len = sqlite3_column_bytes(st, 5);
size_t ct_len = ct ? strnlen(ct, 255) : 0;
uint8_t plain[DM_DATA_MAX];
size_t plen = 0;
if (enc && enc_len >= (int)DM_TAG_SIZE) {
uint8_t nonce[DM_NONCE_SIZE];
dm_build_nonce(conv_num, author, seq, nonce);
if (dm_decrypt(c.content_key, nonce, enc, (size_t)enc_len, plain, &plen) != 0)
plen = 0;
}
size_t need = 1 + 8 + 8 + 8 + 1 + ct_len + 2 + plen;
if (msgs_len + need > msgs_cap) {
msgs_cap = msgs_cap ? msgs_cap * 2 : 4096;
while (msgs_cap < msgs_len + need) msgs_cap *= 2;
uint8_t* tmp = u_realloc(msgs, msgs_cap);
if (!tmp) break;
msgs = tmp;
}
uint16_t plen16 = (uint16_t)plen;
msgs[msgs_len] = dir; msgs_len += 1;
memcpy(msgs + msgs_len, &seq, 8); msgs_len += 8;
memcpy(msgs + msgs_len, &ts, 8); msgs_len += 8;
memcpy(msgs + msgs_len, &author, 8); msgs_len += 8;
msgs[msgs_len] = (uint8_t)ct_len; msgs_len += 1;
if (ct_len) { memcpy(msgs + msgs_len, ct, ct_len); msgs_len += ct_len; }
memcpy(msgs + msgs_len, &plen16, 2); msgs_len += 2;
if (plen) { memcpy(msgs + msgs_len, plain, plen); msgs_len += plen; }
count++;
}
sqlite3_finalize(st);
}
size_t cl = strlen(conv_id);
size_t evt_sz = 1 + cl + 2 + msgs_len;
uint8_t* evt = u_malloc(evt_sz);
if (!evt) { if (msgs) u_free(msgs); return; }
uint16_t u16cnt = (uint16_t)count;
char conv[64];
snprintf(conv, sizeof(conv), "%s", req->conv_id);
u_free(req);
if (!dm_of(inst)) { DEBUG_WARN(DEBUG_CATEGORY_DM, "%s: snapshot after stop", DM_ID); return; }
size_t cl = strlen(conv), cap = 1024 * 1024, len = 0;
uint8_t* evt = u_malloc(1 + cl + cap);
if (!evt) { DEBUG_ERROR(DEBUG_CATEGORY_DM, "%s: snapshot allocation failed", DM_ID); return; }
evt[0] = (uint8_t)cl;
memcpy(evt + 1, conv_id, cl);
memcpy(evt + 1 + cl, &u16cnt, 2);
if (msgs_len) memcpy(evt + 1 + cl + 2, msgs, msgs_len);
chat_event_post(inst, CHAT_EVT_DM_MESSAGES, evt, (int)evt_sz);
memcpy(evt + 1, conv, cl);
if (!dm_list_messages_json(inst, conv, 200, 0, (char*)evt + 1 + cl, cap, &len))
chat_event_post(inst, CHAT_EVT_DM_MESSAGES, evt, (int)(1 + cl + len));
u_free(evt);
if (msgs) u_free(msgs);
}

4
src/dm/dm_core.h

@ -113,8 +113,8 @@ struct dm_send_req {
};
void dm_send_trampoline(void* arg);
/* Запросить расшифрованные сообщения беседы. Результат — CHAT_EVT_DM_MESSAGES:
* [conv_id_len:1][conv_id:var][count:2]([dir:1][seq:8][ts:8][author:8][ct_len:1][ct:var][data_len:2][data:var])* */
/* Запросить последние 200 сообщений со статусами. CHAT_EVT_DM_MESSAGES:
* [conv_id_len:1][conv_id:var][JSON:utf8]. Все чтения в uasync. */
struct dm_messages_req {
struct UTUN_INSTANCE* inst;
char conv_id[64];

Loading…
Cancel
Save