diff --git a/src/chat/chat_core.h b/src/chat/chat_core.h index f1801187..3f7407a4 100644 --- a/src/chat/chat_core.h +++ b/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); diff --git a/src/chat/chat_msg.c b/src/chat/chat_msg.c index 758e23ea..4468fe5c 100644 --- a/src/chat/chat_msg.c +++ b/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 @@ -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; +} diff --git a/src/chat/chat_setting.c b/src/chat/chat_setting.c index f8b0e391..be9c1ed9 100644 --- a/src/chat/chat_setting.c +++ b/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}, diff --git a/src/media_delivery/media_delivery.c b/src/media_delivery/media_delivery.c index e90d88ea..e2316ecf 100644 --- a/src/media_delivery/media_delivery.c +++ b/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,47 +1475,110 @@ 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; } + 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; + } + 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_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; + 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", MD_ID); + return; + } + 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; } - 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; - - uint8_t smsg[64]; size_t soff = 0; - memcpy(smsg + soff, hb.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: announce Ed25519 sign failed block=%d", MD_ID, n); - continue; - } - md_send(inst, group_id, supers[s], (const uint8_t*)&hb, sizeof(hb)); - } - } + 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; diff --git a/src/media_delivery/media_delivery.h b/src/media_delivery/media_delivery.h index 61bde270..f4645d4f 100644 --- a/src/media_delivery/media_delivery.h +++ b/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, diff --git a/src/media_delivery/media_download.c b/src/media_delivery/media_download.c index 67f6f932..5f266e29 100644 --- a/src/media_delivery/media_download.c +++ b/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) { diff --git a/src/utun_instance.c b/src/utun_instance.c index 6c2f5bd6..4bf90ec7 100644 --- a/src/utun_instance.c +++ b/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); } }