diff --git a/src/chat/chat_core.h b/src/chat/chat_core.h index ca8a5496..f1801187 100644 --- a/src/chat/chat_core.h +++ b/src/chat/chat_core.h @@ -20,6 +20,10 @@ struct UTUN_INSTANCE; int chat_core_init(struct UTUN_INSTANCE* inst, const char* db_path); void chat_core_destroy(struct UTUN_INSTANCE* inst); +/* Backfill: регистрирует уже скачанные медиафайлы (лежат на диске, но отсутствуют + * в media_files) — вызывается после загрузки каналов/сообщений при старте. */ +void chat_media_backfill(void); + struct sqlite3* chat_core_get_db(void); struct UTUN_INSTANCE* chat_core_get_inst(void); int chat_core_is_initialized(void); diff --git a/src/chat/chat_msg.c b/src/chat/chat_msg.c index 03dfaabe..758e23ea 100644 --- a/src/chat/chat_msg.c +++ b/src/chat/chat_msg.c @@ -14,6 +14,7 @@ #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/platform_compat.h" @@ -641,65 +642,81 @@ static void extract_base_filename(const char* body, char* out, size_t out_sz) { sanitize_filename(out, out_sz); } -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; - const char* p = data_str; +/* Разбор тела "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 != '|') return -1; p++; + if (*p != '|' || fsize <= 0) return -1; p++; long long bsize = strtoll(p, (char**)&p, 10); - if (*p != '|') return -1; p++; + if (*p != '|' || bsize <= 0) return -1; p++; int nb = (int)strtol(p, (char**)&p, 10); while (*p && *p != '|') p++; - if (!*p) return -1; p++; - if (nb <= 0 || nb > 100 || fsize <= 0) return -1; + if (!*p || nb <= 0 || nb > 100) return -1; p++; - struct media_index_result result; - memset(&result, 0, sizeof(result)); - result.file_size = fsize; result.block_size = bsize; result.num_blocks = nb; + 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); + 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); + 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; } + 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); + result->block_ids[n * 16 + i] = (uint8_t)strtol(b, NULL, 16); } - p += 32; if (*p != ',') { media_index_result_free(&result); return -1; } p++; + 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); + result->block_sigs[n * 64 + i] = (uint8_t)strtol(b, NULL, 16); } p += 128; if (n < nb - 1 && *p == ',') p++; } + return 0; +} - /* для голосовых имя генерируем из media_id: msg_data — это wf-строка, а не имя файла */ - char final_name[256]; +/* Имя файла на диске: голосовое — 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(final_name, sizeof(final_name), "voice_%s.opus", hex); + 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(final_name, sizeof(final_name), "%s", base_filename); + 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 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]; { @@ -891,3 +908,116 @@ static void chat_core_attachment_download_trampoline_impl(void* arg) { void chat_core_attachment_download_trampoline(void* arg) { chat_core_attachment_download_trampoline_impl(arg); } + +/* ─── Backfill: регистрация уже скачанных медиафайлов в media_files ─── + * Файлы, скачанные старым билдом (до md_dl_register_servable), лежат на диске, + * но отсутствуют в media_files. При старте проходим по сообщениям каналов, + * парсим тело "d" (метаданные media_id/block_ids/content_hash уже там) и для + * существующих на диске файлов дорегистрируем блоки — узел снова может их отдавать. */ +void chat_media_backfill(void) { + if (!g_cc.initialized || !g_cc.db) return; + + char media_base[512]; + const char* last_slash = strrchr(g_cc.db_path, '/'); + if (last_slash) snprintf(media_base, sizeof(media_base), "%.*s", (int)(last_slash - g_cc.db_path), g_cc.db_path); + else snprintf(media_base, sizeof(media_base), "%s", g_cc.db_path); + + int total = 0; + sqlite3_stmt* cs = NULL; + if (sqlite3_prepare_v2(g_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(g_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(g_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(g_cc.db, result.media_id, result.block_ids + n * 16, + result.content_hash, ch_id, rel, g_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); +} diff --git a/src/media_delivery/media_delivery.c b/src/media_delivery/media_delivery.c index cea3b1e9..1620bef8 100644 --- a/src/media_delivery/media_delivery.c +++ b/src/media_delivery/media_delivery.c @@ -12,6 +12,7 @@ #include "../transport_layer/secure_channel.h" #include "../chat/member_sync.h" #include "../lib/debug_config.h" +#include "../lib/json_flat.h" #include "../lib/mem.h" #include "../lib/u_async.h" #include "../lib/sqlite3.h" @@ -1189,9 +1190,17 @@ static void md_on_bgp_node(struct TOPO_GROUP* group, uint64_t node_id, int event } } +/* adm_tags хранится как JSON: {"supernode":"yes",...}. Проверяем через json_flat_get, + * т.к. strstr(tags,"supernode=yes") не совпадает с JSON-форматом. */ +static int md_adm_tags_is_supernode(const char* adm_tags) { + char buf[16]; + return adm_tags && json_flat_get(adm_tags, "supernode", buf, sizeof(buf)) == 0 + && strcmp(buf, "yes") == 0; +} + static void md_on_props_changed(uint64_t node_id, const char* adm_tags, const char* channel_id, void* arg) { struct media_delivery_ctx* md = (struct media_delivery_ctx*)arg; - int is_super = adm_tags && strstr(adm_tags, "supernode=yes"); + int is_super = md_adm_tags_is_supernode(adm_tags); if (node_id == md->self_node_id) { media_delivery_set_supernode(md->inst, is_super); @@ -1385,7 +1394,7 @@ int media_delivery_init(struct UTUN_INSTANCE* inst) { if (sqlite3_prepare_v2(inst->topo_sqlite_db, sql, -1, &st, NULL) == SQLITE_OK) { if (sqlite3_step(st) == SQLITE_ROW) { const char* tags = (const char*)sqlite3_column_text(st, 0); - if (tags && strstr(tags, "supernode=yes")) { md->is_supernode = 1; } + if (md_adm_tags_is_supernode(tags)) { md->is_supernode = 1; } } sqlite3_finalize(st); } diff --git a/src/media_delivery/media_download.c b/src/media_delivery/media_download.c index 65f26d68..8149765c 100644 --- a/src/media_delivery/media_download.c +++ b/src/media_delivery/media_download.c @@ -418,12 +418,19 @@ static void md_dl_collect_supernodes(struct UTUN_INSTANCE* inst, struct media_do gle = gle->next; } dl->super_current = 0; - if (dl->super_count == 0 && dl->author_node_id && dl->author_node_id != inst->node_id) { - dl->super_nodes[0] = dl->author_node_id; - dl->super_count = 1; - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: no supernodes, using author 0x%016llx as fallback", - MDL_ID, (unsigned long long)dl->author_node_id); - } else if (dl->super_count == 0) { + /* Автор всегда добавляется последним кандидатом на запрос: если суперузлы вернули + * пустой QUERY_RESP, watchdog доберёт автора, который отдаёт собственные media_files. */ + if (dl->author_node_id && dl->author_node_id != inst->node_id && dl->super_count < 10) { + int dup = 0; + for (int i = 0; i < dl->super_count; i++) + if (dl->super_nodes[i] == dl->author_node_id) { dup = 1; break; } + if (!dup) { + dl->super_nodes[dl->super_count++] = dl->author_node_id; + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: author 0x%016llx added as query fallback (supers=%d)", + MDL_ID, (unsigned long long)dl->author_node_id, dl->super_count); + } + } + if (dl->super_count == 0) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: no supernodes found", MDL_ID); } } diff --git a/src/utun_instance.c b/src/utun_instance.c index 85641ae7..6c2f5bd6 100644 --- a/src/utun_instance.c +++ b/src/utun_instance.c @@ -660,8 +660,10 @@ int utun_instance_init(struct UTUN_INSTANCE *instance) { mkdir_recursive(instance->config->global.db_path); char db_file[512]; snprintf(db_file, sizeof(db_file), "%s/chats.db", instance->config->global.db_path); - if (chat_core_init(instance, db_file) == 0) + if (chat_core_init(instance, db_file) == 0) { chat_sync_init(instance); + chat_media_backfill(); + } } // Set TUN interface in routing module