Browse Source

chat: живой счётчик медиакеша + сверка + DESC-проход докачки/вычистки

v2
evgeny 3 weeks ago
parent
commit
07ef7d9a01
  1. 309
      src/chat/chat_msg.c

309
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":"<filename>"} */
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_<ch>.timestamp).
* Сироты (файлы без сообщения, включая .chunk_* остатки) удаляются на старте.
* Лимит: при превышении storage_total_size удаляются самые старые файлы, но только
* старше storage_backfill_days (их backfill уже не докачает на следующем старте). */
* Имя файла на диске берётся из local_attrs "fp" (fallback — по имени из "d"). */
struct media_msg_file {
char rel_path[1536]; // media/<ch_id>/<final_name>
@ -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) {

Loading…
Cancel
Save