Browse Source

chat/media: докачка недостающего медиа после стартовой синхронизации + HAVE_BLOCK анонс локальных блоков

v2
evgeny 3 weeks ago
parent
commit
a823898379
  1. 9
      src/chat/chat_core.h
  2. 139
      src/chat/chat_msg.c
  3. 1
      src/chat/chat_setting.c
  4. 109
      src/media_delivery/media_delivery.c
  5. 10
      src/media_delivery/media_delivery.h
  6. 22
      src/media_delivery/media_download.c
  7. 2
      src/utun_instance.c

9
src/chat/chat_core.h

@ -24,6 +24,15 @@ void chat_core_destroy(struct UTUN_INSTANCE* inst);
* в media_files) — вызывается после загрузки каналов/сообщений при старте. */
void chat_media_backfill(void);
/* Докачка недостающего медиа за последние storage_backfill_days дней. Проходит по
* медиасообщениям каналов и запускает загрузку для тех, чей файл отсутствует на диске. */
void chat_media_autodownload_backfill(void);
/* Регистрирует триггер: после первого SYNC_DONE (+ дебаунс 2с) анонсирует локальные
* блоки суперузлам и запускает chat_media_autodownload_backfill. Один раз за процесс. */
void chat_media_startup_backfill_init(struct UTUN_INSTANCE* inst);
void chat_media_startup_backfill_destroy(struct UTUN_INSTANCE* inst);
struct sqlite3* chat_core_get_db(void);
struct UTUN_INSTANCE* chat_core_get_inst(void);
int chat_core_is_initialized(void);

139
src/chat/chat_msg.c

@ -10,6 +10,7 @@
#include "chat_whisper.h"
#include "../utun_instance.h"
#include "../ntp_time.h"
#include "../transport_layer/secure_channel.h"
#include "../media_delivery/media_index.h"
#include "../media_delivery/media_download.h"
@ -17,6 +18,7 @@
#include "../media_async/media_async.h"
#include "../video/video.h"
#include "../../lib/mem.h"
#include "../../lib/u_async.h"
#include "../../lib/platform_compat.h"
#include <openssl/sha.h>
@ -1021,3 +1023,140 @@ void chat_media_backfill(void) {
if (total > 0)
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;
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);
int scanned = 0, present = 0, started = 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 id, timestamp, node_id, data, author_signature FROM \"%s\" "
"WHERE timestamp>=? AND node_id!=? ORDER BY timestamp ASC", 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);
uint64_t ts = (uint64_t)sqlite3_column_int64(ms, 1);
uint64_t author = (uint64_t)sqlite3_column_int64(ms, 2);
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);
if (!jdata || !sig_blob || sig_len != 64) continue;
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; }
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);
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);
md_auto_download(g_cc.inst, body, dlen, ch_id, base_filename, ts, author_sig, author, msg_id, content_type);
started++;
u_free(body);
media_index_result_free(&result);
}
sqlite3_finalize(ms);
}
sqlite3_finalize(cs);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: autodownload backfill done scanned=%d present=%d started=%d days=%d",
CC_ID, scanned, present, started, days);
}
/* ─── Триггер: после первого SYNC_DONE (+ дебаунс) — анонс локальных блоков + докачка ─── */
static void* g_media_startup_backfill_timer = NULL;
static int g_media_startup_backfill_ran = 0;
static struct UTUN_INSTANCE* g_media_startup_backfill_inst = NULL;
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) докачка недостающего */
media_delivery_announce_local_blocks(inst);
chat_media_autodownload_backfill();
}
static void media_startup_backfill_done_cb(struct DB_SYNC_INSTANCE* si, uint64_t peer_node_id, void* arg) {
(void)si; (void)peer_node_id; (void)arg;
if (g_media_startup_backfill_ran || !g_media_startup_backfill_inst) return;
struct UTUN_INSTANCE* inst = g_media_startup_backfill_inst;
if (g_media_startup_backfill_timer) uasync_cancel_timeout(inst->ua, g_media_startup_backfill_timer);
g_media_startup_backfill_timer = uasync_set_timeout(inst->ua, 20000, inst,
media_startup_backfill_timer_cb, "chat_media_bl");
}
void chat_media_startup_backfill_init(struct UTUN_INSTANCE* inst) {
if (!inst) return;
g_media_startup_backfill_inst = inst;
g_media_startup_backfill_ran = 0;
db_sync_add_done_cbk(inst, media_startup_backfill_done_cb, NULL);
}
void chat_media_startup_backfill_destroy(struct UTUN_INSTANCE* inst) {
if (!inst) return;
db_sync_remove_done_cbk(inst, media_startup_backfill_done_cb, NULL);
if (g_media_startup_backfill_timer) {
uasync_cancel_timeout(inst->ua, g_media_startup_backfill_timer);
g_media_startup_backfill_timer = NULL;
}
g_media_startup_backfill_inst = NULL;
g_media_startup_backfill_ran = 0;
}

