diff --git a/src/chat/chat_msg.c b/src/chat/chat_msg.c index 64eec0ad..93de0f72 100644 --- a/src/chat/chat_msg.c +++ b/src/chat/chat_msg.c @@ -32,6 +32,32 @@ static void wh_transcribe_done_cb(void* arg, const char* ch_id, const char* text uint64_t reply_ts, uint64_t reply_node); static void chat_media_cache_cleanup_orphans(void); static void chat_media_cache_enforce(void); +static void chat_media_cache_reconcile(void); + +/* ─── счётчик объёма медиакеша (сумма файлов, соответствующих сообщениям) ─── */ + +static uint64_t g_media_cache_bytes = 0; + +static void cache_bytes_add(int64_t sz) { + if (sz > 0) g_media_cache_bytes += (uint64_t)sz; +} + +static void cache_bytes_sub(int64_t sz) { + uint64_t v = sz > 0 ? (uint64_t)sz : 0; + g_media_cache_bytes = g_media_cache_bytes >= v ? g_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) ─── */ @@ -203,6 +229,10 @@ static void on_media_registered(void* arg, int err, const struct media_index_res media_delivery_announce_media(g_cc.inst, strtoull(mctx->channel_id, NULL, 10), result->media_id, result->block_ids, result->num_blocks); + /* учесть локально отправленный файл в счётчике и применить лимит медиакеша */ + cache_bytes_add(result->file_size); + chat_media_cache_enforce(); + /* whisper транскрипция для своих голосовых сообщений */ if (mctx->content_type[0] && strncmp(mctx->content_type, "voice", 5) == 0 && g_cc.inst && g_cc.inst->media_async && g_chat_whisper_trigger) { @@ -258,6 +288,17 @@ void chat_core_submit_media_message(struct chat_msg_submit* req) { if (!g_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 = g_cc.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); @@ -577,7 +618,8 @@ static void md_download_done_cb(void* arg, int err) { wh_transcribe_done_cb, NULL); } - /* после докачки — применить лимит медиакеша */ + /* после докачки — учесть файл в счётчике и применить лимит медиакеша */ + cache_bytes_add(ma_file_size(ctx->dest_abs_path)); chat_media_cache_enforce(); } else { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: media download failed ch=%s err=%d", CC_ID, ctx->channel_id, err); @@ -1029,24 +1071,42 @@ void chat_media_backfill(void) { DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: media backfill done, %d blocks registered", CC_ID, total); } -/* ─── Докачка недостающего медиа после стартовой синхронизации ─── - * Проходим по сообщениям каналов за последние storage_backfill_days дней и для - * медиасообщений, чей файл отсутствует на диске, запускаем загрузку. Вызывается - * один раз после завершения первой синхронизации сообщений чата. */ -void chat_media_autodownload_backfill(void) { - if (!g_cc.initialized || !g_cc.db || !g_cc.inst) return; - if (!chat_setting_get_int("storage_autoload", 1)) return; +/* ─── Докачка недостающего медиа + вычистка старых файлов (единый проход) ─── + * Инвариант кеша: на диске лежат самые свежие storage_total_size байт медиа + * (по timestamp сообщения). Один проход DESC по дате: свежие недостающие — + * докачиваем (в пределах окна storage_backfill_days и бюджета), старые сверх + * бюджета — вычищаем. Фронт докачки идёт с нового конца, фронт вычистки — со + * старого; сходятся на границе бюджета (used упирается в limit). Счётчик + * g_media_cache_bytes корректируется на каждой вычистке; докачки добавляют + * через md_download_done_cb. Вызывается один раз после стартовой синхронизации. */ - int days = chat_setting_get_int("storage_backfill_days", 7); - uint64_t now_ms = (uint64_t)ntp_time_get_seconds(g_cc.inst) * 1000ULL; - uint64_t cutoff_ms = now_ms - (uint64_t)days * 86400ULL * 1000ULL; +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; +}; - 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); +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; +} - int scanned = 0, present = 0, started = 0, skipped = 0; +/* Собирает все медиасообщения всех каналов (presence + размеры + контекст докачки). */ +static void bf_collect(const char* media_base, struct bf_list* l) { sqlite3_stmt* cs = NULL; if (sqlite3_prepare_v2(g_cc.db, "SELECT channel_id FROM channels", -1, &cs, NULL) != SQLITE_OK) return; @@ -1057,12 +1117,9 @@ void chat_media_autodownload_backfill(void) { char sql[256]; snprintf(sql, sizeof(sql), - "SELECT id, timestamp, node_id, data, author_signature FROM \"%s\" " - "WHERE timestamp>=? AND node_id!=? ORDER BY timestamp ASC", tbl); + "SELECT id, timestamp, node_id, data, author_signature, local_attrs FROM \"%s\"", tbl); sqlite3_stmt* ms = NULL; if (sqlite3_prepare_v2(g_cc.db, sql, -1, &ms, NULL) != SQLITE_OK) continue; - sqlite3_bind_int64(ms, 1, (sqlite3_int64)cutoff_ms); - sqlite3_bind_int64(ms, 2, (sqlite3_int64)g_cc.my_node_id); while (sqlite3_step(ms) == SQLITE_ROW) { int64_t msg_id = sqlite3_column_int64(ms, 0); @@ -1071,6 +1128,7 @@ void chat_media_autodownload_backfill(void) { 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\":\""); @@ -1098,39 +1156,117 @@ void chat_media_autodownload_backfill(void) { 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); + 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); - scanned++; - int fsz = ma_file_size(path); - if (fsz > 0 && (int64_t)fsz == result.file_size) { present++; u_free(body); media_index_result_free(&result); continue; } - - uint8_t author_sig[64]; memcpy(author_sig, sig_blob, 64); - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: autodownload missing ch=%s id=%lld file=%s size=%lld", - CC_ID, ch_id, (long long)msg_id, final_name, (long long)result.file_size); - if (md_auto_download(g_cc.inst, body, dlen, ch_id, base_filename, ts, author_sig, author, msg_id, content_type)) - started++; - else - skipped++; - u_free(body); media_index_result_free(&result); } sqlite3_finalize(ms); } sqlite3_finalize(cs); +} + +void chat_media_autodownload_backfill(void) { + if (!g_cc.initialized || !g_cc.db || !g_cc.inst) return; + + uint64_t limit = g_cc.inst->config->global.chatserver_storage_total_size; + int autoload = chat_setting_get_int("storage_autoload", 1); + int days = chat_setting_get_int("storage_backfill_days", 7); + uint64_t now_ms = (uint64_t)ntp_time_get_seconds(g_cc.inst) * 1000ULL; + uint64_t cutoff_ms = now_ms - (uint64_t)days * 86400ULL * 1000ULL; + + 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); - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: autodownload backfill done scanned=%d present=%d started=%d skipped=%d days=%d", - CC_ID, scanned, present, started, skipped, days); + struct bf_list list = {0}; + bf_collect(media_base, &list); + if (list.count > 1) qsort(list.items, (size_t)list.count, sizeof(struct bf_entry), bf_cmp_ts_desc); + + /* инициализация счётчика: сумма всех present-файлов */ + g_media_cache_bytes = 0; + for (int i = 0; i < list.count; i++) + if (list.items[i].actual_size > 0) g_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(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 == g_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(g_cc.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)g_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(); } /* ─── Кеш медиафайлов: очистка сирот + лимит storage_total_size ─── * Возраст файла = timestamp исходного сообщения (msg_.timestamp). * Сироты (файлы без сообщения, включая .chunk_* остатки) удаляются на старте. - * Лимит: при превышении storage_total_size удаляются самые старые файлы, но только - * старше storage_backfill_days (их backfill уже не докачает на следующем старте). */ + * Имя файла на диске берётся из local_attrs "fp" (fallback — по имени из "d"). */ struct media_msg_file { char rel_path[1536]; // media// @@ -1158,13 +1294,14 @@ static void chat_media_foreach_media_message(const char* media_base, media_msg_f char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl)); char sql[256]; - snprintf(sql, sizeof(sql), "SELECT timestamp, data FROM \"%s\"", tbl); + snprintf(sql, sizeof(sql), "SELECT timestamp, data, local_attrs 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) { 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\":\""); @@ -1190,10 +1327,16 @@ static void chat_media_foreach_media_message(const char* media_base, media_msg_f 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]; - media_final_name(&result, content_type, base_filename, final_name, sizeof(final_name)); + 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); @@ -1262,6 +1405,33 @@ static void cache_file_collect_cb(const struct media_msg_file* f, void* arg) { 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(void) { + char media_base[512]; chat_media_base(media_base, sizeof(media_base)); + struct cache_file_list files = {0}; + chat_media_foreach_media_message(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(void) { + if (!g_cc.initialized || !g_cc.db || !g_cc.inst) return; + uint64_t actual = chat_media_cache_scan_sum(); + if (actual != g_media_cache_bytes) { + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: cache counter drift: counter=%llu actual=%llu diff=%lld", + CC_ID, (unsigned long long)g_media_cache_bytes, (unsigned long long)actual, + (long long)((int64_t)actual - (int64_t)g_media_cache_bytes)); + g_media_cache_bytes = actual; + } else { + DEBUG_DEBUG(DEBUG_CATEGORY_GENERAL, "%s: cache reconcile ok bytes=%llu", + CC_ID, (unsigned long long)g_media_cache_bytes); + } +} + struct orphan_walk_ctx { const char* media_base; size_t mb_len; @@ -1315,52 +1485,44 @@ static void chat_media_cache_cleanup_orphans(void) { static void chat_media_cache_enforce(void) { if (!g_cc.initialized || !g_cc.db || !g_cc.inst) return; uint64_t limit = g_cc.inst->config->global.chatserver_storage_total_size; - if (limit == 0) return; + + if (limit == 0 || g_media_cache_bytes <= limit) { + DEBUG_DEBUG(DEBUG_CATEGORY_GENERAL, "%s: cache limit ok bytes=%llu limit=%llu", + CC_ID, (unsigned long long)g_media_cache_bytes, (unsigned long long)limit); + chat_media_cache_reconcile(); + return; + } char media_base[512]; chat_media_base(media_base, sizeof(media_base)); struct cache_file_list files = {0}; chat_media_foreach_media_message(media_base, cache_file_collect_cb, &files); - - uint64_t total = 0; - for (int i = 0; i < files.count; i++) total += (uint64_t)(files.files[i].size > 0 ? files.files[i].size : 0); - - if (total <= limit) { - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: cache limit ok total=%llu limit=%llu files=%d", - CC_ID, (unsigned long long)total, (unsigned long long)limit, files.count); - u_free(files.files); - return; - } - if (files.count > 1) qsort(files.files, (size_t)files.count, sizeof(struct cache_file), cache_file_cmp_ts); - int days = chat_setting_get_int("storage_backfill_days", 7); - uint64_t now_ms = (uint64_t)ntp_time_get_seconds(g_cc.inst) * 1000ULL; - uint64_t cutoff_ms = now_ms - (uint64_t)days * 86400ULL * 1000ULL; - int deleted = 0; uint64_t deleted_bytes = 0; - for (int i = 0; i < files.count && total > limit; i++) { - if (files.files[i].ts_ms >= cutoff_ms) continue; + for (int i = 0; i < files.count && g_media_cache_bytes > limit; i++) { + if (files.files[i].size <= 0) continue; if (remove(files.files[i].full_path) == 0) { + cache_bytes_sub(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); - total -= (uint64_t)(files.files[i].size > 0 ? files.files[i].size : 0); - deleted_bytes += (uint64_t)(files.files[i].size > 0 ? files.files[i].size : 0); - deleted++; } else { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: cache evict failed: %s", CC_ID, files.files[i].full_path); } } - if (total > limit) - DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: cache still over limit total=%llu limit=%llu (older-than-%d-days files exhausted)", - CC_ID, (unsigned long long)total, (unsigned long long)limit, days); + if (g_media_cache_bytes > limit) + DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: cache still over limit bytes=%llu limit=%llu", + CC_ID, (unsigned long long)g_media_cache_bytes, (unsigned long long)limit); else - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: cache limit enforced total=%llu limit=%llu deleted=%d bytes=%llu", - CC_ID, (unsigned long long)total, (unsigned long long)limit, deleted, (unsigned long long)deleted_bytes); + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: cache limit enforced bytes=%llu limit=%llu deleted=%d bytes=%llu", + CC_ID, (unsigned long long)g_media_cache_bytes, (unsigned long long)limit, + deleted, (unsigned long long)deleted_bytes); u_free(files.files); + chat_media_cache_reconcile(); } /* ─── Триггер: после первого SYNC_DONE (+ дебаунс) — анонс локальных блоков + докачка ─── */ @@ -1373,12 +1535,11 @@ static void media_startup_backfill_timer_cb(void* arg) { struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)arg; g_media_startup_backfill_timer = NULL; g_media_startup_backfill_ran = 1; - /* 1) переанонс локальных блоков суперузлам (block_availability), 2) докачка недостающего */ + /* 1) переанонс локальных блоков суперузлам (block_availability), + * 2) чистка сирот, 3) докачка недостающего + вычистка старого (единый проход) */ media_delivery_announce_local_blocks(inst); - chat_media_autodownload_backfill(); - /* 3) сверка кеша: удалить сироты и применить лимит storage_total_size */ chat_media_cache_cleanup_orphans(); - chat_media_cache_enforce(); + chat_media_autodownload_backfill(); } static void media_startup_backfill_done_cb(struct DB_SYNC_INSTANCE* si, uint64_t peer_node_id, void* arg) {