Browse Source

media: фикс детекта суперузла (json_flat_get), backfill осиротевших медиа, fallback на автора при пустом QUERY_RESP

v2
evgeny 3 weeks ago
parent
commit
0055c2c6fe
  1. 4
      src/chat/chat_core.h
  2. 186
      src/chat/chat_msg.c
  3. 13
      src/media_delivery/media_delivery.c
  4. 19
      src/media_delivery/media_download.c
  5. 4
      src/utun_instance.c

4
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);

186
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. Формат:
* <name>|<file_size>|<block_size>|<num_blocks>|<media_id_hex>|<content_hash_hex>|<id_hex>,<sig_hex>,...
* Выделяет 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_<media_id>.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);
}

13
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);
}

19
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);
}
}

4
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

Loading…
Cancel
Save