/* * chat_msg.c — сообщения: отправка, DB-операции для chat_sync * * Вынесено из chat_core.c для уменьшения размера модуля. */ #include "chat_core_priv.h" #include "chat_event.h" #include "chat_setting.h" #include "chat_whisper.h" #include "../utun_instance.h" #include "../ntp_time.h" #include "../transport_layer/secure_channel.h" #include "../media_delivery/media_index.h" #include "../media_delivery/media_download.h" #include "../media_delivery/media_delivery.h" #include "../media_async/media_async.h" #include "../video/video.h" #include "../../lib/mem.h" #include "../../lib/u_async.h" #include "../../lib/platform_compat.h" #include #include 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); 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); /* ─── счётчик объёма медиакеша (сумма файлов, соответствующих сообщениям) ─── */ static void cache_bytes_add(struct chat_core_ctx* cc, int64_t sz) { if (sz > 0) cc->media_cache_bytes += (uint64_t)sz; } static void cache_bytes_sub(struct chat_core_ctx* cc, int64_t sz) { uint64_t v = sz > 0 ? (uint64_t)sz : 0; cc->media_cache_bytes = cc->media_cache_bytes >= v ? cc->media_cache_bytes - v : 0; } /* Извлекает имя файла из local_attrs: {"st":"fl","fp":""} */ static void extract_attrs_fp(const char* local_attrs, char* out, size_t out_sz) { out[0] = '\0'; if (!local_attrs) return; const char* p = strstr(local_attrs, "\"fp\":\""); if (!p) return; p += 6; int i = 0; while (p[i] && p[i] != '"' && i < (int)out_sz - 1) out[i++] = p[i]; out[i] = '\0'; } /* ─── отправка сообщения (GUI → uasync) ─── */ static int msg_get_prev_chain_hash(struct chat_core_ctx* cc, const char* tbl, uint8_t out[32]) { char sql[128]; snprintf(sql, sizeof(sql), "SELECT chain_hash FROM \"%s\" ORDER BY timestamp DESC LIMIT 1", tbl); sqlite3_stmt* st = NULL; if (sqlite3_prepare_v2(cc->db, sql, -1, &st, NULL) != SQLITE_OK) { memset(out,0,32); return -1; } if (sqlite3_step(st) == SQLITE_ROW) { const void* blob = sqlite3_column_blob(st, 0); int blen = sqlite3_column_bytes(st, 0); if (blob && blen == 32) memcpy(out, blob, 32); else memset(out, 0, 32); } else { memset(out, 0, 32); } sqlite3_finalize(st); return 0; } void chat_core_submit_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: submit before init", CC_ID); return; } if (!req) return; DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "%s: SUBMIT ch=%s ct=%s len=%u", CC_ID, req->channel_id, req->content_type, req->data_len); struct DB_SYNC_INSTANCE* si = si_find(inst, req->channel_id); if (!si) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: no db_sync instance for ch=%s", CC_ID, req->channel_id); return; } char json[4096]; int off = snprintf(json, sizeof(json), "{\"n\":%llu,\"ch\":\"%s\",\"ct\":\"%s\"", (unsigned long long)cc->my_node_id, req->channel_id, req->content_type); if (req->reply_to_ts) { off += snprintf(json + off, sizeof(json) - off, ",\"rt\":%llu,\"rn\":%llu", (unsigned long long)req->reply_to_ts, (unsigned long long)req->reply_to_node_id); } off += snprintf(json + off, sizeof(json) - off, ",\"d\":\"%.*s\"}", (int)req->data_len, (const char*)req->data); uint64_t ts = db_sync_next_timestamp(si); uint8_t sig_msg[8192]; size_t soff = 0; memcpy(sig_msg + soff, &ts, 8); soff += 8; size_t jl = strlen(json); memcpy(sig_msg + soff, json, jl); soff += jl; uint8_t sig[64]; if (sc_ed25519_sign(inst->my_ed25519_privkey, sig_msg, soff, sig) != SC_OK) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: Ed25519 sign failed", CC_ID); return; } int ret = db_sync_insert_signed(si, json, strlen(json), sig, 64, ts, NULL); if (ret != 0) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: db_sync_insert_signed failed ret=%d", CC_ID, ret); return; } } void chat_core_submit_trampoline(void* arg) { struct chat_msg_submit* req = (struct chat_msg_submit*)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); u_free(req); } /* ─── Медиа-сообщения: async регистрация в media_files + отправка ─── */ struct media_submit_ctx { struct UTUN_INSTANCE* inst; char channel_id[64]; char content_type[32]; char media_dest[1024]; uint8_t msg_data[4096]; uint32_t msg_data_len; uint64_t timestamp; int64_t duration_ms; /* видео: длительность после транскода */ int width; /* видео: размеры после транскода */ int height; }; 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); 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); 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); return; } /* build full_data: |||||| */ size_t base_len = mctx->msg_data_len; int nb = result->num_blocks; size_t pairs_len = (size_t)nb * (32 + 1 + 128 + 1); /* id_hex, sig_hex, commas */ 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; } size_t off = 0; if (base_len > 0) { char b64[2048]; size_t b64len = b64_encode(mctx->msg_data, base_len, b64, sizeof(b64)); full_data[off++] = '*'; memcpy(full_data + off, b64, b64len); off += b64len; } off += snprintf(full_data + off, fdcap - off, "|%lld|%lld|%d|", (long long)result->file_size, (long long)result->block_size, nb); for (int i = 0; i < 16; i++) off += snprintf(full_data + off, fdcap - off, "%02x", result->media_id[i]); off += snprintf(full_data + off, fdcap - off, "|"); for (int i = 0; i < 32; i++) off += snprintf(full_data + off, fdcap - off, "%02x", result->content_hash[i]); off += snprintf(full_data + off, fdcap - off, "|"); for (int n = 0; n < nb; n++) { if (n > 0) full_data[off++] = ','; for (int i = 0; i < 16; i++) off += snprintf(full_data + off, fdcap - off, "%02x", result->block_ids[n * 16 + i]); full_data[off++] = ','; for (int i = 0; i < 64; i++) off += snprintf(full_data + off, fdcap - off, "%02x", result->block_sigs[n * 64 + i]); } /* для видео дописываем длительность/размеры в конец (хвост игнорируется старыми парсерами) */ if (mctx->content_type[0] && strcmp(mctx->content_type, "video/mp4") == 0 && mctx->duration_ms > 0) { off += snprintf(full_data + off, fdcap - off, "|dur=%lld|w=%d|h=%d", (long long)mctx->duration_ms, mctx->width, mctx->height); } /* для фото дописываем разрешение в конец */ else if (mctx->content_type[0] && strncmp(mctx->content_type, "image/", 6) == 0 && (mctx->width > 0 || mctx->height > 0)) { off += snprintf(full_data + off, fdcap - off, "|w=%d|h=%d", mctx->width, mctx->height); } size_t full_data_len = off; size_t json_cap = full_data_len + 256; 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; } snprintf(json, json_cap, "{\"n\":%llu,\"ch\":\"%s\",\"ct\":\"%s\",\"d\":\"%.*s\"}", (unsigned long long)cc->my_node_id, mctx->channel_id, mctx->content_type, (int)full_data_len, full_data); 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; } 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; } int ret = db_sync_insert_signed(si, json, strlen(json), json_sig, 64, ts, NULL); if (ret != 0) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: media insert failed ret=%d", CC_ID, ret); } else { DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "%s: media message inserted ch=%s ts=%llu", CC_ID, mctx->channel_id, (unsigned long long)ts); const char* base = strrchr(mctx->media_dest, '/'); const char* filename = base ? base + 1 : mctx->media_dest; char attrs[1100]; snprintf(attrs, sizeof(attrs), "{\"st\":\"fl\",\"fp\":\"%s\"}", filename); chat_core_update_local_attrs(mctx->inst, mctx->channel_id, ts, json_sig, attrs); /* анонсируем блоки суперузлам — иначе первый скачивающий не найдёт держателя */ media_delivery_announce_media(mctx->inst, strtoull(mctx->channel_id, NULL, 10), result->media_id, result->block_ids, result->num_blocks); /* учесть локально отправленный файл в счётчике и применить лимит медиакеша */ cache_bytes_add(cc, result->file_size); chat_media_cache_enforce(mctx->inst); /* 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, mctx->media_dest, mctx->channel_id, ts, cc->my_node_id, wh_transcribe_done_cb, mctx->inst); } } u_free(sig_msg); u_free(json); u_free(full_data); u_free(mctx); } static void on_video_prepared(void* arg, int err) { struct media_submit_ctx* mctx = (struct media_submit_ctx*)arg; struct chat_core_ctx* cc = CC(mctx->inst); 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); return; } /* после транскода знаем итоговую длительность и размеры — для тела сообщения */ struct video_info vi; if (video_get_info(mctx->media_dest, &vi) == 0 && vi.is_video) { mctx->duration_ms = vi.duration_ms; mctx->width = vi.width; mctx->height = vi.height; DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "%s: video prepared %dx%d dur=%lldms", CC_ID, vi.width, vi.height, (long long)vi.duration_ms); } char media_base[1024]; { const char* last_slash = strrchr(cc->db_path, '/'); if (last_slash) snprintf(media_base, sizeof(media_base), "%.*s", (int)(last_slash - cc->db_path), cc->db_path); else snprintf(media_base, sizeof(media_base), "%s", cc->db_path); } media_index_register_async( mctx->inst->media_async, mctx->inst->ua, cc->db, cc->my_node_id, mctx->inst->my_ed25519_privkey, mctx->channel_id, mctx->media_dest, mctx->media_dest, 0, media_base, on_media_registered, mctx); } 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; /* лимит одного медиафайла применяется и к локальным отправкам */ uint64_t unit = inst->config->global.chatserver_storage_unit_size; if (unit > 0) { int64_t sz = (int64_t)ma_file_size(req->media_src); 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); return; } } DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "%s: MEDIA SUBMIT ch=%s ct=%s src=%s dst=%s copy=%d video=%d", CC_ID, req->channel_id, req->content_type, req->media_src, req->media_dest, req->media_copy, req->media_video); if (!inst->media_async) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: media_async not initialized", CC_ID); return; } /* build media_base from db_path: dirname(db_path) e.g. /path/to/data */ char media_base[1024]; { const char* last_slash = strrchr(cc->db_path, '/'); if (last_slash) snprintf(media_base, sizeof(media_base), "%.*s", (int)(last_slash - cc->db_path), cc->db_path); else snprintf(media_base, sizeof(media_base), "%s", cc->db_path); } /* copy fields for async callback (req is freed by trampoline immediately) */ struct media_submit_ctx* mctx = u_calloc(1, sizeof(*mctx)); if (!mctx) return; mctx->inst = inst; 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); mctx->timestamp = req->timestamp; mctx->duration_ms = req->duration_ms; mctx->width = req->width; mctx->height = req->height; if (req->data && req->data_len > 0) { mctx->msg_data_len = req->data_len < sizeof(mctx->msg_data) ? req->data_len : sizeof(mctx->msg_data) - 1; memcpy(mctx->msg_data, req->data, mctx->msg_data_len); } if (req->media_video) { video_transcode_start(inst->ua, req->media_src, req->media_dest, on_video_prepared, mctx); return; } media_index_register_async( inst->media_async, inst->ua, cc->db, cc->my_node_id, inst->my_ed25519_privkey, req->channel_id, req->media_src, req->media_dest, req->media_copy ? 1 : 0, media_base, on_media_registered, mctx); } /* ─── Обновление local_attrs (для будущей приёмной стороны) ─── */ int chat_core_update_local_attrs(struct UTUN_INSTANCE* inst, const char* ch_id, uint64_t ts, const uint8_t* author_sig, const char* local_attrs_json) { struct chat_core_ctx* cc = CC(inst); if (!cc || !cc->initialized || !ch_id || !author_sig || !local_attrs_json) return -1; char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl)); char sql[256]; snprintf(sql, sizeof(sql), "UPDATE \"%s\" SET local_attrs=? WHERE timestamp=? AND author_signature=?", tbl); sqlite3_stmt* stmt = NULL; if (sqlite3_prepare_v2(cc->db, sql, -1, &stmt, NULL) != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: update local_attrs prep fail: %s", CC_ID, sqlite3_errmsg(cc->db)); return -1; } sqlite3_bind_text(stmt, 1, local_attrs_json, -1, SQLITE_STATIC); sqlite3_bind_int64(stmt, 2, (sqlite3_int64)ts); sqlite3_bind_blob(stmt, 3, author_sig, DB_SIG_SIZE, SQLITE_STATIC); int rc = sqlite3_step(stmt); sqlite3_finalize(stmt); if (rc != SQLITE_DONE) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: update local_attrs step fail rc=%d", CC_ID, rc); return -1; } DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "%s: local_attrs updated ch=%s ts=%llu", CC_ID, ch_id, (unsigned long long)ts); return 0; } /* ─── Отметка голосового как проигранного ─── */ void chat_core_mark_voice_played(struct UTUN_INSTANCE* inst, const char* ch_id, uint64_t ts, uint64_t author_node_id) { struct chat_core_ctx* cc = CC(inst); if (!cc || !cc->initialized || !ch_id || ts == 0) return; char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl)); char sql[256]; snprintf(sql, sizeof(sql), "SELECT local_attrs, author_signature FROM \"%s\" WHERE timestamp=? AND node_id=? LIMIT 1", tbl); sqlite3_stmt* st = NULL; if (sqlite3_prepare_v2(cc->db, sql, -1, &st, NULL) != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_CHAT_SYNC, "%s: mark_voice_played prep fail: %s", CC_ID, sqlite3_errmsg(cc->db)); return; } sqlite3_bind_int64(st, 1, (sqlite3_int64)ts); sqlite3_bind_int64(st, 2, (sqlite3_int64)author_node_id); if (sqlite3_step(st) != SQLITE_ROW) { sqlite3_finalize(st); return; } const char* la = (const char*)sqlite3_column_text(st, 0); const uint8_t* sig = sqlite3_column_blob(st, 1); int sig_len = sqlite3_column_bytes(st, 1); if (!sig || sig_len != DB_SIG_SIZE) { sqlite3_finalize(st); return; } uint8_t author_sig[DB_SIG_SIZE]; memcpy(author_sig, sig, DB_SIG_SIZE); if (la && strstr(la, "\"pl\":1")) { sqlite3_finalize(st); return; } /* уже помечено */ /* вставляем "pl":1 перед последней '}' */ char new_attrs[256]; if (la && la[0]) { size_t la_len = strlen(la); if (la_len > 0 && la[la_len - 1] == '}') { snprintf(new_attrs, sizeof(new_attrs), "%.*s,\"pl\":1}", (int)(la_len - 1), la); } else { snprintf(new_attrs, sizeof(new_attrs), "%s,\"pl\":1}", la); } } else { snprintf(new_attrs, sizeof(new_attrs), "{\"pl\":1}"); } sqlite3_finalize(st); chat_core_update_local_attrs(inst, ch_id, ts, author_sig, new_attrs); DEBUG_INFO(DEBUG_CATEGORY_CHAT_SYNC, "%s: voice marked played ch=%s ts=%llu", CC_ID, ch_id, (unsigned long long)ts); } void chat_core_mark_voice_played_trampoline(void* arg) { struct voice_played_req* req = (struct voice_played_req*)arg; if (!req) return; chat_core_mark_voice_played(req->inst, req->channel_id, req->ts, req->author_node_id); u_free(req); } /* ─── DB-операции для chat_sync ─── */ uint32_t chat_core_count(struct UTUN_INSTANCE* inst, const char* ch_id) { struct chat_core_ctx* cc = CC(inst); if (!cc || !cc->initialized) return 0; char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl)); char sql[128]; snprintf(sql, sizeof(sql), "SELECT COUNT(*) FROM \"%s\"", tbl); sqlite3_stmt* stmt = NULL; if (sqlite3_prepare_v2(cc->db, sql, -1, &stmt, NULL) != SQLITE_OK) return 0; uint32_t cnt = 0; if (sqlite3_step(stmt) == SQLITE_ROW) cnt = (uint32_t)sqlite3_column_int64(stmt, 0); sqlite3_finalize(stmt); return cnt; } int chat_core_chain_hash_at(struct UTUN_INSTANCE* inst, const char* ch_id, uint32_t pos, uint8_t* hash_out) { struct chat_core_ctx* cc = CC(inst); if (!cc || !cc->initialized || !hash_out) return -1; char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl)); char sql[200]; snprintf(sql, sizeof(sql), "SELECT chain_hash FROM \"%s\" ORDER BY timestamp, node_id ASC LIMIT 1 OFFSET ?", tbl); sqlite3_stmt* stmt = NULL; if (sqlite3_prepare_v2(cc->db, sql, -1, &stmt, NULL) != SQLITE_OK) { memset(hash_out, 0, 32); return -1; } sqlite3_bind_int64(stmt, 1, (sqlite3_int64)pos); if (sqlite3_step(stmt) == SQLITE_ROW) { const void* blob = sqlite3_column_blob(stmt, 0); int len = sqlite3_column_bytes(stmt, 0); if (blob && len == 32) memcpy(hash_out, blob, 32); else memset(hash_out, 0, 32); } else { memset(hash_out, 0, 32); } sqlite3_finalize(stmt); return 0; } int chat_core_list_channels(struct UTUN_INSTANCE* inst, uint8_t* buf, size_t buf_size, size_t* out_len) { struct chat_core_ctx* cc = CC(inst); if (!cc || !cc->initialized || !buf || !out_len) return -1; sqlite3_stmt* stmt = NULL; if (sqlite3_prepare_v2(cc->db, "SELECT channel_id FROM channels ORDER BY created_at ASC", -1, &stmt, NULL) != SQLITE_OK) return -1; uint8_t* out = buf; uint8_t* start = out; out += 2; uint16_t cnt = 0; while (sqlite3_step(stmt) == SQLITE_ROW) { const char* ch_id = (const char*)sqlite3_column_text(stmt, 0); int ch_len = sqlite3_column_bytes(stmt, 0); if (!ch_id || ch_len <= 0 || ch_len > 63) continue; size_t need = (size_t)(out - buf) + 1 + (size_t)ch_len; if (need > buf_size) break; *out++ = (uint8_t)ch_len; memcpy(out, ch_id, (size_t)ch_len); out += ch_len; cnt++; } sqlite3_finalize(stmt); memcpy(start, &cnt, 2); *out_len = (size_t)(out - buf); return 0; } int chat_core_list_peers(struct UTUN_INSTANCE* inst, const char* ch_id, uint8_t* buf, size_t buf_size, size_t* out_len) { struct chat_core_ctx* cc = CC(inst); if (!cc || !cc->initialized || !buf || !out_len) return -1; char tbl[80]; peers_table_name(ch_id, tbl, sizeof(tbl)); char sql[128]; snprintf(sql, sizeof(sql), "SELECT node_id FROM \"%s\" WHERE deleted=0", tbl); sqlite3_stmt* stmt = NULL; if (sqlite3_prepare_v2(cc->db, sql, -1, &stmt, NULL) != SQLITE_OK) { uint16_t z = 0; memcpy(buf, &z, 2); *out_len = 2; return -1; } uint16_t cnt = 0; uint8_t* out = buf + 2; while (sqlite3_step(stmt) == SQLITE_ROW) { if ((size_t)(out - buf) + 8 > buf_size) break; uint64_t nid = (uint64_t)sqlite3_column_int64(stmt, 0); memcpy(out, &nid, 8); out += 8; cnt++; } sqlite3_finalize(stmt); memcpy(buf, &cnt, 2); *out_len = (size_t)(out - buf); return 0; } int chat_core_load_nodeinfo(struct UTUN_INSTANCE* inst, uint64_t node_id, uint8_t* buf, size_t buf_size, size_t* out_len) { struct chat_core_ctx* cc = CC(inst); if (!cc || !cc->initialized || !buf || !out_len) return -1; sqlite3_stmt* stmt = NULL; if (sqlite3_prepare_v2(cc->db, "SELECT x25519_pubkey, ed25519_pubkey FROM nodes WHERE node_id=?", -1, &stmt, NULL) != SQLITE_OK) return -1; sqlite3_bind_int64(stmt, 1, (sqlite3_int64)node_id); if (sqlite3_step(stmt) != SQLITE_ROW) { sqlite3_finalize(stmt); return -1; } const uint8_t* x25519 = (const uint8_t*)sqlite3_column_blob(stmt, 0); int x25519_len = sqlite3_column_bytes(stmt, 0); const uint8_t* ed25519 = (const uint8_t*)sqlite3_column_blob(stmt, 1); int ed25519_len = sqlite3_column_bytes(stmt, 1); uint8_t* out = buf; if (x25519 && x25519_len == 32) memcpy(out, x25519, 32); else memset(out, 0, 32); out += 32; if (ed25519 && ed25519_len == 32) memcpy(out, ed25519, 32); else memset(out, 0, 32); out += 32; sqlite3_finalize(stmt); if (sqlite3_prepare_v2(cc->db, "SELECT family, protocol, address, port, rtt FROM node_addresses WHERE node_id=?", -1, &stmt, NULL) != SQLITE_OK) { *out++ = 0; *out_len = (size_t)(out - buf); return 0; } sqlite3_bind_int64(stmt, 1, (sqlite3_int64)node_id); uint8_t* cnt_pos = out; *out++ = 0; uint8_t addr_cnt = 0; while (sqlite3_step(stmt) == SQLITE_ROW && addr_cnt < 255) { int family = sqlite3_column_int(stmt, 0); int proto = sqlite3_column_int(stmt, 1); const uint8_t* addr = (const uint8_t*)sqlite3_column_blob(stmt, 2); int addr_len = sqlite3_column_bytes(stmt, 2); uint16_t port = (uint16_t)sqlite3_column_int(stmt, 3); int16_t rtt = (int16_t)sqlite3_column_int(stmt, 4); if ((size_t)(out - buf) + 7 + (size_t)addr_len > buf_size) break; *out++ = (uint8_t)family; *out++ = (uint8_t)proto; *out++ = (uint8_t)addr_len; if (addr_len > 0) { memcpy(out, addr, (size_t)addr_len); out += addr_len; } memcpy(out, &port, 2); out += 2; memcpy(out, &rtt, 2); out += 2; addr_cnt++; } *cnt_pos = addr_cnt; sqlite3_finalize(stmt); *out_len = (size_t)(out - buf); return 0; } /* ─── download completion ─── */ struct md_done_ctx { struct UTUN_INSTANCE* inst; char channel_id[64]; uint64_t ts; int64_t msg_id; uint8_t author_sig[64]; char dest_relpath[512]; char dest_abs_path[1024]; char content_type[32]; uint64_t author_node_id; }; static void md_download_progress_cb(void* arg, int blocks_done, int num_blocks) { struct md_done_ctx* ctx = (struct md_done_ctx*)arg; uint8_t evt[256]; uint8_t cl = (uint8_t)strlen(ctx->channel_id); evt[0] = cl; memcpy(evt + 1, ctx->channel_id, cl); memcpy(evt + 1 + cl, &ctx->msg_id, 8); uint32_t bd = (uint32_t)blocks_done, nb = (uint32_t)num_blocks; memcpy(evt + 1 + cl + 8, &bd, 4); memcpy(evt + 1 + cl + 12, &nb, 4); chat_event_post(ctx->inst, CHAT_EVT_DOWNLOAD_PROGRESS, evt, 1 + cl + 16); } static void md_download_done_cb(void* arg, int err) { struct md_done_ctx* ctx = (struct md_done_ctx*)arg; struct chat_core_ctx* cc = CC(ctx->inst); if (!err) { char attrs[1024]; snprintf(attrs, sizeof(attrs), "{\"st\":\"fl\",\"fp\":\"%s\"}", ctx->dest_relpath); chat_core_update_local_attrs(ctx->inst, ctx->channel_id, ctx->ts, ctx->author_sig, attrs); DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: media downloaded ch=%s", CC_ID, ctx->channel_id); /* 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, ctx->dest_abs_path, ctx->channel_id, ctx->ts, ctx->author_node_id, wh_transcribe_done_cb, ctx->inst); } /* после докачки — учесть файл в счётчике и применить лимит медиакеша */ cache_bytes_add(cc, ma_file_size(ctx->dest_abs_path)); chat_media_cache_enforce(ctx->inst); } else { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: media download failed ch=%s err=%d", CC_ID, ctx->channel_id, err); /* пометить сообщение ошибкой, чтобы UI показал «повторить» */ char attrs[64]; snprintf(attrs, sizeof(attrs), "{\"st\":\"er\"}"); chat_core_update_local_attrs(ctx->inst, ctx->channel_id, ctx->ts, ctx->author_sig, attrs); } uint8_t evt[80]; uint8_t cl = (uint8_t)strlen(ctx->channel_id); evt[0] = cl; memcpy(evt + 1, ctx->channel_id, cl); memcpy(evt + 1 + cl, &ctx->msg_id, 8); chat_event_post(ctx->inst, CHAT_EVT_ATTACHMENT_DOWNLOADED, evt, 1 + cl + 8); u_free(ctx); } 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) { 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); return; } 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); struct chat_msg_submit* req = u_calloc(1, sizeof(*req)); if (!req) 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; chat_core_submit_message(inst, req); u_free(req); } static void sanitize_filename(char* buf, size_t size) { char orig[256]; snprintf(orig, sizeof(orig), "%s", buf); size_t w = 0; for (size_t r = 0; buf[r]; r++) { char c = buf[r]; if ((unsigned char)c < 0x20) continue; if (c == '/' || c == '\\' || c == ':' || c == '*' || c == '?' || c == '"' || c == '<' || c == '>' || c == '|') continue; if (w < size - 1) buf[w++] = c; } buf[w] = '\0'; if (w == 0) snprintf(buf, size, "file"); if (strcmp(orig, buf) != 0) DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: filename sanitized [%s] -> [%s]", CC_ID, orig, buf); } /* Извлекает имя файла из тела "d"-поля медиасообщения: берёт сегмент до '|', декодирует base64url-префикс '*', если он есть, и санитизирует. */ static void extract_base_filename(const char* body, char* out, size_t out_sz) { size_t bfi = 0; while (body[bfi] && body[bfi] != '|' && bfi < out_sz - 1) out[bfi++] = body[bfi]; out[bfi] = '\0'; if (out[0] == '*' && bfi > 1) { uint8_t dec[256]; int dlen = (int)b64_decode(out + 1, bfi - 1, dec, sizeof(dec)); if (dlen > 0 && dlen < (int)out_sz) { memcpy(out, dec, (size_t)dlen); out[dlen] = '\0'; } } if (bfi == 0) snprintf(out, out_sz, "file"); sanitize_filename(out, out_sz); } /* Разбор тела "d" медиасообщения в media_index_result. Формат: * ||||||,,... * Выделяет result->block_ids/block_sigs (освобождать через media_index_result_free). */ static int chat_msg_parse_media_body(const char* body, struct media_index_result* result) { memset(result, 0, sizeof(*result)); if (!body) return -1; const char* p = body; while (*p && *p != '|') p++; if (!*p) return -1; p++; long long fsize = strtoll(p, (char**)&p, 10); if (*p != '|' || fsize <= 0) return -1; p++; long long bsize = strtoll(p, (char**)&p, 10); if (*p != '|' || bsize <= 0) return -1; p++; int nb = (int)strtol(p, (char**)&p, 10); while (*p && *p != '|') p++; if (!*p || nb <= 0 || nb > 100) return -1; p++; result->file_size = fsize; result->block_size = bsize; result->num_blocks = nb; for (int i = 0; i < 16; i++) { char b[3] = { p[i*2], p[i*2+1], 0 }; result->media_id[i] = (uint8_t)strtol(b, NULL, 16); } p += 32; if (*p != '|') return -1; p++; for (int i = 0; i < 32; i++) { char b[3] = { p[i*2], p[i*2+1], 0 }; result->content_hash[i] = (uint8_t)strtol(b, NULL, 16); } p += 64; if (*p != '|') return -1; p++; result->block_ids = u_malloc((size_t)nb * 16); result->block_sigs = u_malloc((size_t)nb * 64); if (!result->block_ids || !result->block_sigs) { media_index_result_free(result); return -1; } for (int n = 0; n < nb; n++) { for (int i = 0; i < 16; i++) { char b[3] = { p[i*2], p[i*2+1], 0 }; result->block_ids[n * 16 + i] = (uint8_t)strtol(b, NULL, 16); } p += 32; if (*p != ',') { media_index_result_free(result); return -1; } p++; for (int i = 0; i < 64; i++) { char b[3] = { p[i*2], p[i*2+1], 0 }; result->block_sigs[n * 64 + i] = (uint8_t)strtol(b, NULL, 16); } p += 128; if (n < nb - 1 && *p == ',') p++; } return 0; } /* Имя файла на диске: голосовое — voice_.opus, иначе — базовое имя из "d". */ static void media_final_name(const struct media_index_result* result, const char* content_type, const char* base_filename, char* out, size_t out_sz) { if (content_type && strcmp(content_type, "audio/opus") == 0) { char hex[33]; for (int i = 0; i < 16; i++) snprintf(hex + i * 2, 3, "%02x", result->media_id[i]); snprintf(out, out_sz, "voice_%s.opus", hex); } else { snprintf(out, out_sz, "%s", base_filename); } } static int md_start_download(struct UTUN_INSTANCE* inst, const char* data_str, size_t data_len, const char* ch_id, const char* base_filename, uint64_t ts, const uint8_t* author_sig, uint64_t author_node_id, int64_t msg_id, const char* content_type) { (void)data_len; struct chat_core_ctx* cc = CC(inst); struct media_index_result result; if (chat_msg_parse_media_body(data_str, &result) != 0) return -1; int nb = result.num_blocks; long long fsize = result.file_size; /* для голосовых имя генерируем из media_id: msg_data — это wf-строка, а не имя файла */ char final_name[256]; media_final_name(&result, content_type, base_filename, final_name, sizeof(final_name)); char dest[1024], media_base[512]; { const char* last_slash = strrchr(cc->db_path, '/'); if (last_slash) { snprintf(media_base, sizeof(media_base), "%.*s", (int)(last_slash - cc->db_path), cc->db_path); } else { snprintf(media_base, sizeof(media_base), "%s", cc->db_path); } } snprintf(dest, sizeof(dest), "%s/media/%s/%s", media_base, ch_id, final_name); struct md_done_ctx* ctx = u_calloc(1, sizeof(*ctx)); if (!ctx) { media_index_result_free(&result); return -1; } ctx->inst = inst; snprintf(ctx->channel_id, sizeof(ctx->channel_id), "%s", ch_id); ctx->ts = ts; ctx->msg_id = msg_id; snprintf(ctx->dest_relpath, sizeof(ctx->dest_relpath), "%s", final_name); snprintf(ctx->dest_abs_path, sizeof(ctx->dest_abs_path), "%s", dest); ctx->author_node_id = author_node_id; if (content_type) snprintf(ctx->content_type, sizeof(ctx->content_type), "%s", content_type); else ctx->content_type[0] = '\0'; if (author_sig) memcpy(ctx->author_sig, author_sig, 64); DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: download start ch=%s ct=%s file=%s blocks=%d size=%lld dest=%s", CC_ID, ch_id, content_type ? content_type : "?", final_name, nb, (long long)fsize, dest); uint64_t gid = strtoull(ch_id, NULL, 10); media_download_start(inst, gid, &result, dest, media_base, author_node_id, md_download_done_cb, ctx, md_download_progress_cb, ctx); media_index_result_free(&result); return 0; } static int md_auto_download(struct UTUN_INSTANCE* inst, const char* data_str, size_t data_len, const char* ch_id, const char* base_filename, uint64_t ts, const uint8_t* author_sig, uint64_t author, int64_t msg_id, const char* content_type) { if (!chat_setting_get_int(inst, "storage_autoload", 1)) return 0; uint64_t max_size = inst->config->global.chatserver_storage_unit_size; const char* p = data_str; while (*p && *p != '|') p++; if (!*p) return 0; p++; long long fsize = strtoll(p, (char**)&p, 10); if (max_size > 0 && fsize > (long long)max_size) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: media too large (%lld > %llu bytes), skip auto-download ch=%s file=%s", CC_ID, fsize, (unsigned long long)max_size, ch_id, base_filename); return 0; } return md_start_download(inst, data_str, data_len, ch_id, base_filename, ts, author_sig, author, msg_id, content_type) == 0 ? 1 : 0; } /* ─── db_sync callback ─── */ void on_msg_inserted(struct DB_SYNC_INSTANCE* si, uint64_t record_ts, const char* data, size_t len, uint64_t author, void* arg) { (void)si; struct msg_insert_arg* ia = (struct msg_insert_arg*)arg; struct UTUN_INSTANCE* inst = ia ? ia->inst : NULL; struct chat_core_ctx* cc = CC(inst); if (!ia || !cc) return; const char* ch_id = ia->ch_id; uint8_t evt[73]; uint8_t cl = (uint8_t)strlen(ch_id); evt[0] = cl; memcpy(evt + 1, ch_id, cl); memcpy(evt + 1 + cl, &author, 8); chat_event_post(inst, CHAT_EVT_MSG_RECEIVED, evt, 1 + cl + 8); /* auto-download media if message contains media metadata */ if (data && len > 0 && author != cc->my_node_id) { /* data from db_sync may not be null-terminated; copy to local buffer */ char buf[4096]; size_t blen = len < sizeof(buf) - 1 ? len : sizeof(buf) - 1; memcpy(buf, data, blen); buf[blen] = '\0'; const char* ct = strstr(buf, "\"ct\":\""); const char* dpos = strstr(buf, "\"d\":\""); if (ct && dpos) { /* extract content_type value */ char content_type[32] = {0}; const char* ct_val = ct + 6; int cti = 0; while (ct_val[cti] && ct_val[cti] != '\"' && cti < (int)sizeof(content_type) - 1) content_type[cti++] = ct_val[cti]; content_type[cti] = '\0'; const char* d_start = dpos + 5; char body[4096]; size_t bi = 0; while (*d_start && *d_start != '\"' && bi < sizeof(body) - 1) body[bi++] = *d_start++; body[bi] = '\0'; char base_filename[256]; extract_base_filename(body, base_filename, sizeof(base_filename)); uint8_t author_sig[64] = {0}; /* We don't have author_sig in the callback; derive from body content. Instead, read it from the DB using the just-inserted record */ char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl)); char sql[256]; snprintf(sql, sizeof(sql), "SELECT author_signature, id FROM \"%s\" WHERE timestamp=? AND node_id=? LIMIT 1", tbl); sqlite3_stmt* st = NULL; if (sqlite3_prepare_v2(cc->db, sql, -1, &st, NULL) == SQLITE_OK) { sqlite3_bind_int64(st, 1, (sqlite3_int64)record_ts); sqlite3_bind_int64(st, 2, (sqlite3_int64)author); int64_t msg_id = 0; if (sqlite3_step(st) == SQLITE_ROW) { const void* sblob = sqlite3_column_blob(st, 0); int sblen = sqlite3_column_bytes(st, 0); if (sblob && sblen == 64) memcpy(author_sig, sblob, 64); msg_id = sqlite3_column_int64(st, 1); } sqlite3_finalize(st); if (msg_id > 0) md_auto_download(inst, body, bi, ch_id, base_filename, record_ts, author_sig, author, msg_id, content_type); } } } sqlite3_stmt* st = NULL; sqlite3_prepare_v2(cc->db, "UPDATE channels SET last_msg_at = ? WHERE channel_id = ?", -1, &st, NULL); if (st) { sqlite3_bind_int64(st, 1, (sqlite3_int64)record_ts); sqlite3_bind_text(st, 2, ch_id, -1, SQLITE_STATIC); sqlite3_step(st); sqlite3_finalize(st); } } /* ─── manual attachment download (called from JNI bridge via uasync post) ─── */ void chat_core_attachment_download(struct UTUN_INSTANCE* inst, const char* channel_id, int64_t msg_id) { struct chat_core_ctx* cc = CC(inst); if (!cc || !channel_id || msg_id <= 0) return; char tbl[80]; msg_table_name(channel_id, tbl, sizeof(tbl)); char sql[256]; snprintf(sql, sizeof(sql), "SELECT data, timestamp, author_signature, node_id FROM \"%s\" WHERE id=?", tbl); sqlite3_stmt* st = NULL; if (sqlite3_prepare_v2(cc->db, sql, -1, &st, NULL) != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: download sql prep failed ch=%s", CC_ID, channel_id); return; } sqlite3_bind_int64(st, 1, (sqlite3_int64)msg_id); if (sqlite3_step(st) != SQLITE_ROW) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: download msg not found ch=%s id=%lld", CC_ID, channel_id, (long long)msg_id); sqlite3_finalize(st); return; } const uint8_t* jdata = sqlite3_column_text(st, 0); int jlen = sqlite3_column_bytes(st, 0); int64_t db_ts = sqlite3_column_int64(st, 1); const uint8_t* sig_blob = (const uint8_t*)sqlite3_column_blob(st, 2); int sig_len = sqlite3_column_bytes(st, 2); int64_t author_node_id = sqlite3_column_int64(st, 3); if (!jdata || jlen <= 0 || !sig_blob || sig_len != 64) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: download bad msg data ch=%s id=%lld", CC_ID, channel_id, (long long)msg_id); sqlite3_finalize(st); return; } const char* ds = strstr((const char*)jdata, "\"d\":\""); if (!ds) { sqlite3_finalize(st); return; } const char* d_start = ds + 5; char body[4096]; size_t bi = 0; while (*d_start && *d_start != '\"' && bi < sizeof(body) - 1) body[bi++] = *d_start++; body[bi] = '\0'; char content_type[32] = {0}; { const char* ct = strstr((const char*)jdata, "\"ct\":\""); if (ct) { const char* cv = ct + 6; int i = 0; while (cv[i] && cv[i] != '\"' && i < (int)sizeof(content_type) - 1) content_type[i++] = cv[i]; content_type[i] = '\0'; } } char base_filename[256]; extract_base_filename(body, base_filename, sizeof(base_filename)); uint8_t author_sig[64]; memcpy(author_sig, sig_blob, 64); sqlite3_finalize(st); DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: manual download ch=%s id=%lld ct=%s file=%s", CC_ID, channel_id, (long long)msg_id, content_type, base_filename); md_start_download(inst, body, bi, channel_id, base_filename, (uint64_t)db_ts, author_sig, (uint64_t)author_node_id, msg_id, content_type); } static void chat_core_attachment_download_trampoline_impl(void* arg) { struct attachment_dl_req* req = (struct attachment_dl_req*)arg; if (!req) return; chat_core_attachment_download(req->inst, req->channel_id, req->msg_id); u_free(req); } void chat_core_attachment_download_trampoline(void* arg) { chat_core_attachment_download_trampoline_impl(arg); } /* ─── Backfill: регистрация уже скачанных медиафайлов в media_files ─── */ void chat_media_backfill(struct UTUN_INSTANCE* inst) { struct chat_core_ctx* cc = CC(inst); if (!cc || !cc->initialized || !cc->db) return; char media_base[512]; const char* last_slash = strrchr(cc->db_path, '/'); if (last_slash) snprintf(media_base, sizeof(media_base), "%.*s", (int)(last_slash - cc->db_path), cc->db_path); else snprintf(media_base, sizeof(media_base), "%s", cc->db_path); int total = 0; sqlite3_stmt* cs = NULL; if (sqlite3_prepare_v2(cc->db, "SELECT channel_id FROM channels", -1, &cs, NULL) != SQLITE_OK) return; while (sqlite3_step(cs) == SQLITE_ROW) { const char* ch_id = (const char*)sqlite3_column_text(cs, 0); if (!ch_id || !ch_id[0]) continue; char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl)); char sql[256]; snprintf(sql, sizeof(sql), "SELECT data FROM \"%s\"", tbl); sqlite3_stmt* ms = NULL; if (sqlite3_prepare_v2(cc->db, sql, -1, &ms, NULL) != SQLITE_OK) continue; while (sqlite3_step(ms) == SQLITE_ROW) { const char* jdata = (const char*)sqlite3_column_text(ms, 0); if (!jdata) continue; /* тело "d" (значение не содержит кавычек/слэшей) */ const char* ds = strstr(jdata, "\"d\":\""); if (!ds) continue; const char* d_start = ds + 5; const char* d_end = strchr(d_start, '"'); if (!d_end) continue; size_t dlen = (size_t)(d_end - d_start); if (dlen == 0 || dlen > 65536) continue; char* body = u_malloc(dlen + 1); if (!body) continue; memcpy(body, d_start, dlen); body[dlen] = '\0'; struct media_index_result result; if (chat_msg_parse_media_body(body, &result) != 0) { u_free(body); continue; } /* content_type нужен только чтобы отличить голосовое имя файла */ char content_type[32] = {0}; const char* ct = strstr(jdata, "\"ct\":\""); if (ct) { const char* cv = ct + 6; int i = 0; while (cv[i] && cv[i] != '"' && i < (int)sizeof(content_type) - 1) content_type[i++] = cv[i]; content_type[i] = '\0'; } char base_filename[256]; extract_base_filename(body, base_filename, sizeof(base_filename)); char final_name[256]; media_final_name(&result, content_type, base_filename, final_name, sizeof(final_name)); char path[1536]; snprintf(path, sizeof(path), "%s/media/%s/%s", media_base, ch_id, final_name); int fsz = ma_file_size(path); if (fsz <= 0 || (int64_t)fsz != result.file_size) { u_free(body); media_index_result_free(&result); continue; } int registered = 0; for (int n = 0; n < result.num_blocks; n++) { int exists = 0; sqlite3_stmt* ex = NULL; if (sqlite3_prepare_v2(cc->db, "SELECT 1 FROM media_files WHERE block_id=? LIMIT 1", -1, &ex, NULL) == SQLITE_OK) { sqlite3_bind_blob(ex, 1, result.block_ids + n * 16, 16, SQLITE_STATIC); exists = (sqlite3_step(ex) == SQLITE_ROW); sqlite3_finalize(ex); } if (exists) continue; char rel[512]; size_t bl = strlen(media_base); if (strncmp(path, media_base, bl) == 0 && path[bl] == '/') snprintf(rel, sizeof(rel), "%s", path + bl + 1); else snprintf(rel, sizeof(rel), "%s", path); int64_t chunk_size = (n == result.num_blocks - 1) ? result.file_size - (int64_t)n * result.block_size : result.block_size; int64_t offset = (int64_t)n * result.block_size; if (media_index_register_downloaded(cc->db, result.media_id, result.block_ids + n * 16, result.content_hash, ch_id, rel, cc->my_node_id, result.file_size, chunk_size, n, offset) == 0) { registered++; } } if (registered > 0) { total += registered; DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: backfill registered %d blocks ch=%s file=%s", CC_ID, registered, ch_id, final_name); } u_free(body); media_index_result_free(&result); } sqlite3_finalize(ms); } sqlite3_finalize(cs); if (total > 0) DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: media backfill done, %d blocks registered", CC_ID, total); } /* ─── Докачка недостающего медиа + вычистка старых файлов (единый проход) ─── */ struct bf_entry { uint64_t ts_ms; char full_path[1536]; int64_t claimed_size; /* размер из тела сообщения */ int64_t actual_size; /* >0 — файл на диске */ uint64_t author; int64_t msg_id; char ch_id[64]; char content_type[32]; char base_filename[256]; uint8_t author_sig[64]; char* body; /* "d"-поле (NULL для present) */ size_t body_len; }; struct bf_list { struct bf_entry* items; int count, cap; }; static int bf_cmp_ts_desc(const void* a, const void* b) { const struct bf_entry* x = (const struct bf_entry*)a; const struct bf_entry* y = (const struct bf_entry*)b; if (x->ts_ms > y->ts_ms) return -1; if (x->ts_ms < y->ts_ms) return 1; return 0; } /* Собирает все медиасообщения всех каналов (presence + размеры + контекст докачки). */ static void bf_collect(struct chat_core_ctx* cc, const char* media_base, struct bf_list* l) { sqlite3_stmt* cs = NULL; if (sqlite3_prepare_v2(cc->db, "SELECT channel_id FROM channels", -1, &cs, NULL) != SQLITE_OK) return; while (sqlite3_step(cs) == SQLITE_ROW) { const char* ch_id = (const char*)sqlite3_column_text(cs, 0); if (!ch_id || !ch_id[0]) continue; char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl)); char sql[256]; snprintf(sql, sizeof(sql), "SELECT id, timestamp, node_id, data, author_signature, local_attrs FROM \"%s\"", tbl); sqlite3_stmt* ms = NULL; if (sqlite3_prepare_v2(cc->db, sql, -1, &ms, NULL) != SQLITE_OK) continue; while (sqlite3_step(ms) == SQLITE_ROW) { int64_t msg_id = sqlite3_column_int64(ms, 0); uint64_t ts = (uint64_t)sqlite3_column_int64(ms, 1); uint64_t author = (uint64_t)sqlite3_column_int64(ms, 2); const char* jdata = (const char*)sqlite3_column_text(ms, 3); const uint8_t* sig_blob = (const uint8_t*)sqlite3_column_blob(ms, 4); int sig_len = sqlite3_column_bytes(ms, 4); const char* la = (const char*)sqlite3_column_text(ms, 5); if (!jdata || !sig_blob || sig_len != 64) continue; const char* ds = strstr(jdata, "\"d\":\""); if (!ds) continue; const char* d_start = ds + 5; const char* d_end = strchr(d_start, '"'); if (!d_end) continue; size_t dlen = (size_t)(d_end - d_start); if (dlen == 0 || dlen > 65536) continue; char* body = u_malloc(dlen + 1); if (!body) continue; memcpy(body, d_start, dlen); body[dlen] = '\0'; struct media_index_result result; if (chat_msg_parse_media_body(body, &result) != 0) { u_free(body); continue; } char content_type[32] = {0}; const char* ct = strstr(jdata, "\"ct\":\""); if (ct) { const char* cv = ct + 6; int i = 0; while (cv[i] && cv[i] != '"' && i < (int)sizeof(content_type) - 1) content_type[i++] = cv[i]; content_type[i] = '\0'; } char base_filename[256]; extract_base_filename(body, base_filename, sizeof(base_filename)); char fp[256]; extract_attrs_fp(la, fp, sizeof(fp)); char final_name[256]; if (fp[0]) snprintf(final_name, sizeof(final_name), "%s", fp); else media_final_name(&result, content_type, base_filename, final_name, sizeof(final_name)); char full_path[1536]; snprintf(full_path, sizeof(full_path), "%s/media/%s/%s", media_base, ch_id, final_name); int64_t actual = (int64_t)ma_file_size(full_path); if (actual < 0) actual = 0; if (l->count == l->cap) { int nc = l->cap ? l->cap * 2 : 64; struct bf_entry* ni = u_realloc(l->items, (size_t)nc * sizeof(struct bf_entry)); if (!ni) { media_index_result_free(&result); u_free(body); continue; } l->items = ni; l->cap = nc; } struct bf_entry* e = &l->items[l->count++]; memset(e, 0, sizeof(*e)); e->ts_ms = ts; snprintf(e->full_path, sizeof(e->full_path), "%s", full_path); e->claimed_size = result.file_size; e->actual_size = actual; e->author = author; e->msg_id = msg_id; snprintf(e->ch_id, sizeof(e->ch_id), "%s", ch_id); snprintf(e->content_type, sizeof(e->content_type), "%s", content_type); snprintf(e->base_filename, sizeof(e->base_filename), "%s", base_filename); memcpy(e->author_sig, sig_blob, 64); if (actual <= 0) { e->body = body; e->body_len = dlen; } else u_free(body); media_index_result_free(&result); } sqlite3_finalize(ms); } sqlite3_finalize(cs); } void chat_media_autodownload_backfill(struct UTUN_INSTANCE* inst) { struct chat_core_ctx* cc = CC(inst); if (!cc || !cc->initialized || !cc->db || !inst) return; uint64_t limit = inst->config->global.chatserver_storage_total_size; int autoload = chat_setting_get_int(inst, "storage_autoload", 1); int days = chat_setting_get_int(inst, "storage_backfill_days", 7); uint64_t now_ms = (uint64_t)ntp_time_get_seconds(inst) * 1000ULL; uint64_t cutoff_ms = now_ms - (uint64_t)days * 86400ULL * 1000ULL; char media_base[512]; const char* last_slash = strrchr(cc->db_path, '/'); if (last_slash) snprintf(media_base, sizeof(media_base), "%.*s", (int)(last_slash - cc->db_path), cc->db_path); else snprintf(media_base, sizeof(media_base), "%s", cc->db_path); struct bf_list list = {0}; bf_collect(cc, media_base, &list); if (list.count > 1) qsort(list.items, (size_t)list.count, sizeof(struct bf_entry), bf_cmp_ts_desc); /* инициализация счётчика: сумма всех present-файлов */ cc->media_cache_bytes = 0; for (int i = 0; i < list.count; i++) if (list.items[i].actual_size > 0) cc->media_cache_bytes += (uint64_t)list.items[i].actual_size; uint64_t used = 0; int present = 0, downloaded = 0, evicted = 0, skipped = 0; uint64_t evicted_bytes = 0; for (int i = 0; i < list.count; i++) { struct bf_entry* e = &list.items[i]; if (e->actual_size > 0) { if (limit == 0 || used + (uint64_t)e->actual_size <= limit) { used += (uint64_t)e->actual_size; present++; } else { if (remove(e->full_path) == 0) { cache_bytes_sub(cc, e->actual_size); evicted++; evicted_bytes += (uint64_t)e->actual_size; DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: cache evict: %s ts=%llu size=%lld", CC_ID, e->full_path, (unsigned long long)e->ts_ms, (long long)e->actual_size); } else { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: cache evict failed: %s", CC_ID, e->full_path); } } } else { if (!autoload || e->author == cc->my_node_id) { skipped++; continue; } if (e->ts_ms < cutoff_ms) { skipped++; continue; } if (limit > 0 && used + (uint64_t)e->claimed_size > limit) { skipped++; continue; } if (md_auto_download(inst, e->body, e->body_len, e->ch_id, e->base_filename, e->ts_ms, e->author_sig, e->author, e->msg_id, e->content_type)) { downloaded++; used += (uint64_t)e->claimed_size; } else { skipped++; } } } DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: cache backfill done present=%d downloaded=%d skipped=%d evicted=%d evicted_bytes=%llu bytes=%llu limit=%llu", CC_ID, present, downloaded, skipped, evicted, (unsigned long long)evicted_bytes, (unsigned long long)cc->media_cache_bytes, (unsigned long long)limit); for (int i = 0; i < list.count; i++) if (list.items[i].body) u_free(list.items[i].body); u_free(list.items); chat_media_cache_reconcile(inst); } /* ─── Кеш медиафайлов: очистка сирот + лимит storage_total_size ─── */ struct media_msg_file { char rel_path[1536]; // media// char full_path[1536]; // /media// uint64_t ts_ms; }; typedef void (*media_msg_file_cb)(const struct media_msg_file* f, void* arg); static void chat_media_base(struct chat_core_ctx* cc, char* out, size_t out_sz) { const char* ls = strrchr(cc->db_path, '/'); if (ls) snprintf(out, out_sz, "%.*s", (int)(ls - cc->db_path), cc->db_path); else snprintf(out, out_sz, "%s", cc->db_path); } /* Итератор по всем медиасообщениям всех каналов (без фильтра по node_id/времени). */ static void chat_media_foreach_media_message(struct chat_core_ctx* cc, const char* media_base, media_msg_file_cb cb, void* arg) { if (!cc->db || !cb) return; sqlite3_stmt* cs = NULL; if (sqlite3_prepare_v2(cc->db, "SELECT channel_id FROM channels", -1, &cs, NULL) != SQLITE_OK) return; while (sqlite3_step(cs) == SQLITE_ROW) { const char* ch_id = (const char*)sqlite3_column_text(cs, 0); if (!ch_id || !ch_id[0]) continue; char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl)); char sql[256]; snprintf(sql, sizeof(sql), "SELECT timestamp, data, local_attrs FROM \"%s\"", tbl); sqlite3_stmt* ms = NULL; if (sqlite3_prepare_v2(cc->db, sql, -1, &ms, NULL) != SQLITE_OK) continue; while (sqlite3_step(ms) == SQLITE_ROW) { uint64_t ts = (uint64_t)sqlite3_column_int64(ms, 0); const char* jdata = (const char*)sqlite3_column_text(ms, 1); const char* la = (const char*)sqlite3_column_text(ms, 2); if (!jdata) continue; const char* ds = strstr(jdata, "\"d\":\""); if (!ds) continue; const char* d_start = ds + 5; const char* d_end = strchr(d_start, '"'); if (!d_end) continue; size_t dlen = (size_t)(d_end - d_start); if (dlen == 0 || dlen > 65536) continue; char* body = u_malloc(dlen + 1); if (!body) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: cache scan malloc failed ch=%s", CC_ID, ch_id); continue; } memcpy(body, d_start, dlen); body[dlen] = '\0'; struct media_index_result result; if (chat_msg_parse_media_body(body, &result) != 0) { u_free(body); continue; } char content_type[32] = {0}; const char* ct = strstr(jdata, "\"ct\":\""); if (ct) { const char* cv = ct + 6; int i = 0; while (cv[i] && cv[i] != '"' && i < (int)sizeof(content_type) - 1) content_type[i++] = cv[i]; content_type[i] = '\0'; } char fp[256]; extract_attrs_fp(la, fp, sizeof(fp)); char final_name[256]; if (fp[0]) { snprintf(final_name, sizeof(final_name), "%s", fp); } else { char base_filename[256]; extract_base_filename(body, base_filename, sizeof(base_filename)); media_final_name(&result, content_type, base_filename, final_name, sizeof(final_name)); } struct media_msg_file f; snprintf(f.rel_path, sizeof(f.rel_path), "media/%s/%s", ch_id, final_name); snprintf(f.full_path, sizeof(f.full_path), "%s/media/%s/%s", media_base, ch_id, final_name); f.ts_ms = ts; cb(&f, arg); u_free(body); media_index_result_free(&result); } sqlite3_finalize(ms); } sqlite3_finalize(cs); } struct cache_file { char full_path[1536]; uint64_t ts_ms; int64_t size; }; struct cache_rel_list { char** items; int count, cap; }; struct cache_file_list { struct cache_file* files; int count, cap; }; static int cache_str_cmp(const void* a, const void* b) { return strcmp(*(const char* const*)a, *(const char* const*)b); } static int cache_file_cmp_ts(const void* a, const void* b) { const struct cache_file* x = (const struct cache_file*)a; const struct cache_file* y = (const struct cache_file*)b; if (x->ts_ms < y->ts_ms) return -1; if (x->ts_ms > y->ts_ms) return 1; return 0; } static int cache_rel_add(struct cache_rel_list* l, const char* rel) { if (l->count == l->cap) { int nc = l->cap ? l->cap * 2 : 64; char** ni = u_realloc(l->items, (size_t)nc * sizeof(char*)); if (!ni) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: cache rel list realloc failed (%d)", CC_ID, nc); return -1; } l->items = ni; l->cap = nc; } l->items[l->count] = u_strdup(rel); if (!l->items[l->count]) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: cache rel strdup failed", CC_ID); return -1; } l->count++; return 0; } static int cache_file_add(struct cache_file_list* l, const char* full_path, uint64_t ts_ms, int64_t size) { if (l->count == l->cap) { int nc = l->cap ? l->cap * 2 : 64; struct cache_file* ni = u_realloc(l->files, (size_t)nc * sizeof(struct cache_file)); if (!ni) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: cache file list realloc failed (%d)", CC_ID, nc); return -1; } l->files = ni; l->cap = nc; } snprintf(l->files[l->count].full_path, sizeof(l->files[l->count].full_path), "%s", full_path); l->files[l->count].ts_ms = ts_ms; l->files[l->count].size = size; l->count++; return 0; } static void cache_rel_collect_cb(const struct media_msg_file* f, void* arg) { cache_rel_add((struct cache_rel_list*)arg, f->rel_path); } static void cache_file_collect_cb(const struct media_msg_file* f, void* arg) { int fsz = ma_file_size(f->full_path); if (fsz > 0) cache_file_add((struct cache_file_list*)arg, f->full_path, f->ts_ms, (int64_t)fsz); } /* Фактическая сумма размеров файлов, соответствующих сообщениям. */ static uint64_t chat_media_cache_scan_sum(struct chat_core_ctx* cc) { char media_base[512]; chat_media_base(cc, media_base, sizeof(media_base)); struct cache_file_list files = {0}; chat_media_foreach_media_message(cc, media_base, cache_file_collect_cb, &files); uint64_t sum = 0; for (int i = 0; i < files.count; i++) if (files.files[i].size > 0) sum += (uint64_t)files.files[i].size; u_free(files.files); return sum; } /* Сверка счётчика с фактической суммой: при расхождении — WARN + самокоррекция. */ static void chat_media_cache_reconcile(struct UTUN_INSTANCE* inst) { struct chat_core_ctx* cc = CC(inst); if (!cc || !cc->initialized || !cc->db || !inst) return; uint64_t actual = chat_media_cache_scan_sum(cc); if (actual != cc->media_cache_bytes) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: cache counter drift: counter=%llu actual=%llu diff=%lld", CC_ID, (unsigned long long)cc->media_cache_bytes, (unsigned long long)actual, (long long)((int64_t)actual - (int64_t)cc->media_cache_bytes)); cc->media_cache_bytes = actual; } else { DEBUG_DEBUG(DEBUG_CATEGORY_GENERAL, "%s: cache reconcile ok bytes=%llu", CC_ID, (unsigned long long)cc->media_cache_bytes); } } struct orphan_walk_ctx { const char* media_base; size_t mb_len; char** rels; int rel_count; int removed; uint64_t removed_bytes; }; static void orphan_walk_cb(void* arg, const char* full_path, int64_t size) { struct orphan_walk_ctx* ow = (struct orphan_walk_ctx*)arg; const char* rel = full_path; if (strncmp(full_path, ow->media_base, ow->mb_len) == 0 && full_path[ow->mb_len] == '/') rel = full_path + ow->mb_len + 1; else return; char* found = bsearch(&rel, ow->rels, (size_t)ow->rel_count, sizeof(char*), cache_str_cmp); if (found) return; if (remove(full_path) == 0) { ow->removed++; ow->removed_bytes += (uint64_t)(size > 0 ? size : 0); DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: cache orphan removed: %s (%lld bytes)", CC_ID, rel, (long long)size); } else { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: cache orphan remove failed: %s", CC_ID, full_path); } } static void chat_media_cache_cleanup_orphans(struct UTUN_INSTANCE* inst) { struct chat_core_ctx* cc = CC(inst); if (!cc || !cc->initialized || !cc->db || !inst) return; char media_base[512]; chat_media_base(cc, media_base, sizeof(media_base)); char media_dir[600]; snprintf(media_dir, sizeof(media_dir), "%s/media", media_base); struct cache_rel_list rels = {0}; chat_media_foreach_media_message(cc, media_base, cache_rel_collect_cb, &rels); if (rels.count > 1) qsort(rels.items, (size_t)rels.count, sizeof(char*), cache_str_cmp); struct orphan_walk_ctx ow; ow.media_base = media_base; ow.mb_len = strlen(media_base); ow.rels = rels.items; ow.rel_count = rels.count; ow.removed = 0; ow.removed_bytes = 0; ma_dir_walk(media_dir, orphan_walk_cb, &ow); DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: cache orphan cleanup done removed=%d bytes=%llu", CC_ID, ow.removed, (unsigned long long)ow.removed_bytes); for (int i = 0; i < rels.count; i++) u_free(rels.items[i]); u_free(rels.items); } static void chat_media_cache_enforce(struct UTUN_INSTANCE* inst) { struct chat_core_ctx* cc = CC(inst); if (!cc || !cc->initialized || !cc->db || !inst) return; uint64_t limit = inst->config->global.chatserver_storage_total_size; if (limit == 0 || cc->media_cache_bytes <= limit) { DEBUG_DEBUG(DEBUG_CATEGORY_GENERAL, "%s: cache limit ok bytes=%llu limit=%llu", CC_ID, (unsigned long long)cc->media_cache_bytes, (unsigned long long)limit); chat_media_cache_reconcile(inst); return; } char media_base[512]; chat_media_base(cc, media_base, sizeof(media_base)); struct cache_file_list files = {0}; chat_media_foreach_media_message(cc, media_base, cache_file_collect_cb, &files); if (files.count > 1) qsort(files.files, (size_t)files.count, sizeof(struct cache_file), cache_file_cmp_ts); int deleted = 0; uint64_t deleted_bytes = 0; for (int i = 0; i < files.count && cc->media_cache_bytes > limit; i++) { if (files.files[i].size <= 0) continue; if (remove(files.files[i].full_path) == 0) { cache_bytes_sub(cc, files.files[i].size); deleted++; deleted_bytes += (uint64_t)files.files[i].size; DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: cache evict: %s ts=%llu size=%lld", CC_ID, files.files[i].full_path, (unsigned long long)files.files[i].ts_ms, (long long)files.files[i].size); } else { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: cache evict failed: %s", CC_ID, files.files[i].full_path); } } if (cc->media_cache_bytes > limit) DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: cache still over limit bytes=%llu limit=%llu", CC_ID, (unsigned long long)cc->media_cache_bytes, (unsigned long long)limit); else DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: cache limit enforced bytes=%llu limit=%llu deleted=%d bytes=%llu", CC_ID, (unsigned long long)cc->media_cache_bytes, (unsigned long long)limit, deleted, (unsigned long long)deleted_bytes); u_free(files.files); chat_media_cache_reconcile(inst); } /* ─── Триггер: после первого SYNC_DONE (+ дебаунс) — анонс локальных блоков + докачка ─── */ static void media_startup_backfill_timer_cb(void* arg) { struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)arg; struct chat_core_ctx* cc = CC(inst); cc->media_backfill_timer = NULL; cc->media_backfill_ran = 1; /* 1) переанонс локальных блоков суперузлам (block_availability), * 2) чистка сирот, 3) докачка недостающего + вычистка старого (единый проход) */ media_delivery_announce_local_blocks(inst); chat_media_cache_cleanup_orphans(inst); chat_media_autodownload_backfill(inst); } static void media_startup_backfill_done_cb(struct DB_SYNC_INSTANCE* si, uint64_t peer_node_id, void* arg) { struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)arg; struct chat_core_ctx* cc = CC(inst); (void)si; (void)peer_node_id; if (!cc || cc->media_backfill_ran || !cc->media_backfill_inst) return; if (cc->media_backfill_timer) uasync_cancel_timeout(inst->ua, cc->media_backfill_timer); cc->media_backfill_timer = uasync_set_timeout(inst->ua, 20000, inst, media_startup_backfill_timer_cb, "chat_media_bl"); } void chat_media_startup_backfill_init(struct UTUN_INSTANCE* inst) { struct chat_core_ctx* cc = CC(inst); if (!inst || !cc) return; cc->media_backfill_inst = inst; cc->media_backfill_ran = 0; db_sync_add_done_cbk(inst, media_startup_backfill_done_cb, inst); } void chat_media_startup_backfill_destroy(struct UTUN_INSTANCE* inst) { struct chat_core_ctx* cc = CC(inst); if (!inst || !cc) return; db_sync_remove_done_cbk(inst, media_startup_backfill_done_cb, inst); if (cc->media_backfill_timer) { uasync_cancel_timeout(inst->ua, cc->media_backfill_timer); cc->media_backfill_timer = NULL; } cc->media_backfill_inst = NULL; cc->media_backfill_ran = 0; }