Browse Source

Store signed replies by message ID and confirm channel submissions

master
evgeny 4 days ago
parent
commit
68037b25a5
  1. 127
      src/chat/chat_msg.c
  2. 39
      src/chat/chat_whisper.c
  3. 4
      src/chat/chat_whisper.h

127
src/chat/chat_msg.c

@ -31,7 +31,7 @@ chat_whisper_trigger_fn g_chat_whisper_trigger = NULL;
/* forward decl */
void chat_core_submit_media_message(struct UTUN_INSTANCE* inst, struct chat_msg_submit* req);
static void wh_transcribe_done_cb(void* arg, const char* ch_id, const char* text, int err,
uint64_t reply_ts, uint64_t reply_node);
const uint8_t reply_id[32]);
static void chat_media_cache_cleanup_orphans(struct UTUN_INSTANCE* inst);
static void chat_media_cache_enforce(struct UTUN_INSTANCE* inst);
static void chat_media_cache_reconcile(struct UTUN_INSTANCE* inst);
@ -73,6 +73,26 @@ static int msg_get_prev_chain_hash(struct chat_core_ctx* cc, const char* tbl, ui
sqlite3_finalize(st); return 0;
}
/* Локальный ответ ссылается только на существующий оригинал своего канала; wire хранит общий ID. */
static int msg_reply_hex(struct chat_core_ctx* cc, const char* channel, const uint8_t reply[32], char hex[65]) {
static const uint8_t empty[32] = {0};
hex[0] = 0;
if (!memcmp(reply, empty, 32)) return 0;
for (unsigned i = 0; i < 32; i++) snprintf(hex + i * 2, 3, "%02x", reply[i]);
char table[80], sql[160]; msg_table_name(channel, table, sizeof(table));
snprintf(sql, sizeof(sql), "SELECT 1 FROM \"%s\" WHERE message_id=?", table);
sqlite3_stmt* st = NULL;
int rc = sqlite3_prepare_v2(cc->db, sql, -1, &st, NULL);
if (rc == SQLITE_OK) { sqlite3_bind_blob(st, 1, reply, 32, SQLITE_STATIC); rc = sqlite3_step(st); }
sqlite3_finalize(st);
if (rc != SQLITE_ROW) {
DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "reply: original unavailable ch=%s id=%s rc=%d: %s",
channel, hex, rc, sqlite3_errmsg(cc->db)); return -1;
}
DEBUG_DEBUG(DEBUG_CATEGORY_CHAT_SYNC, "reply: original resolved ch=%s id=%s", channel, hex);
return 0;
}
int chat_core_submit_message(struct UTUN_INSTANCE* inst, struct chat_msg_submit* req) {
struct chat_core_ctx* cc = inst ? CC(inst) : NULL;
if (!cc || !cc->initialized || !req || (req->data_len && !req->data) || req->data_len > 4000) {
@ -83,8 +103,10 @@ int chat_core_submit_message(struct UTUN_INSTANCE* inst, struct chat_msg_submit*
if (!si) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: no history for ch=%s", CC_ID, req->channel_id); return -1; }
/* SQLite формирует JSON без обрезки и с экранированием текста/Unicode. */
const char* query = req->reply_to_ts
? "SELECT json_object('n',?1,'ch',?2,'ct',?3,'d',?4,'rt',?5,'rn',?6)"
char reply_hex[65];
if (msg_reply_hex(cc, req->channel_id, req->reply_to, reply_hex) < 0) return -1;
const char* query = reply_hex[0]
? "SELECT json_object('n',?1,'ch',?2,'ct',?3,'d',?4,'reply_to',?5)"
: "SELECT json_object('n',?1,'ch',?2,'ct',?3,'d',?4)";
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(cc->db, query, -1, &st, NULL) != SQLITE_OK) {
@ -94,10 +116,7 @@ int chat_core_submit_message(struct UTUN_INSTANCE* inst, struct chat_msg_submit*
sqlite3_bind_text(st, 2, req->channel_id, -1, SQLITE_STATIC);
sqlite3_bind_text(st, 3, req->content_type, -1, SQLITE_STATIC);
sqlite3_bind_text(st, 4, req->data ? (const char*)req->data : "", req->data_len, SQLITE_STATIC);
if (req->reply_to_ts) {
sqlite3_bind_int64(st, 5, (sqlite3_int64)req->reply_to_ts);
sqlite3_bind_int64(st, 6, (sqlite3_int64)req->reply_to_node_id);
}
if (reply_hex[0]) sqlite3_bind_text(st, 5, reply_hex, -1, SQLITE_STATIC);
int rc = sqlite3_step(st), bytes = rc == SQLITE_ROW ? sqlite3_column_bytes(st, 0) : 0;
char json[4001];
if (rc != SQLITE_ROW || bytes <= 0 || bytes >= (int)sizeof(json)) {
@ -127,8 +146,10 @@ void chat_core_submit_trampoline(void* arg) {
if (!req) return;
if (req->media_src[0] != '\0')
chat_core_submit_media_message(req->inst, req);
else
chat_core_submit_message(req->inst, req);
else {
int result = chat_core_submit_message(req->inst, req);
chat_event_submitted(req->inst, req->request_id, result == 0);
}
u_free(req);
}
@ -145,8 +166,16 @@ struct media_submit_ctx {
int64_t duration_ms; /* видео: длительность после транскода */
int width; /* видео: размеры после транскода */
int height;
char reply_hex[65];
uint64_t request_id;
};
/* Завершение владеющей асинхронной отправки, включая все ошибки подготовки/регистрации. */
static void media_submit_finish(struct media_submit_ctx* ctx, int success) {
chat_event_submitted(ctx->inst, ctx->request_id, success);
u_free(ctx);
}
static void on_media_registered(void* arg, int err, const struct media_index_result* result) {
struct media_submit_ctx* mctx = (struct media_submit_ctx*)arg;
struct chat_core_ctx* cc = CC(mctx->inst);
@ -154,14 +183,14 @@ static void on_media_registered(void* arg, int err, const struct media_index_res
if (err || !result) {
DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: media registration failed ch=%s err=%d",
CC_ID, mctx->channel_id, err);
u_free(mctx);
media_submit_finish(mctx, 0);
return;
}
struct DB_SYNC_INSTANCE* si = si_find(mctx->inst, mctx->channel_id);
if (!si) {
DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: no db_sync instance for ch=%s", CC_ID, mctx->channel_id);
u_free(mctx);
media_submit_finish(mctx, 0);
return;
}
@ -172,7 +201,7 @@ static void on_media_registered(void* arg, int err, const struct media_index_res
size_t extra = 128 + 32 + 64 + pairs_len + 256;
size_t fdcap = base_len + extra;
char* full_data = u_malloc(fdcap);
if (!full_data) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: data alloc fail", CC_ID); u_free(mctx); return; }
if (!full_data) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: data alloc fail", CC_ID); media_submit_finish(mctx, 0); return; }
size_t off = 0;
if (base_len > 0) {
@ -206,24 +235,28 @@ static void on_media_registered(void* arg, int err, const struct media_index_res
}
size_t full_data_len = off;
size_t json_cap = full_data_len + 256;
size_t json_cap = full_data_len + 384;
char* json = u_malloc(json_cap);
if (!json) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: json alloc fail", CC_ID); u_free(full_data); u_free(mctx); return; }
if (!json) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: json alloc fail", CC_ID); u_free(full_data); media_submit_finish(mctx, 0); return; }
snprintf(json, json_cap,
"{\"n\":%llu,\"ch\":\"%s\",\"ct\":\"%s\",\"d\":\"%.*s\"}",
"{\"n\":%llu,\"ch\":\"%s\",\"ct\":\"%s\",\"d\":\"%.*s\"%s%s%s}",
(unsigned long long)cc->my_node_id, mctx->channel_id,
mctx->content_type, (int)full_data_len, full_data);
mctx->content_type, (int)full_data_len, full_data,
mctx->reply_hex[0] ? ",\"reply_to\":\"" : "", mctx->reply_hex, mctx->reply_hex[0] ? "\"" : "");
uint64_t ts = db_sync_next_timestamp(si);
size_t sig_msg_len = 8 + strlen(json);
uint8_t* sig_msg = u_malloc(sig_msg_len);
if (!sig_msg) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: sig_msg alloc fail", CC_ID); u_free(json); u_free(full_data); u_free(mctx); return; }
if (!sig_msg) {
DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: sig_msg alloc fail", CC_ID);
u_free(json); u_free(full_data); media_submit_finish(mctx, 0); return;
}
memcpy(sig_msg, &ts, 8);
memcpy(sig_msg + 8, json, sig_msg_len - 8);
uint8_t json_sig[64];
if (sc_ed25519_sign(mctx->inst->my_ed25519_privkey, sig_msg, sig_msg_len, json_sig) != SC_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: Ed25519 sign JSON failed", CC_ID);
u_free(sig_msg); u_free(json); u_free(full_data); u_free(mctx); return;
u_free(sig_msg); u_free(json); u_free(full_data); media_submit_finish(mctx, 0); return;
}
/* local_attrs готовим заранее и пишем атомарно при вставке — иначе GUI по
@ -251,14 +284,16 @@ static void on_media_registered(void* arg, int err, const struct media_index_res
/* whisper транскрипция для своих голосовых сообщений */
if (mctx->content_type[0] && strncmp(mctx->content_type, "voice", 5) == 0
&& mctx->inst && mctx->inst->media_async && g_chat_whisper_trigger) {
g_chat_whisper_trigger(mctx->inst, mctx->inst->media_async, mctx->inst->ua,
uint8_t original[32];
if (chat_message_id(db_sync_instance_group_id(si), ts, json_sig, original) == 0)
g_chat_whisper_trigger(mctx->inst, mctx->inst->media_async, mctx->inst->ua,
mctx->media_dest, mctx->channel_id,
ts, cc->my_node_id,
original,
wh_transcribe_done_cb, mctx->inst);
}
}
u_free(sig_msg); u_free(json); u_free(full_data); u_free(mctx);
u_free(sig_msg); u_free(json); u_free(full_data); media_submit_finish(mctx, ret == 0);
}
static void on_video_prepared(void* arg, int err, const struct video_info* vi) {
@ -268,7 +303,7 @@ static void on_video_prepared(void* arg, int err, const struct video_info* vi) {
if (err) {
DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: video prepare failed ch=%s err=%d",
CC_ID, mctx->channel_id, err);
u_free(mctx);
media_submit_finish(mctx, 0);
return;
}
@ -300,9 +335,16 @@ static void on_video_prepared(void* arg, int err, const struct video_info* vi) {
}
void chat_core_submit_media_message(struct UTUN_INSTANCE* inst, struct chat_msg_submit* req) {
struct chat_core_ctx* cc = CC(inst);
if (!cc || !cc->initialized) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: media submit before init", CC_ID); return; }
if (!req || req->media_src[0] == '\0') return;
struct chat_core_ctx* cc = inst ? CC(inst) : NULL;
if (!cc || !cc->initialized || !req || !req->media_src[0]) {
DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: invalid media submit/stopped core", CC_ID);
if (req) chat_event_submitted(inst, req->request_id, 0);
return;
}
char reply_hex[65];
if (msg_reply_hex(cc, req->channel_id, req->reply_to, reply_hex) < 0) {
chat_event_submitted(inst, req->request_id, 0); return;
}
/* лимит одного медиафайла применяется и к локальным отправкам */
uint64_t unit = inst->config->global.chatserver_storage_unit_size;
@ -311,6 +353,7 @@ void chat_core_submit_media_message(struct UTUN_INSTANCE* inst, struct chat_msg_
if (sz > 0 && (uint64_t)sz > unit) {
DEBUG_WARN(DEBUG_CATEGORY_CHAT_SYNC, "%s: media too large for local send (%lld > %llu bytes), rejected ch=%s",
CC_ID, (long long)sz, (unsigned long long)unit, req->channel_id);
chat_event_submitted(inst, req->request_id, 0);
return;
}
}
@ -321,6 +364,7 @@ void chat_core_submit_media_message(struct UTUN_INSTANCE* inst, struct chat_msg_
if (!inst->media_async) {
DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: media_async not initialized", CC_ID);
chat_event_submitted(inst, req->request_id, 0);
return;
}
@ -336,8 +380,13 @@ void chat_core_submit_media_message(struct UTUN_INSTANCE* inst, struct chat_msg_
/* copy fields for async callback (req is freed by trampoline immediately) */
struct media_submit_ctx* mctx = u_calloc(1, sizeof(*mctx));
if (!mctx) return;
if (!mctx) {
DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: media context allocation failed", CC_ID);
chat_event_submitted(inst, req->request_id, 0); return;
}
mctx->inst = inst;
mctx->request_id = req->request_id;
memcpy(mctx->reply_hex, reply_hex, strlen(reply_hex) + 1);
snprintf(mctx->channel_id, sizeof(mctx->channel_id), "%s", req->channel_id);
snprintf(mctx->content_type, sizeof(mctx->content_type), "%s", req->content_type);
snprintf(mctx->media_dest, sizeof(mctx->media_dest), "%s", req->media_dest);
@ -639,9 +688,11 @@ static void md_download_done_cb(void* arg, int err) {
/* whisper транскрипция для голосовых сообщений */
if (ctx->content_type[0] && strncmp(ctx->content_type, "voice", 5) == 0
&& ctx->inst && ctx->inst->media_async && g_chat_whisper_trigger) {
g_chat_whisper_trigger(ctx->inst, ctx->inst->media_async, ctx->inst->ua,
uint8_t original[32];
if (chat_message_id(strtoull(ctx->channel_id, NULL, 10), ctx->ts, ctx->author_sig, original) == 0)
g_chat_whisper_trigger(ctx->inst, ctx->inst->media_async, ctx->inst->ua,
ctx->dest_abs_path, ctx->channel_id,
ctx->ts, ctx->author_node_id,
original,
wh_transcribe_done_cb, ctx->inst);
}
@ -663,7 +714,7 @@ static void md_download_done_cb(void* arg, int err) {
}
static void wh_transcribe_done_cb(void* arg, const char* ch_id, const char* text, int err,
uint64_t reply_ts, uint64_t reply_node) {
const uint8_t reply_id[32]) {
struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)arg;
if (err || !text || !text[0]) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: transcription failed ch=%s err=%d", CC_ID, ch_id, err);
@ -672,17 +723,27 @@ static void wh_transcribe_done_cb(void* arg, const char* ch_id, const char* text
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: transcription result ch=%s: [%s]", CC_ID, ch_id, text);
/* mark original voice message as played */
chat_core_mark_voice_played(inst, ch_id, reply_ts, reply_node);
char table[80], sql[192]; msg_table_name(ch_id, table, sizeof(table));
snprintf(sql, sizeof(sql), "SELECT timestamp,node_id FROM \"%s\" WHERE message_id=?", table);
sqlite3_stmt* st = NULL;
int rc = sqlite3_prepare_v2(CC(inst)->db, sql, -1, &st, NULL);
if (rc == SQLITE_OK) { sqlite3_bind_blob(st, 1, reply_id, 32, SQLITE_STATIC); rc = sqlite3_step(st); }
if (rc != SQLITE_ROW) {
DEBUG_ERROR(DEBUG_CATEGORY_CHAT, "transcription: original unavailable ch=%s rc=%d", ch_id, rc);
sqlite3_finalize(st); return;
}
uint64_t ts = (uint64_t)sqlite3_column_int64(st, 0), author = (uint64_t)sqlite3_column_int64(st, 1);
sqlite3_finalize(st);
chat_core_mark_voice_played(inst, ch_id, ts, author);
struct chat_msg_submit* req = u_calloc(1, sizeof(*req));
if (!req) return;
if (!req) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT, "transcription: submit allocation failed"); return; }
req->inst = inst;
snprintf(req->channel_id, sizeof(req->channel_id), "%s", ch_id);
snprintf(req->content_type, sizeof(req->content_type), "voice_transcription");
req->data = (uint8_t*)text;
req->data_len = (uint32_t)strlen(text);
req->reply_to_ts = reply_ts;
req->reply_to_node_id = reply_node;
memcpy(req->reply_to, reply_id, 32);
chat_core_submit_message(inst, req);
u_free(req);
}

39
src/chat/chat_whisper.c

@ -41,8 +41,7 @@ struct wh_job {
struct UTUN_INSTANCE* inst;
char audio_path[1024];
char channel_id[64];
uint64_t reply_to_ts;
uint64_t reply_to_node;
uint8_t reply_id[32];
chat_whisper_done_fn done_cb;
void* done_arg;
struct media_async* ma;
@ -235,14 +234,14 @@ static void chat_whisper_transcribe_async_impl(
struct UTUN_INSTANCE* inst,
struct media_async* ma, struct UASYNC* ua,
const char* audio_path, const char* channel_id,
uint64_t reply_to_ts, uint64_t reply_to_node,
const uint8_t reply_id[32],
chat_whisper_done_fn done_cb, void* done_arg);
static void chat_whisper_transcribe_async(
struct UTUN_INSTANCE* inst,
struct media_async* ma, struct UASYNC* ua,
const char* audio_path, const char* channel_id,
uint64_t reply_to_ts, uint64_t reply_to_node,
const uint8_t reply_id[32],
chat_whisper_done_fn done_cb, void* done_arg);
struct wh_work_ctx {
@ -353,7 +352,7 @@ static void wh_done_fn(void* raw, int err) {
if (w->job.done_cb) {
w->job.done_cb(w->job.done_arg, w->job.channel_id,
w->text, w->err,
w->job.reply_to_ts, w->job.reply_to_node);
w->job.reply_id);
}
u_free(w->text);
@ -364,7 +363,7 @@ static void wh_done_fn(void* raw, int err) {
struct wh_job* nj = g_pending_job;
g_pending_job = NULL;
chat_whisper_transcribe_async_impl(nj->inst, nj->ma, nj->ua, nj->audio_path, nj->channel_id,
nj->reply_to_ts, nj->reply_to_node,
nj->reply_id,
nj->done_cb, nj->done_arg);
u_free(nj);
}
@ -461,11 +460,11 @@ void chat_whisper_transcribe_async(
struct UTUN_INSTANCE* inst,
struct media_async* ma, struct UASYNC* ua,
const char* audio_path, const char* channel_id,
uint64_t reply_to_ts, uint64_t reply_to_node,
const uint8_t reply_id[32],
chat_whisper_done_fn done_cb, void* done_arg)
{
chat_whisper_transcribe_async_impl(inst, ma, ua, audio_path, channel_id,
reply_to_ts, reply_to_node,
reply_id,
done_cb, done_arg);
}
@ -473,12 +472,12 @@ static void chat_whisper_transcribe_async_impl(
struct UTUN_INSTANCE* inst,
struct media_async* ma, struct UASYNC* ua,
const char* audio_path, const char* channel_id,
uint64_t reply_to_ts, uint64_t reply_to_node,
const uint8_t reply_id[32],
chat_whisper_done_fn done_cb, void* done_arg)
{
if (!audio_path || !channel_id || !done_cb || !ma || !ua) {
if (!audio_path || !channel_id || !reply_id || !done_cb || !ma || !ua) {
DEBUG_ERROR(DEBUG_CATEGORY_CHAT, "%s: invalid args", CW_ID);
if (done_cb) done_cb(done_arg, channel_id ? channel_id : "", NULL, -1, reply_to_ts, reply_to_node);
if (done_cb) done_cb(done_arg, channel_id ? channel_id : "", NULL, -1, reply_id);
return;
}
@ -486,7 +485,7 @@ static void chat_whisper_transcribe_async_impl(
int rc = chat_whisper_init(inst);
if (rc != 0) {
DEBUG_WARN(DEBUG_CATEGORY_CHAT, "%s: whisper not available, skip transcription", CW_ID);
done_cb(done_arg, channel_id, NULL, -1, reply_to_ts, reply_to_node);
done_cb(done_arg, channel_id, NULL, -1, reply_id);
return;
}
}
@ -495,12 +494,14 @@ static void chat_whisper_transcribe_async_impl(
/* поставить в очередь */
if (g_pending_job) { u_free(g_pending_job); } /* заменяем */
g_pending_job = u_malloc(sizeof(*g_pending_job));
if (!g_pending_job) { done_cb(done_arg, channel_id, NULL, -1, reply_to_ts, reply_to_node); return; }
if (!g_pending_job) {
DEBUG_ERROR(DEBUG_CATEGORY_CHAT, "%s: pending job allocation failed", CW_ID);
done_cb(done_arg, channel_id, NULL, -1, reply_id); return;
}
g_pending_job->inst = inst;
snprintf(g_pending_job->audio_path, sizeof(g_pending_job->audio_path), "%s", audio_path);
snprintf(g_pending_job->channel_id, sizeof(g_pending_job->channel_id), "%s", channel_id);
g_pending_job->reply_to_ts = reply_to_ts;
g_pending_job->reply_to_node = reply_to_node;
memcpy(g_pending_job->reply_id, reply_id, 32);
g_pending_job->done_cb = done_cb;
g_pending_job->done_arg = done_arg;
g_pending_job->ma = ma;
@ -512,12 +513,14 @@ static void chat_whisper_transcribe_async_impl(
g_processing = 1;
struct wh_work_ctx* w = u_calloc(1, sizeof(*w));
if (!w) { g_processing = 0; done_cb(done_arg, channel_id, NULL, -1, reply_to_ts, reply_to_node); return; }
if (!w) {
DEBUG_ERROR(DEBUG_CATEGORY_CHAT, "%s: worker allocation failed", CW_ID);
g_processing = 0; done_cb(done_arg, channel_id, NULL, -1, reply_id); return;
}
snprintf(w->job.audio_path, sizeof(w->job.audio_path), "%s", audio_path);
snprintf(w->job.channel_id, sizeof(w->job.channel_id), "%s", channel_id);
w->job.reply_to_ts = reply_to_ts;
w->job.reply_to_node = reply_to_node;
memcpy(w->job.reply_id, reply_id, 32);
w->job.done_cb = done_cb;
w->job.done_arg = done_arg;
w->job.ma = ma;

4
src/chat/chat_whisper.h

@ -36,13 +36,13 @@ void chat_whisper_destroy(void); /* выгрузка модел
typedef void (*chat_whisper_done_fn)(void* arg, const char* channel_id,
const char* text, int err,
uint64_t reply_ts, uint64_t reply_node);
const uint8_t reply_id[32]);
typedef void (*chat_whisper_trigger_fn)(
struct UTUN_INSTANCE* inst,
struct media_async* ma, struct UASYNC* ua,
const char* audio_path, const char* channel_id,
uint64_t reply_ts, uint64_t reply_node,
const uint8_t reply_id[32],
chat_whisper_done_fn done_cb, void* done_arg);
/* Глобальный триггер: устанавливается chat_whisper_init(), вызывается из chat_msg.c */

Loading…
Cancel
Save