1
src/chat/chat_setting.c

@ -47,6 +47,7 @@ static const struct chat_setting_def g_setting_defs[] = {
{"compressor_max_gain_db", CHAT_SETTING_INT, 25, 1, 100},
{"compressor_rise_rate", CHAT_SETTING_INT, 10, 1, 100},
{"media_download_max_peers", CHAT_SETTING_INT, 3, 1, 16},
{"storage_backfill_days", CHAT_SETTING_INT, 7, 1, 365},
{"standby_active_sec", CHAT_SETTING_INT, 2, 1, 3600},
{"standby_sleep_sec", CHAT_SETTING_INT, 60, 1, 86400},
{"standby_min_sleep_sec", CHAT_SETTING_INT, 15, 1, 86400},

109
src/media_delivery/media_delivery.c

@ -1463,19 +1463,9 @@ int media_delivery_bind(struct UTUN_INSTANCE* inst) {
return media_delivery_init(inst);
}
void media_delivery_announce_media(struct UTUN_INSTANCE* inst, uint64_t group_id,
const uint8_t* media_id, const uint8_t* block_ids,
int num_blocks) {
if (!inst || !media_id || !block_ids || num_blocks <= 0) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: announce invalid args inst=%p media=%p blocks=%p n=%d",
MD_ID, (void*)inst, (void*)media_id, (void*)block_ids, num_blocks);
return;
}
struct media_delivery_ctx* md = &inst->md;
if (!md->initialized) return;
/* собрать суперузлы всех CHAT-групп (node_type=4, не сам) — как md_dl_collect_supernodes */
uint64_t supers[10]; int nsup = 0;
/* Собирает уникальных суперузлов (node_type=4, не сам) всех CHAT-групп. */
static int md_collect_supernodes(struct UTUN_INSTANCE* inst, uint64_t supers[10]) {
int nsup = 0;
struct ll_entry* gle = inst->topo_groups ? inst->topo_groups->group_list->head : NULL;
while (gle) {
struct TOPO_GROUP* g = (struct TOPO_GROUP*)gle;
@ -1485,30 +1475,35 @@ void media_delivery_announce_media(struct UTUN_INSTANCE* inst, uint64_t group_id
g->channel_id, (unsigned long long)inst->node_id);
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(inst->topo_sqlite_db, sql, -1, &st, NULL) == SQLITE_OK) {
while (sqlite3_step(st) == SQLITE_ROW && nsup < 10)
supers[nsup++] = (uint64_t)sqlite3_column_int64(st, 0);
while (sqlite3_step(st) == SQLITE_ROW && nsup < 10) {
uint64_t nid = (uint64_t)sqlite3_column_int64(st, 0);
int dup = 0;
for (int i = 0; i < nsup; i++) if (supers[i] == nid) { dup = 1; break; }
if (!dup) supers[nsup++] = nid;
}
sqlite3_finalize(st);
}
}
gle = gle->next;
}
if (nsup == 0) {
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: announce — no supernodes for group 0x%016llx",
MD_ID, (unsigned long long)group_id);
return nsup;
}
void media_delivery_send_have_block(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t dst_node_id,
const uint8_t* media_id, const uint8_t* block_id, uint32_t chunk) {
if (!inst || !media_id || !block_id) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: send_have_block invalid args inst=%p media=%p block=%p",
MD_ID, (void*)inst, (void*)media_id, (void*)block_id);
return;
}
int64_t ts = (int64_t)time(NULL);
for (int s = 0; s < nsup; s++) {
for (int n = 0; n < num_blocks; n++) {
struct media_pkt_have_block hb;
memset(&hb, 0, sizeof(hb));
hb.subcmd = MEDIA_SUBCMD_HAVE_BLOCK;
hb.group_id = group_id;
memcpy(hb.media_id, media_id, 16);
memcpy(hb.block_id, block_ids + (size_t)n * 16, 16);
hb.chunk = (uint32_t)n;
hb.timestamp = ts;
memcpy(hb.block_id, block_id, 16);
hb.chunk = chunk;
hb.timestamp = (int64_t)time(NULL);
uint8_t smsg[64]; size_t soff = 0;
memcpy(smsg + soff, hb.block_id, 16); soff += 16;
@ -1516,16 +1511,74 @@ void media_delivery_announce_media(struct UTUN_INSTANCE* inst, uint64_t group_id
memcpy(smsg + soff, &hb.timestamp, 8); soff += 8;
uint64_t self = inst->node_id; memcpy(smsg + soff, &self, 8); soff += 8;
if (sc_ed25519_sign(inst->my_ed25519_privkey, smsg, soff, hb.node_sign) != SC_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: announce Ed25519 sign failed block=%d", MD_ID, n);
continue;
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: Ed25519 sign failed for HAVE_BLOCK", MD_ID);
return;
}
md_send(inst, group_id, supers[s], (const uint8_t*)&hb, sizeof(hb));
md_send(inst, group_id, dst_node_id, (const uint8_t*)&hb, sizeof(hb));
}
void media_delivery_announce_media(struct UTUN_INSTANCE* inst, uint64_t group_id,
const uint8_t* media_id, const uint8_t* block_ids,
int num_blocks) {
if (!inst || !media_id || !block_ids || num_blocks <= 0) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: announce invalid args inst=%p media=%p blocks=%p n=%d",
MD_ID, (void*)inst, (void*)media_id, (void*)block_ids, num_blocks);
return;
}
struct media_delivery_ctx* md = &inst->md;
if (!md->initialized) return;
uint64_t supers[10]; int nsup = md_collect_supernodes(inst, supers);
if (nsup == 0) {
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: announce — no supernodes for group 0x%016llx",
MD_ID, (unsigned long long)group_id);
return;
}
for (int s = 0; s < nsup; s++)
for (int n = 0; n < num_blocks; n++)
media_delivery_send_have_block(inst, group_id, supers[s], media_id,
block_ids + (size_t)n * 16, (uint32_t)n);
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: announce media=%02x%02x... blocks=%d supers=%d",
MD_ID, media_id[0], media_id[1], num_blocks, nsup);
}
void media_delivery_announce_local_blocks(struct UTUN_INSTANCE* inst) {
if (!inst) return;
struct media_delivery_ctx* md = &inst->md;
if (!md->initialized) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: announce_local_blocks — media_delivery not initialized", MD_ID); return; }
sqlite3* db = inst->topo_sqlite_db;
if (!db) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: announce_local_blocks — no DB", MD_ID); return; }
uint64_t supers[10]; int nsup = md_collect_supernodes(inst, supers);
if (nsup == 0) { DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: announce_local_blocks — no supernodes", MD_ID); return; }
sqlite3_stmt* st = NULL;
if (sqlite3_prepare_v2(db, "SELECT media_id, block_id, chunk, chat_id FROM media_files WHERE node_id=?",
-1, &st, NULL) != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: announce_local_blocks prepare failed: %s", MD_ID, sqlite3_errmsg(db));
return;
}
sqlite3_bind_int64(st, 1, (sqlite3_int64)inst->node_id);
int announced = 0;
while (sqlite3_step(st) == SQLITE_ROW) {
const uint8_t* media_id = (const uint8_t*)sqlite3_column_blob(st, 0);
const uint8_t* block_id = (const uint8_t*)sqlite3_column_blob(st, 1);
uint32_t chunk = (uint32_t)sqlite3_column_int(st, 2);
const char* chat_id = (const char*)sqlite3_column_text(st, 3);
if (!media_id || !block_id || !chat_id || !chat_id[0]) continue;
uint64_t gid = strtoull(chat_id, NULL, 10);
if (gid == 0) continue;
for (int s = 0; s < nsup; s++)
media_delivery_send_have_block(inst, gid, supers[s], media_id, block_id, chunk);
announced++;
}
sqlite3_finalize(st);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: announce_local_blocks done, blocks=%d supers=%d",
MD_ID, announced, nsup);
}
void media_delivery_set_supernode(struct UTUN_INSTANCE* inst, int is_supernode) {
if (!inst) return;
struct media_delivery_ctx* md = &inst->md;

10
src/media_delivery/media_delivery.h

@ -136,6 +136,16 @@ void media_delivery_announce_media(struct UTUN_INSTANCE* inst, uint64_t group_id
const uint8_t* media_id, const uint8_t* block_ids,
int num_blocks);
/* Отправить один HAVE_BLOCK суперузлу (подпись Ed25519 + md_send). Общий helper
* для announce_media / announce_local_blocks / media_download. */
void media_delivery_send_have_block(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t dst_node_id,
const uint8_t* media_id, const uint8_t* block_id, uint32_t chunk);
/* Переанонсировать суперузлам ВСЕ локальные блоки (media_files WHERE node_id=self).
* Вызывается при старте после синхронизации, чтобы суперузел восстановил свой
* block_availability для уже скачанных ранее медиа. Идемпотентно на суперузле. */
void media_delivery_announce_local_blocks(struct UTUN_INSTANCE* inst);
/* relay block context helpers (used by media_download.c) */
struct relay_block_ctx* md_relay_find(struct media_delivery_ctx* md, const uint8_t* block_id);
struct relay_block_ctx* md_relay_add(struct media_delivery_ctx* md,

22
src/media_delivery/media_download.c

@ -132,28 +132,10 @@ static void md_dl_send_have_block(struct UTUN_INSTANCE* inst, struct media_downl
uint64_t super = dl->super_count > 0 ? dl->super_nodes[dl->super_current] : 0;
if (!super) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: HAVE_BLOCK but no supernode available", MDL_ID); return; }
struct media_pkt_have_block hb;
memset(&hb, 0, sizeof(hb));
hb.subcmd = MEDIA_SUBCMD_HAVE_BLOCK;
hb.group_id = dl->group_id;
memcpy(hb.media_id, dl->media_id, 16);
memcpy(hb.block_id, dl->blocks[bi].block_id, 16);
hb.chunk = (uint32_t)bi;
hb.timestamp = (int64_t)time(NULL);
uint8_t smsg[64]; size_t soff = 0;
memcpy(smsg + soff, dl->blocks[bi].block_id, 16); soff += 16;
memcpy(smsg + soff, &hb.chunk, 4); soff += 4;
memcpy(smsg + soff, &hb.timestamp, 8); soff += 8;
uint64_t self = inst->node_id; memcpy(smsg + soff, &self, 8); soff += 8;
if (sc_ed25519_sign(inst->my_ed25519_privkey, smsg, soff, hb.node_sign) != SC_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: Ed25519 sign failed for HAVE_BLOCK", MDL_ID);
return;
}
DEBUG_INFO(DEBUG_CATEGORY_MEDIA, "%s: HAVE_BLOCK to super 0x%016llx block=%d",
MDL_ID, (unsigned long long)super, bi);
md_dl_send(inst, dl->group_id, super, (const uint8_t*)&hb, sizeof(hb));
media_delivery_send_have_block(inst, dl->group_id, super, dl->media_id,
dl->blocks[bi].block_id, (uint32_t)bi);
}
static void md_dl_send_block_processing(struct UTUN_INSTANCE* inst, struct media_download* dl, int bi) {

2
src/utun_instance.c

@ -478,6 +478,7 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) {
/* Phase H: chat */
chat_headless_control_destroy();
chat_media_startup_backfill_destroy(instance);
chat_sync_destroy(instance);
chat_core_destroy(instance);
DEBUG_INFO(DEBUG_CATEGORY_SYS, "[DESTROY] H done — chat");
@ -663,6 +664,7 @@ int utun_instance_init(struct UTUN_INSTANCE *instance) {
if (chat_core_init(instance, db_file) == 0) {
chat_sync_init(instance);
chat_media_backfill();
chat_media_startup_backfill_init(instance);
}
}

Loading…
Cancel
Save