Browse Source

chat: defer media auto-download to per-channel backfill after sync

- db_sync insert-cb gets initial_sync flag (SEND_DATA=1, PUSH/local=0)
- on_msg_inserted skips media download during initial sync
- new chat_media_autodownload_backfill_channel (newest-first) triggered
  on SYNC_DONE per channel (2s debounce)
proxy
evgeny 1 week ago
parent
commit
6ee9bad17d
  1. 5
      src/chat/chat_core_priv.h
  2. 99
      src/chat/chat_msg.c
  3. 15
      src/chat/db_sync.c
  4. 6
      src/chat/db_sync.h
  5. 4
      tests/test_media_delivery_chat.c

5
src/chat/chat_core_priv.h

@ -92,6 +92,9 @@ struct chat_core_ctx {
void* media_backfill_timer;/* was chat_msg.c g_media_startup_backfill_timer */
int media_backfill_ran; /* was chat_msg.c g_media_startup_backfill_ran */
struct UTUN_INSTANCE* media_backfill_inst; /* was chat_msg.c g_media_startup_backfill_inst */
void* media_channel_backfill_timer; /* debounce per-channel докачки */
char media_channel_backfill[16][64]; /* pending ch_id */
int media_channel_backfill_pending;
};
/* Per-instance chat context. inst->chat_core обязателен после chat_core_init. */
@ -111,7 +114,7 @@ void si_register(struct UTUN_INSTANCE* inst, struct DB_SYNC_INSTANCE* si, const
/* ── db_sync callback (chat_msg.c) ── */
struct msg_insert_arg { struct UTUN_INSTANCE* inst; char ch_id[64]; };
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 on_msg_inserted(struct DB_SYNC_INSTANCE* si, uint64_t record_ts, const char* data, size_t len, uint64_t author, int initial_sync, void* arg);
/* ── применить настройку group_autoconnect ко всем CHAT-группам (chat_channel.c) ── */

99
src/chat/chat_msg.c

@ -838,7 +838,7 @@ static int md_auto_download(struct UTUN_INSTANCE* inst,
/* ─── 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 on_msg_inserted(struct DB_SYNC_INSTANCE* si, uint64_t record_ts, const char* data, size_t len, uint64_t author, int initial_sync, void* arg) {
(void)si;
struct msg_insert_arg* ia = (struct msg_insert_arg*)arg;
struct UTUN_INSTANCE* inst = ia ? ia->inst : NULL;
@ -849,8 +849,14 @@ void on_msg_inserted(struct DB_SYNC_INSTANCE* si, uint64_t record_ts, const char
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) {
/* auto-download media if message contains media metadata (live PUSH/local only).
initial-sync (SEND_DATA) media докачивается после sync через per-channel backfill. */
if (initial_sync && data && len > 0 && author != cc->my_node_id) {
char dbg[4096]; size_t dbgl = len < sizeof(dbg) - 1 ? len : sizeof(dbg) - 1;
memcpy(dbg, data, dbgl); dbg[dbgl] = '\0';
if (strstr(dbg, "\"ct\":\"") && strstr(dbg, "\"d\":\""))
DEBUG_DEBUG(DEBUG_CATEGORY_MEDIA, "%s: media download deferred (initial sync) ch=%s", CC_ID, ch_id);
} else 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';
@ -1114,14 +1120,15 @@ static int bf_cmp_ts_desc(const void* a, const void* b) {
return 0;
}
/* Собирает все медиасообщения всех каналов (presence + размеры + контекст докачки). */
static void bf_collect(struct chat_core_ctx* cc, const char* media_base, struct bf_list* l) {
/* Собирает медиасообщения (presence + размеры + контекст докачки). only_ch_id != NULL — только этот канал. */
static void bf_collect(struct chat_core_ctx* cc, const char* media_base, const char* only_ch_id, 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;
if (only_ch_id && strcmp(only_ch_id, ch_id) != 0) continue;
char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl));
char sql[256];
@ -1221,7 +1228,7 @@ void chat_media_autodownload_backfill(struct UTUN_INSTANCE* inst) {
else snprintf(media_base, sizeof(media_base), "%s", cc->db_path);
struct bf_list list = {0};
bf_collect(cc, media_base, &list);
bf_collect(cc, media_base, NULL, &list);
if (list.count > 1) qsort(list.items, (size_t)list.count, sizeof(struct bf_entry), bf_cmp_ts_desc);
/* инициализация счётчика: сумма всех present-файлов */
@ -1273,6 +1280,51 @@ void chat_media_autodownload_backfill(struct UTUN_INSTANCE* inst) {
chat_media_cache_reconcile(inst);
}
/* Докачка медиа одного канала после завершения его chat_sync.
Порядок — от последних сообщений назад (список сортируется desc по ts). */
void chat_media_autodownload_backfill_channel(struct UTUN_INSTANCE* inst, const char* ch_id) {
struct chat_core_ctx* cc = CC(inst);
if (!cc || !cc->initialized || !cc->db || !inst || !ch_id || !ch_id[0]) 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, ch_id, &list);
if (list.count > 1) qsort(list.items, (size_t)list.count, sizeof(struct bf_entry), bf_cmp_ts_desc);
uint64_t used = cc->media_cache_bytes;
int downloaded = 0, skipped = 0;
for (int i = 0; i < list.count; i++) {
struct bf_entry* e = &list.items[i];
if (e->actual_size > 0) continue;
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: channel backfill done ch=%s downloaded=%d skipped=%d",
CC_ID, ch_id, downloaded, skipped);
for (int i = 0; i < list.count; i++) if (list.items[i].body) u_free(list.items[i].body);
u_free(list.items);
}
/* ─── Кеш медиафайлов: очистка сирот + лимит storage_total_size ─── */
struct media_msg_file {
@ -1559,22 +1611,57 @@ static void media_startup_backfill_done_cb(struct DB_SYNC_INSTANCE* si, uint64_t
media_startup_backfill_timer_cb, "chat_media_bl");
}
/* per-channel: после SYNC_DONE канала — докачка его медиа (дебаунс 2s, от новых к старым) */
static void media_channel_backfill_timer_cb(void* arg) {
struct UTUN_INSTANCE* inst = (struct UTUN_INSTANCE*)arg;
struct chat_core_ctx* cc = CC(inst);
cc->media_channel_backfill_timer = NULL;
for (int i = 0; i < cc->media_channel_backfill_pending; i++)
chat_media_autodownload_backfill_channel(inst, cc->media_channel_backfill[i]);
cc->media_channel_backfill_pending = 0;
}
static void media_channel_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)peer_node_id;
if (!cc || !si) return;
char ch_id[32]; snprintf(ch_id, sizeof(ch_id), "%llu", (unsigned long long)db_sync_instance_group_id(si));
for (int i = 0; i < cc->media_channel_backfill_pending; i++)
if (strcmp(cc->media_channel_backfill[i], ch_id) == 0) return;
if (cc->media_channel_backfill_pending < (int)(sizeof(cc->media_channel_backfill) / sizeof(cc->media_channel_backfill[0])))
snprintf(cc->media_channel_backfill[cc->media_channel_backfill_pending++],
sizeof(cc->media_channel_backfill[0]), "%s", ch_id);
if (cc->media_channel_backfill_timer) uasync_cancel_timeout(inst->ua, cc->media_channel_backfill_timer);
cc->media_channel_backfill_timer = uasync_set_timeout(inst->ua, 20000, inst,
media_channel_backfill_timer_cb, "chat_md_chbl");
}
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;
cc->media_channel_backfill_pending = 0;
cc->media_channel_backfill_timer = NULL;
db_sync_add_done_cbk(inst, media_startup_backfill_done_cb, inst);
db_sync_add_done_cbk(inst, media_channel_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);
db_sync_remove_done_cbk(inst, media_channel_backfill_done_cb, inst);
if (cc->media_backfill_timer) {
uasync_cancel_timeout(inst->ua, cc->media_backfill_timer);
cc->media_backfill_timer = NULL;
}
if (cc->media_channel_backfill_timer) {
uasync_cancel_timeout(inst->ua, cc->media_channel_backfill_timer);
cc->media_channel_backfill_timer = NULL;
}
cc->media_channel_backfill_pending = 0;
cc->media_backfill_inst = NULL;
cc->media_backfill_ran = 0;
}

15
src/chat/db_sync.c

@ -237,6 +237,11 @@ struct DB_SYNC_INSTANCE* db_sync_instance_find(struct UTUN_INSTANCE* inst, uint6
return db_instance_find(inst->db_sync, group_id);
}
uint64_t db_sync_instance_group_id(struct DB_SYNC_INSTANCE* si)
{
return si ? si->group_id : 0;
}
// Выделяет новый инстанс в куче (стабильный указатель) и регистрирует его в массиве
static struct DB_SYNC_INSTANCE* db_instance_alloc(struct DB_SYNC* db)
{
@ -646,12 +651,12 @@ static int db_record_insert_cascade(struct DB_SYNC_INSTANCE* si,
uint64_t rid, uint64_t rts, uint64_t rauthor,
const char* rdata, uint32_t rdlen,
const uint8_t* rsig, uint64_t source_peer,
const char* local_attrs)
const char* local_attrs, int initial_sync)
{
int ret = db_record_insert(si, rid, rts, rauthor, rdata, rdlen, rsig, local_attrs);
if (ret != 0) return ret;
if (si->on_insert) si->on_insert(si, rts, rdata, rdlen, rauthor, si->on_insert_arg);
if (si->on_insert) si->on_insert(si, rts, rdata, rdlen, rauthor, initial_sync, si->on_insert_arg);
uint8_t pbuf[4096]; uint32_t poff = 1; pbuf[0] = DB_MSG_PUSH;
memcpy(pbuf + poff, &rid, 8); poff += 8;
@ -1046,7 +1051,7 @@ static void db_handle_send_data(struct DB_SYNC_INSTANCE* si, uint64_t src, const
for (uint16_t i = 0; i < count && i < 32; i++) {
uint64_t rid, rts, rauthor; uint32_t rdlen; const uint8_t* rdata, *rsig; int rsiglen;
if (si_parse_record(&ptr, p + len, &rid, &rts, &rauthor, &rdlen, &rdata, &rsig, &rsiglen) != 0) break;
int ret = db_record_insert_cascade(si, rid, rts, rauthor, (const char*)rdata, rdlen, rsig, src, NULL);
int ret = db_record_insert_cascade(si, rid, rts, rauthor, (const char*)rdata, rdlen, rsig, src, NULL, 1);
if (ret >= 0) received++;
if (ret == 1) duplicates++;
if (ret == -2) bad_sigs++;
@ -1178,7 +1183,7 @@ static void db_handle_push(struct DB_SYNC_INSTANCE* si, uint64_t src, const uint
if (si_parse_record(&ptr, p + len, &rid, &rts, &rauthor, &rdlen, &rdata, &rsig, &rsiglen) != 0) return;
DEBUG_TRACE(DEBUG_CATEGORY_CHAT_SYNC, "from=%016llx id=%llu author=%016llx ts=%llu len=%u", (unsigned long long)src, (unsigned long long)rid, (unsigned long long)rauthor, (unsigned long long)rts, rdlen);
int ret = db_record_insert_cascade(si, rid, rts, rauthor, (const char*)rdata, rdlen, rsig, src, NULL);
int ret = db_record_insert_cascade(si, rid, rts, rauthor, (const char*)rdata, rdlen, rsig, src, NULL, 0);
if (ret == 0) {
uint32_t ins_pos = si_find_pos(si, rts, rsig);
uint32_t total = db_count(si);
@ -1881,7 +1886,7 @@ int db_sync_insert_signed(struct DB_SYNC_INSTANCE* si, const char* json_data, si
uint64_t id = si->next_id;
uint64_t author_node_id = si->db_sync->inst->node_id;
return db_record_insert_cascade(si, id, ts, author_node_id, json_data, (uint32_t)len, sig, author_node_id, local_attrs);
return db_record_insert_cascade(si, id, ts, author_node_id, json_data, (uint32_t)len, sig, author_node_id, local_attrs, 0);
}
// Возвращает количество записей в инстансе

6
src/chat/db_sync.h

@ -99,6 +99,8 @@ struct DB_SYNC_INSTANCE* db_sync_instance_add(struct UTUN_INSTANCE* inst, const
void db_sync_instance_remove(struct DB_SYNC_INSTANCE* si);
// Найти живой инстанс по group_id (поиск в текущем массиве db->instances — указатель всегда актуален).
struct DB_SYNC_INSTANCE* db_sync_instance_find(struct UTUN_INSTANCE* inst, uint64_t group_id);
// group_id инстанса (для chat-таблиц равен числовому channel_id). Возвращает 0 при si==NULL.
uint64_t db_sync_instance_group_id(struct DB_SYNC_INSTANCE* si);
// Data operations (per-instance)
// db_sync_insert_signed: sig MUST be non-NULL, 64 bytes.
@ -122,10 +124,12 @@ int db_sync_select(struct DB_SYNC_INSTANCE* si, uint32_t offset, uint32_t l
// Insert callback: fired after local or peer insert succeeds.
// author_node_id = self for local inserts, peer node_id for remote.
// initial_sync = 1 если запись пришла в рамках bulk-обмена SEND_DATA (initial/re-sync),
// 0 для live-PUSH и локальных вставок.
typedef void (*db_sync_insert_cb)(struct DB_SYNC_INSTANCE* si,
uint64_t record_timestamp,
const char* json_data, size_t len,
uint64_t author_node_id, void* arg);
uint64_t author_node_id, int initial_sync, void* arg);
void db_sync_set_insert_cb(struct DB_SYNC_INSTANCE* si, db_sync_insert_cb cb, void* arg);
// Verify chain integrity: returns 0 if all chain_hashes are correct, 1 if any mismatch found

4
tests/test_media_delivery_chat.c

@ -283,8 +283,8 @@ static int parse_media_body(const char* body, struct media_index_result* out) {
/* ── insert-callback: storage-узлы запускают auto-download ── */
static void storage_on_insert(struct DB_SYNC_INSTANCE* si, uint64_t ts,
const char* data, size_t len, uint64_t author, void* arg) {
(void)si; (void)ts;
const char* data, size_t len, uint64_t author, int initial_sync, void* arg) {
(void)si; (void)ts; (void)initial_sync;
int idx = (int)(intptr_t)arg;
struct UTUN_INSTANCE* inst = g_inst[idx];
if (author == inst->node_id) return;

Loading…
Cancel
Save