diff --git a/AGENTS.md b/AGENTS.md index fbba7ca4..a9b467bd 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -635,6 +635,7 @@ void lottie_animation_destroy(Lottie_Animation *anim); Android-версия чатгуи — P2P чат на STCP (TCP), UI на Jetpack Compose (Kotlin). C-ядро: `libutun_lite` (выборочная компиляция нужных .c из `lib/` и `src/`). Сборка: CMake (headless, Linux) + Gradle/NDK (Android APK). +- 'fw' - собрать и обновить chatgui-andriod на телефоне **Подробная инструкция:** `tools/chatgui-android/AGENTS.md` **Chat-модули:** все файлы из `src/chat/` (описаны выше в секции «Chat») компилируются в `libutun_lite`. diff --git a/src/chat/chat_core.c b/src/chat/chat_core.c index ae935155..dcdcbf25 100644 --- a/src/chat/chat_core.c +++ b/src/chat/chat_core.c @@ -166,6 +166,17 @@ int chat_core_init(struct UTUN_INSTANCE* inst, const char* db_path) { " value TEXT" ");" ); + /* store media_base for block_req handler */ + { + char mbase[512]; + const char* ls = strrchr(db_path, '/'); + if (ls) snprintf(mbase, sizeof(mbase), "%.*s", (int)(ls - db_path), db_path); + else snprintf(mbase, sizeof(mbase), "%s", db_path); + char sql[1024]; snprintf(sql, sizeof(sql), + "INSERT OR REPLACE INTO ui_state(key,value) VALUES('media_base','%s')", mbase); + sqlite3_exec(g_cc.db, sql, NULL, NULL, NULL); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: media_base=%s", CC_ID, mbase); + } /* seed stub accounts if empty */ { diff --git a/src/chat/chat_event.h b/src/chat/chat_event.h index 742b1dfb..a8785964 100644 --- a/src/chat/chat_event.h +++ b/src/chat/chat_event.h @@ -33,6 +33,7 @@ extern "C" { #define CHAT_EVT_SERVICE_STARTED 12 /* data: none */ #define CHAT_EVT_SERVICE_STOPPED 13 /* data: none */ #define CHAT_EVT_ATTACHMENT_DOWNLOADED 14 /* [ch_id_len:1][ch_id:var][msg_id:8] */ +#define CHAT_EVT_DOWNLOAD_PROGRESS 15 /* [ch_id_len:1][ch_id:var][msg_id:8][blocks_done:4][num_blocks:4] */ typedef void (*chat_event_handler_fn)(int type, const uint8_t* data, int len); diff --git a/src/chat/chat_msg.c b/src/chat/chat_msg.c index 82e5548a..ead6dae5 100644 --- a/src/chat/chat_msg.c +++ b/src/chat/chat_msg.c @@ -397,27 +397,35 @@ int chat_core_load_nodeinfo(uint64_t node_id, uint8_t* buf, size_t buf_size, struct md_done_ctx { char channel_id[64]; uint64_t ts; + int64_t msg_id; uint8_t author_sig[64]; + char dest_relpath[512]; }; +static void md_download_progress_cb(void* arg, int blocks_done, int num_blocks) { + struct md_done_ctx* ctx = (struct md_done_ctx*)arg; + uint8_t evt[256]; + uint8_t cl = (uint8_t)strlen(ctx->channel_id); + evt[0] = cl; memcpy(evt + 1, ctx->channel_id, cl); + memcpy(evt + 1 + cl, &ctx->msg_id, 8); + uint32_t bd = (uint32_t)blocks_done, nb = (uint32_t)num_blocks; + memcpy(evt + 1 + cl + 8, &bd, 4); memcpy(evt + 1 + cl + 12, &nb, 4); + chat_event_post(CHAT_EVT_DOWNLOAD_PROGRESS, evt, 1 + cl + 16); +} + static void md_download_done_cb(void* arg, int err) { struct md_done_ctx* ctx = (struct md_done_ctx*)arg; if (!err) { - char dest_relpath[1024]; - { - const char* last_slash = strrchr(g_cc.db_path, '/'); - const char* base = last_slash ? last_slash + 1 : g_cc.db_path; - snprintf(dest_relpath, sizeof(dest_relpath), "media/%s", base); - } char attrs[1024]; - snprintf(attrs, sizeof(attrs), "{\"st\":\"fl\",\"fp\":\"%s\"}", dest_relpath); + snprintf(attrs, sizeof(attrs), "{\"st\":\"fl\",\"fp\":\"%s\"}", ctx->dest_relpath); chat_core_update_local_attrs(ctx->channel_id, ctx->ts, ctx->author_sig, attrs); DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: media downloaded ch=%s", CC_ID, ctx->channel_id); } else { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: media download failed ch=%s err=%d", CC_ID, ctx->channel_id, err); } - uint8_t evt[73]; uint8_t cl = (uint8_t)strlen(ctx->channel_id); - evt[0] = cl; memcpy(evt + 1, ctx->channel_id, cl); { uint64_t z = 0; memcpy(evt + 1 + cl, &z, 8); } + uint8_t evt[80]; uint8_t cl = (uint8_t)strlen(ctx->channel_id); + evt[0] = cl; memcpy(evt + 1, ctx->channel_id, cl); + memcpy(evt + 1 + cl, &ctx->msg_id, 8); chat_event_post(CHAT_EVT_ATTACHMENT_DOWNLOADED, evt, 1 + cl + 8); u_free(ctx); } @@ -425,7 +433,8 @@ static void md_download_done_cb(void* arg, int err) { static int md_start_download(struct UTUN_INSTANCE* inst, const char* data_str, size_t data_len, const char* ch_id, const char* base_filename, - uint64_t ts, const uint8_t* author_sig) { + uint64_t ts, const uint8_t* author_sig, + uint64_t author_node_id, int64_t msg_id) { (void)data_len; const char* p = data_str; while (*p && *p != '|') p++; @@ -481,15 +490,18 @@ static int md_start_download(struct UTUN_INSTANCE* inst, } snprintf(dest, sizeof(dest), "%s/media/%s/%s", media_base, ch_id, base_filename); + const char* fp_name = strrchr(dest, '/'); fp_name = fp_name ? fp_name + 1 : dest; + struct md_done_ctx* ctx = u_calloc(1, sizeof(*ctx)); if (!ctx) { media_index_result_free(&result); return -1; } snprintf(ctx->channel_id, sizeof(ctx->channel_id), "%s", ch_id); - ctx->ts = ts; + ctx->ts = ts; ctx->msg_id = msg_id; + snprintf(ctx->dest_relpath, sizeof(ctx->dest_relpath), "%s", fp_name); if (author_sig) memcpy(ctx->author_sig, author_sig, 64); DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: download start ch=%s media=%02x%02x... blocks=%d size=%lld dest=%s", CC_ID, ch_id, result.media_id[0], result.media_id[1], nb, (long long)fsize, dest); - media_download_start(inst, 0, &result, dest, media_base, md_download_done_cb, ctx); + media_download_start(inst, 0, &result, dest, media_base, author_node_id, md_download_done_cb, ctx, md_download_progress_cb, ctx); media_index_result_free(&result); return 0; } @@ -497,7 +509,8 @@ static int md_start_download(struct UTUN_INSTANCE* inst, static void md_auto_download(struct UTUN_INSTANCE* inst, const char* data_str, size_t data_len, const char* ch_id, const char* base_filename, - uint64_t ts, const uint8_t* author_sig) { + uint64_t ts, const uint8_t* author_sig, + uint64_t author, int64_t msg_id) { if (!chat_setting_get_int("storage_autoload", 1)) return; int max_mb = chat_setting_get_int("storage_autoload_maxsize_mb", 10); @@ -510,7 +523,7 @@ static void md_auto_download(struct UTUN_INSTANCE* inst, CC_ID, fsize, max_mb); return; } - md_start_download(inst, data_str, data_len, ch_id, base_filename, ts, author_sig); + md_start_download(inst, data_str, data_len, ch_id, base_filename, ts, author_sig, author, msg_id); } /* ─── db_sync callback ─── */ @@ -549,19 +562,21 @@ void on_msg_inserted(struct DB_SYNC_INSTANCE* si, uint64_t record_ts, const char Instead, read it from the DB using the just-inserted record */ char tbl[80]; msg_table_name(ch_id, tbl, sizeof(tbl)); char sql[256]; - snprintf(sql, sizeof(sql), "SELECT author_signature FROM \"%s\" WHERE timestamp=? AND node_id=? LIMIT 1", tbl); + snprintf(sql, sizeof(sql), "SELECT author_signature, id FROM \"%s\" WHERE timestamp=? AND node_id=? LIMIT 1", tbl); sqlite3_stmt* st = NULL; if (sqlite3_prepare_v2(g_cc.db, sql, -1, &st, NULL) == SQLITE_OK) { sqlite3_bind_int64(st, 1, (sqlite3_int64)record_ts); sqlite3_bind_int64(st, 2, (sqlite3_int64)author); + int64_t msg_id = 0; if (sqlite3_step(st) == SQLITE_ROW) { const void* sblob = sqlite3_column_blob(st, 0); int sblen = sqlite3_column_bytes(st, 0); if (sblob && sblen == 64) memcpy(author_sig, sblob, 64); + msg_id = sqlite3_column_int64(st, 1); } sqlite3_finalize(st); + if (msg_id > 0) md_auto_download(g_cc.inst, body, bi, ch_id, base_filename, record_ts, author_sig, author, msg_id); } - md_auto_download(g_cc.inst, body, bi, ch_id, base_filename, record_ts, author_sig); } } @@ -580,7 +595,7 @@ void chat_core_attachment_download(const char* channel_id, int64_t msg_id) { if (!channel_id || msg_id <= 0) return; char tbl[80]; msg_table_name(channel_id, tbl, sizeof(tbl)); char sql[256]; - snprintf(sql, sizeof(sql), "SELECT data, timestamp, author_signature FROM \"%s\" WHERE id=?", tbl); + snprintf(sql, sizeof(sql), "SELECT data, timestamp, author_signature, node_id FROM \"%s\" WHERE id=?", tbl); sqlite3_stmt* st = NULL; if (sqlite3_prepare_v2(g_cc.db, sql, -1, &st, NULL) != SQLITE_OK) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: download sql prep failed ch=%s", CC_ID, channel_id); @@ -597,6 +612,7 @@ void chat_core_attachment_download(const char* channel_id, int64_t msg_id) { int64_t db_ts = sqlite3_column_int64(st, 1); const uint8_t* sig_blob = (const uint8_t*)sqlite3_column_blob(st, 2); int sig_len = sqlite3_column_bytes(st, 2); + int64_t author_node_id = sqlite3_column_int64(st, 3); if (!jdata || jlen <= 0 || !sig_blob || sig_len != 64) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: download bad msg data ch=%s id=%lld", CC_ID, channel_id, (long long)msg_id); sqlite3_finalize(st); @@ -621,7 +637,7 @@ void chat_core_attachment_download(const char* channel_id, int64_t msg_id) { DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: manual download ch=%s id=%lld file=%s", CC_ID, channel_id, (long long)msg_id, base_filename); - md_start_download(g_cc.inst, body, bi, channel_id, base_filename, (uint64_t)db_ts, author_sig); + md_start_download(g_cc.inst, body, bi, channel_id, base_filename, (uint64_t)db_ts, author_sig, (uint64_t)author_node_id, msg_id); } static void chat_core_attachment_download_trampoline_impl(void* arg) { diff --git a/src/chat/db_sync_arch.md b/src/chat/db_sync_arch.md new file mode 100644 index 00000000..6f0cf48d --- /dev/null +++ b/src/chat/db_sync_arch.md @@ -0,0 +1,530 @@ +# db_sync — децентрализованная реплицируемая таблица с цепным хешированием + +## 1. Архитектура + +### 1.1 Обзор + +``` + UTUN_INSTANCE + │ + ┌────▼────┐ + │ DB_SYNC │ (глобальный контекст, один на инстанс) + └────┬────┘ + │ + ┌────────────────┼────────────────┐ + ▼ ▼ ▼ + ┌──────────────┐ ┌──────────────┐ ┌──────────────┐ + │DB_SYNC │ │DB_SYNC │ │DB_SYNC │ + │INSTANCE │ │INSTANCE │ │INSTANCE ... │ + │msg_GENERAL │ │msg_SECRET │ │ │ + └──────┬───────┘ └──────┬───────┘ └──────────────┘ + │ │ + ┌────────┼────────┐ │ + ▼ ▼ ▼ ▼ + ┌────────┐┌────────┐┌──────┐ + │SI_PEER ││SI_PEER ││SI_...│ (per-instance peer state) + │node_A ││node_B ││ │ + └────────┘└────────┘└──────┘ +``` + +**DB_SYNC** — глобальный модуль (struct, один на UTUN_INSTANCE): +- `inst` — указатель на UTUN_INSTANCE +- `db` — SQLite-соединение (своё или shared `inst->topo_sqlite_db`) +- `instances[]` — динамический массив DB_SYNC_INSTANCE +- `peer_check_timer` — периодическая проверка (5с), ищет несинхронизированных пиров +- `done_cbks` — связанный список коллбэков завершения синхронизации + +**DB_SYNC_INSTANCE** — одна таблица синхронизации: +- `table_name` — имя SQLite-таблицы (например `msg_GENERAL`) +- `hash` — instance_hash = первые 8 байт SHA256(name || id_be) для маршрутизации сообщений +- `next_id` — автоинкрементный ID записей +- `last_timestamp_ms` — последний использованный timestamp (монотонный) +- `peers[]` — динамический массив SI_PEER (состояние синхронизации с каждым пиром) +- `recalc` — отложенный пересчёт цепных хешей (асинхронно, батчами по 50) +- `ttl_timer` — периодическая TTL-очистка (каждый час) +- `on_insert` / `on_insert_arg` — коллбэк при вставке записи + +**SI_PEER** — состояние пира в контексте инстанса: +- `node_id` — идентификатор пира +- `synced_pos` — позиция (0-based) до которой данные синхронизированы +- `verified_pos` — последняя подтверждённая позиция (совпадение хешей) +- `last_peer_count` — последнее известное количество записей у пира +- `sync_state` — 0=idle, 1=syncing, 2=synced +- `sync_start_tb` — время начала синхронизации (для детекции таймаутов) +- `retry_count` — счётчик ретраев (зарезервировано, не используется) + +### 1.2 Структура данных (SQLite) + +Каждый инстанс — это SQLite-таблица: + +```sql +CREATE TABLE "msg_GENERAL" ( + timestamp INTEGER NOT NULL, -- мс, монотонный (NTP или local) + node_id INTEGER NOT NULL, -- author node_id (8 байт) + id INTEGER NOT NULL, -- автоинкремент внутри инстанса + chain_hash BLOB NOT NULL, -- SHA256(prev_hash || id || ts || author || sig) + flags INTEGER DEFAULT 0, -- DB_REC_FLAG_WAS_SENT (0x01) + data BLOB, -- JSON-данные записи + author_signature BLOB NOT NULL, -- Ed25519(ts || data), 64 байта + local_attrs TEXT DEFAULT '', -- локальные атрибуты (не реплицируются) + delivered_peers INTEGER DEFAULT 0, -- счётчик пиров, принявших PUSH + delivery_chain TEXT DEFAULT '', -- цепочка peer_id hex через запятую + PRIMARY KEY (timestamp, author_signature) +); +CREATE INDEX idx_msg_GENERAL_ttl ON msg_GENERAL (node_id, timestamp); +``` + +**Цепной хеш:** +``` +chain_hash[pos] = SHA256(chain_hash[pos-1] || id || timestamp || author || author_signature) +chain_hash[0] = SHA256(0x00..00 [32] || id_0 || ts_0 || author_0 || sig_0) +``` + +Упорядочение: `ORDER BY timestamp, author_signature`. + +### 1.3 Ключевые механизмы + +| Механизм | Описание | +|----------|----------| +| **Цепной хеш** | Каждая запись ссылается на хеш предыдущей — защита от вставки/удаления в середине цепочки | +| **Ed25519-подписи** | Каждая запись подписана автором. Проверяется при вставке (и локальной, и от пира) | +| **Async recalc** | Пересчёт `chain_hash` батчами по 50 записей (через `uasync_set_timeout`) — не блокирует event loop | +| **Master/Slave** | Узел с бОльшим `node_id` — мастер (шлёт `INIT_SYNC`). Узел с меньшим — слейв (шлёт `REQUEST_SYNC`) | +| **Периодическая проверка** | Каждые 5 секунд — поиск пиров с sync_state=0 и запуск синхронизации, детекция таймаутов (15с) | +| **TTL-очистка** | Каждый час — удаление неподтверждённых (`flags&1==0`) записей локального автора старше `db_sync_ttl` | +| **Lazy-register** | При получении сообщения с неизвестным instance_hash — поиск таблицы `msg_` с совпадающим хешем | +| **Push** | Новая локальная запись немедленно рассылается всем пирам в sync_state >= 1 | +| **Delivery tracking** | `delivered_peers` / `delivery_chain` — отслеживание доставки PUSH каждому пиру | + +### 1.4 Интеграция с ETCP + +- **Service ID:** `ETCP_RT_ID_DB_SYNC = 0x20` +- **Биндинг:** `etcp_bind(inst, 0x20, db_sync_recv_cb)` — приём сообщений +- **Connection callbacks:** `etcp_add_conn_status_cbk(inst, db_sync_on_conn_status, NULL)` — отслеживание UP/DOWN соединений +- **Отправка:** через `etcp_send()` в прямое P2P-соединение +- **Маршрутизация:** каждый пакет начинается с `[svc:1][hash_be:8]`, hash идентифицирует инстанс на приёмной стороне + +--- + +## 2. Диаграмма состояний синхронизации + +### 2.1 Состояния пира (sync_state) + +``` + ┌─────────────────────────────────────────────┐ + │ │ + ▼ │ + ┌───────┐ INIT_SYNC/REQUEST_SYNC ┌─────────┐ │ + │ IDLE │ ──────────────────────────► │ SYNCING │ │ + │ (0) │ │ (1) │ │ + └───────┘ └────┬────┘ │ + ▲ │ │ + │ ┌─────────┼──────┘ + │ │ │ + │ ┌───────────▼──┐ ┌───▼──────────┐ + │ │ SYNC_DONE │ │ TIMEOUT+ERROR │ + │ │ matched ✓ │ │ (15s) │ + │ │ (2) │ │ reset → IDLE │ + │ └──────────────┘ └───────────────┘ + │ │ + └─────────────────────────────────────────────┘ + CONN DOWN / ERROR (DB_ERR_NOT_FOUND, DB_ERR_DISABLED) +``` + +### 2.2 Диаграмма протокола синхронизации + +``` +MASTER (node_id > peer) SLAVE (node_id < peer) +─────────────────────── ────────────────────── + +conn UP → initiate sync +│ +├─ INIT_SYNC ──────────────────────────► +│ [my_count:4] │ flush recalc +│ │ определяет tp = min(peer_c, my_c) - 1 +│ │ вычисляет ch8_at_tp +│ │ строит sparse hashes (интервалы: 1,1,1,2,2,4×8) +│ │ +│ ◄────────────────────────── INIT_RESP │ +│ [tp:4][ch8:8][has_tail:1] │ +│ [sparse_cnt:1][(pos,ch8)*] │ +│ │ +│ flush recalc │ +│ ┌─ peer ch8 == 0? ──► send ALL data │ +│ ├─ ch8 match + no_tail? ──► DONE │ +│ ├─ ch8 match + has_tail? ──► send │ +│ │ from tp+1, want_from=tp+1 │ +│ └─ ch8 mismatch? ──► находим │ +│ первую разницу через sparse │ +│ hashes → fm (first mismatch) │ +│ vp = fm-1, send from fm │ +│ │ +│ SEND_DATA ──────────────────────────► │ +│ [from:4][count:2][vp:4] │ вставляет записи, проверяет подписи +│ [want_from:4][records...] │ проверяет want_from: +│ │ want≠NONE → отправляет встречный SEND_DATA +│ │ want=NONE + ss=1 → SYNC_DONE +│ │ +│ ◄────────────────────────── SEND_DATA │ +│ (встречные данные при want→NONE) │ +│ │ +│ ┌─ want=NONE + ss=1 → SYNC_DONE ────►│ +│ │ [mc:4][ch8_last:8] │ сверяет mc==pc и ch8==pch8 +│ │ │ match → ss=2, DONE +│ │ │ mcpc → send tail +│ └─────────────────────────────────────│ +``` + +### 2.3 Временная диаграмма жизненного цикла + +``` +┌─ init ─────────────────────────────────────────────────────────────────────┐ +│ db_sync_init(inst) │ +│ ├─ открывает SQLite (shared или свой chats.db) │ +│ ├─ etcp_bind(inst, 0x20, recv_cb) │ +│ ├─ etcp_add_conn_status_cbk(conn_status) │ +│ └─ запускает peer_check_timer (5s) │ +│ │ +│ ... позже, при создании канала ... │ +│ │ +│ db_sync_instance_add(inst, "msg_GENERAL", hash, 1) │ +│ ├─ CREATE TABLE IF NOT EXISTS... │ +│ ├─ загружает next_id = MAX(id)+1 │ +│ ├─ db_verify_chain() — проверка целостности цепных хешей │ +│ ├─ запускает TTL-таймер (1 час) │ +│ └─ для всех активных соединений: si_peer_add() + initiate_sync() │ +├────────────────────────────────────────────────────────────────────────────┤ +│ │ +│ ┌─── DB_SYNC_INSTANCE активен ───────────────────────────────────────────┐│ +│ │ ││ +│ │ conn UP ──► on_conn_up() ││ +│ │ └─ для всех enabled инстансов: ││ +│ │ si_peer_add() → initiate_sync() ││ +│ │ ││ +│ │ conn DOWN ──► on_conn_down() ││ +│ │ └─ сброс sync_state для всех инстансов ││ +│ │ ││ +│ │ recv_cb() — приём сообщений ││ +│ │ ├─ извлекает instance_hash, type, payload ││ +│ │ ├─ lazy-register если инстанс не найден ││ +│ │ └─ диспетчеризация по type ││ +│ │ ││ +│ │ peer_check_timer (каждые 5с) ││ +│ │ ├─ находит пиров c sync_state=0, запускает sync ││ +│ │ └─ детектит таймауты (15s), сбрасывает sync_state ││ +│ │ ││ +│ │ ttl_timer (каждый час) ││ +│ │ └─ удаляет неподтверждённые локальные записи старше TTL ││ +│ │ ││ +│ │ recalc_timer (асинхронно, батчи по 50) ││ +│ │ └─ пересчитывает chain_hash для записей после вставки ││ +│ └─────────────────────────────────────────────────────────────────────────┘│ +│ │ +│ db_sync_instance_remove(si) │ +│ ├─ отменяет таймеры │ +│ ├─ освобождает peers[] │ +│ └─ сдвигает массив instances[] (таблица БД не удаляется) │ +│ │ +│ db_sync_destroy(inst) │ +│ ├─ etcp_unbind(0x20) │ +│ ├─ останавливает все таймеры │ +│ ├─ освобождает все инстансы + done_cbks │ +│ ├─ закрывает SQLite │ +│ └─ u_free(db) │ +└─────────────────────────────────────────────────────────────────────────────┘ +``` + +--- + +## 3. Протокол (форматы кодограмм) + +### 3.1 Общий заголовок ETCP + +Каждое сообщение db_sync инкапсулируется в ETCP-дейтаграмму: + +``` +┌─────────────────┬──────────────┬──────────────────────────────┐ +│ svc_id (1 byte) │ hash_be (8) │ payload (variable) │ +│ ETCP_RT_ID=0x20 │ instance_hash│ зависит от type │ +└─────────────────┴──────────────┴──────────────────────────────┘ +``` + +**instance_hash** = первые 8 байт `SHA256(table_name || id_be)` в big-endian. +Это позволяет приёмной стороне найти `DB_SYNC_INSTANCE` по хешу. + +### 3.2 Типы сообщений + +| #define | Value | Направление | Назначение | +|---------|-------|-------------|------------| +| `DB_MSG_INIT_SYNC` | 0x01 | Master→Slave | Инициирует синхронизацию | +| `DB_MSG_INIT_RESP` | 0x02 | Slave→Master | Ответ с хешами на точке расхождения | +| `DB_MSG_SEND_DATA` | 0x04 | ↔ | Передача батча записей | +| `DB_MSG_PUSH` | 0x05 | →Peers | Рассылка новой записи всем synced-пирам | +| `DB_MSG_ACK_PUSH` | 0x06 | →Author | Подтверждение получения PUSH | +| `DB_MSG_SYNC_DONE` | 0x07 | ↔ | Завершение синхронизации | +| `DB_MSG_ERROR` | 0x08 | ↔ | Ошибка (NOT_FOUND, DISABLED) | +| `DB_MSG_REQUEST_SYNC` | 0x09 | Slave→Master | Запрос на инициацию синхронизации | + +### 3.3 DB_MSG_INIT_SYNC (0x01) + +**Отправитель:** мастер (node_id > peer_id) +**Назначение:** сообщить пиру своё количество записей для поиска точки расхождения + +``` +┌──────┬──────────────┐ +│ 0x01 │ count (4 BE) │ +└──────┴──────────────┘ + total: 5 bytes +``` + +- `count` — количество записей в таблице мастера (uint32, big-endian) + +### 3.4 DB_MSG_INIT_RESP (0x02) + +**Отправитель:** слейв (node_id < peer_id) +**Назначение:** вернуть хеш на точке пересечения tp и разреженные хеши для бинарного поиска расхождения + +``` +┌──────┬──────────────┬──────────────┬─────────────┬─────────────┬──────────────────────────────────┐ +│ 0x02 │ tp (4 BE) │ ch8_at_tp(8) │ has_tail(1) │ sparse_cnt │ sparse_entries[sparse_cnt] │ +│ │ │ │ │ (1) │ { pos(4BE), ch8(8) } × sparse_cnt│ +└──────┴──────────────┴──────────────┴─────────────┴─────────────┴──────────────────────────────────┘ + total: 14 + 12×sparse_cnt bytes (max 206) +``` + +- `tp` — точка пересечения: `min(peer_count, my_count) - 1` (uint32, network byte order) +- `ch8_at_tp` — первые 8 байт chain_hash на позиции tp (uint64, host byte order) +- `has_tail` — 1 если у слейва есть записи после tp (uint8) +- `sparse_cnt` — количество разреженных точек (uint8) +- `sparse_entries` — массив `{pos:4, ch8:8}` от tp назад с интервалами: 1, 1, 1, 2, 2, 4, 4, 4, 4, 4, 4, 4, 4, 4, 4, 4 (макс. 16) + +**Пример sparse-интервалов при tp=10:** +``` +tp=10 → entry[0]: pos=9 (tp - 1) + entry[1]: pos=8 (tp - 2) + entry[2]: pos=7 (tp - 3) + entry[3]: pos=5 (tp - 5) + entry[4]: pos=3 (tp - 7) + entry[5]: pos=0 (tp - 11) — только если tp достаточно велико +``` + +### 3.5 DB_MSG_SEND_DATA (0x04) + +**Отправитель:** любой узел +**Назначение:** передать батч записей (до 32) и/или запросить встречные данные + +``` +┌──────┬───────────┬───────────┬───────────┬──────────────┬──────────────────────────┐ +│ 0x04 │ from(4BE) │ count(2) │ vp(4BE) │ want_from(4) │ records[count] │ +│ │ │ │ │ │ ≤32 штук │ +└──────┴───────────┴───────────┴───────────┴──────────────┴──────────────────────────┘ + header: 14 bytes + records +``` + +- `from` — начальная позиция записей (uint32, network byte order) +- `count` — количество записей в этом батче (uint16, network byte order), ≤ 32 +- `vp` — verified position (uint32, BE): последняя позиция, где хеши совпали. `0xFFFFFFFF` если нет +- `want_from` — запрос встречных данных от пира начиная с этой позиции. `0xFFFFFFFF` = `DB_WANT_FROM_NONE` (не запрашивать) +- `records` — массив записей + +#### Формат одной записи (record): + +``` +┌───────┬───────┬──────────┬──────────┬────────────┬──────────┬──────────┐ +│ id(8) │ ts(8) │ author(8)│ dlen(4BE)│ data(dlen) │ sig_len │ sig(64) │ +│ │ │ │ │ │ =64 (1) │ │ +└───────┴───────┴──────────┴──────────┴────────────┴──────────┴──────────┘ + минимальный размер: 28 + 1 + 64 = 93 bytes (при dlen=0) +``` + +- `id` — идентификатор записи (uint64, host byte order) +- `ts` — timestamp в мс (uint64, host byte order) +- `author` — node_id автора (uint64, host byte order) +- `dlen` — длина JSON-данных (uint32, network byte order) +- `data` — JSON-данные (dlen байт) +- `sig_len` — длина подписи, всегда 64 (uint8) +- `sig` — Ed25519-подпись: `Ed25519(ts[8] || data[dlen])`, 64 байта + +**Note:** поля `id`, `ts`, `author` в host byte order (LE на x86), а `dlen` в network byte order (BE). + +### 3.6 DB_MSG_PUSH (0x05) + +**Отправитель:** автор записи → все synced-пиры +**Назначение:** рассылка новой записи в реальном времени + +Формат идентичен одной record из SEND_DATA: + +``` +┌──────┬───────┬───────┬──────────┬──────────┬────────────┬──────────┬──────────┐ +│ 0x05 │ id(8) │ ts(8) │ author(8)│ dlen(4BE)│ data(dlen) │ sig_len │ sig(64) │ +│ │ │ │ │ │ │ =64 (1) │ │ +└──────┴───────┴───────┴──────────┴──────────┴────────────┴──────────┴──────────┘ + минимальный размер: 29 + 1 + 64 = 94 bytes (при dlen=0) +``` + +### 3.7 DB_MSG_ACK_PUSH (0x06) + +**Отправитель:** получатель PUSH → автору +**Назначение:** подтверждение получения PUSH-записи + +``` +┌──────┬──────────┬──────────────┐ +│ 0x06 │ ts (8) │ author (8) │ +└──────┴──────────┴──────────────┘ + total: 17 bytes +``` + +- `ts` — timestamp записи (uint64, host byte order) +- `author` — node_id автора (uint64, host byte order) + +При получении ACK_PUSH автор помечает запись флагом `DB_REC_FLAG_WAS_SENT` (0x01) и обновляет `delivery_chain`/`delivered_peers`. + +### 3.8 DB_MSG_SYNC_DONE (0x07) + +**Отправитель:** любой узел +**Назначение:** финальная сверка — сравнение количества записей и последнего chain_hash + +``` +┌──────┬─────────────┬──────────────────┐ +│ 0x07 │ count (4BE) │ chain_hash8 (8) │ +└──────┴─────────────┴──────────────────┘ + total: 13 bytes +``` + +- `count` — количество записей в таблице отправителя (uint32, network byte order) +- `chain_hash8` — первые 8 байт последнего chain_hash (uint64, host byte order) + +**Обработка:** +- `count == my_count && ch8 == my_ch8` → sync завершён, ss→2, fire `db_sync_done_cb` +- `my_count < peer_count` → запрашиваем хвост (SEND_DATA from=my_count, want=my_count) +- `my_count > peer_count` → отправляем хвост (SEND_DATA from=peer_count) +- `count == my_count && ch8 != my_ch8` → принимаем как сошедшееся (конвергенция) + +### 3.9 DB_MSG_ERROR (0x08) + +**Отправитель:** любой узел +**Назначение:** сигнализировать об ошибке (инстанс не найден или выключен) + +``` +┌──────┬──────────┐ +│ 0x08 │ code (1) │ +└──────┴──────────┘ + total: 2 bytes +``` + +Коды ошибок: +| #define | Value | Значение | +|---------|-------|----------| +| `DB_ERR_NOT_FOUND` | 0x01 | Инстанс с таким hash не найден | +| `DB_ERR_DISABLED` | 0x02 | Инстанс существует, но disabled | + +При получении ERROR сбрасывается `sync_state` пира в 0. + +### 3.10 DB_MSG_REQUEST_SYNC (0x09) + +**Отправитель:** слейв → мастер +**Назначение:** запросить инициацию синхронизации (слейв говорит мастеру «начни sync») + +``` +┌──────┐ +│ 0x09 │ +└──────┘ + total: 1 byte (без payload) +``` + +При получении REQUEST_SYNC мастер вызывает `initiate_sync()` → шлёт `INIT_SYNC`. + +--- + +## 4. Алгоритм синхронизации (по шагам) + +### 4.1 Запуск + +1. При conn UP для каждого enabled инстанса: `si_peer_add()` → `initiate_sync()` +2. `initiate_sync()`: + - Если `node_id > peer_id` (мастер): шлёт `INIT_SYNC` с `my_count` + - Если `node_id < peer_id` (слейв): шлёт `REQUEST_SYNC` + +### 4.2 Обработка INIT_SYNC (слейв) + +1. `db_sync_flush_recalc()` — синхронно завершает все отложенные пересчёты +2. Вычисляет `tp = min(peer_count, my_count) - 1` (точка пересечения) +3. Вычисляет `ch8_at_tp` — chain_hash8 на позиции tp +4. Определяет `has_tail = (my_count > tp + 1)` +5. Строит разреженные хеши: от tp назад с интервалами 1,1,1,2,2,4×8 +6. Отправляет `INIT_RESP` с `tp`, `ch8_at_tp`, `has_tail`, `sparse_hashes` +7. Устанавливает `sync_state = 1` для этого пира + +### 4.3 Обработка INIT_RESP (мастер) + +1. `db_sync_flush_recalc()` — синхронно завершает все отложенные пересчёты +2. Извлекает `tp`, `peer_ch8`, `has_tail`, `sc`, `sparse_hashes` +3. Вычисляет `my_ch8_at_tp` + +**Ветки:** +- **peer_ch8 == 0 && sc == 0** (пир пуст): шлёт **все** свои данные батчами по 32 записи. Последний батч с `want_from = my_count` (запрос встречных данных, которые будут пустыми → SYNC_DONE) +- **my_ch8 == peer_ch8** (совпадение на tp): + - `!has_tail && my_count <= tp+1`: обе стороны идентичны → сразу `SYNC_DONE` + - `has_tail`: отправляем свои записи после tp и запрашиваем записи пира (`want_from = tp+1`) +- **my_ch8 != peer_ch8** (расхождение): ищем первую различающуюся позицию `fm` через sparse-хеши. `vp = fm-1`. Отправляем данные с fm и запрашиваем данные пира с fm. + +### 4.4 Обработка SEND_DATA + +1. `db_sync_flush_recalc()` — синхронно завершает все отложенные пересчёты +2. Парсит `from`, `count`, `vp`, `want_from` +3. Для каждой записи (макс. 32): `si_parse_record()` → `db_record_insert()` + - Проверка Ed25519-подписи (если неверна — удаление записи) + - Проверка дубликата по `(timestamp, author_signature)` + - Вычисление chain_hash + - Вставка в SQLite, запуск async recalc + - Вызов `on_insert` коллбэка +4. Корректирует `synced_pos` (и соседних пиров, если запись вставлена не в конец) +5. Если `want_from != DB_WANT_FROM_NONE`: отправляет встречный SEND_DATA +6. Если `want_from == DB_WANT_FROM_NONE` и `sync_state == 1`: отправляет `SYNC_DONE` + +### 4.5 Обработка SYNC_DONE + +1. `db_sync_flush_recalc()` +2. Сравнивает `my_count`, `my_ch8_last` с `peer_count`, `peer_ch8_last` +3. **Совпало** → sync_state=2, fire `db_sync_done_cb` +4. **my < peer** → запрашиваем хвост SEND_DATA(from=my_count, want=my_count) +5. **my > peer** → отправляем хвост SEND_DATA(from=peer_count) +6. **count совпало но хеши разные** → принимаем как converged + +### 4.6 PUSH (реальное время) + +1. При локальной вставке (`db_sync_insert_signed`): + - Формирует PUSH-сообщение + - Рассылает всем пирам с `sync_state >= 1` + - Обновляет `delivery_chain` для каждого пира +2. При получении PUSH: + - Вставляет запись через `db_record_insert()` + - Шлёт `ACK_PUSH` обратно автору + - Если запись вставлена в середину (не в конец) — сбрасывает `synced_pos` всех пиров до позиции вставки-1 +3. При получении ACK_PUSH: + - Помечает запись флагом `WAS_SENT` + - Обновляет `delivery_chain` + +--- + +## 5. Безопасность + +1. **Ed25519-подписи:** каждая запись подписана `Ed25519(timestamp || json_data)`. Проверяется всегда — и при локальной вставке, и при приёме от пира. Невалидные записи отбрасываются и удаляются из БД. +2. **Цепной хеш:** SHA256-цепочка защищает от вставки/удаления записей в середине истории. При проверке (`db_sync_chain_verify`) или при обнаружении расхождения во время отправки — запускается пересчёт с позиции ошибки. +3. **Публичные ключи:** Ed25519 публичный ключ пира получается из трёх источников (по приоритету): + - Локальный `inst->my_ed25519_pubkey` (для своих записей) + - `topo_node_sqlite_get_ed25519_pubkey()` (из БД узлов) + - `conn->peer_ed25519_pubkey` (из ETCP handshake) +4. **Отсутствие удаления/редактирования:** записи не редактируются и не удаляются явно — только TTL-очистка неподтверждённых. + +--- + +## 6. Конфигурация + +```ini +[global] +db_sync_enabled = 1 # включить модуль (по умолчанию 0) +db_sync_ttl = 86400 # TTL неподтверждённых записей в секундах (по умолчанию 86400) +db_path = /var/lib/utun # путь к БД (по умолчанию /tmp/utun_db_sync) +``` diff --git a/src/media_delivery/media_delivery.c b/src/media_delivery/media_delivery.c index a2fee7a9..fa9caae1 100644 --- a/src/media_delivery/media_delivery.c +++ b/src/media_delivery/media_delivery.c @@ -324,21 +324,37 @@ static void md_handle_query(struct media_delivery_ctx* md, uint64_t from_node, struct media_pkt_query* q = (struct media_pkt_query*)data; uint16_t nb = q->num_blocks; if (nb > 64) nb = 64; - uint8_t* bid = (uint8_t*)data + MEDIA_QUERY_HDR_SIZE; + uint8_t* bids = (uint8_t*)data + MEDIA_QUERY_HDR_SIZE; uint8_t resp[2048]; int off = 0; resp[off++] = MEDIA_SUBCMD_QUERY_RESP; - int num_pos = off; off += 2; /* reserve for num_entries */ + int num_pos = off; off += 2; int count = 0; for (int i = 0; i < nb && off + 30 <= (int)sizeof(resp); i++) { - uint64_t nodes[10]; - int nf = md_ba_find_nodes(md->db, bid + i * 16, nodes, 10); + uint8_t* block_id = bids + i * 16; + uint64_t nodes[10]; int nf = 0; + + /* 1) block_availability (supernode or HAVE_BLOCK reports) */ + nf = md_ba_find_nodes(md->db, block_id, nodes, 10); + + /* 2) media_files: check if we own this block */ + if (nf == 0) { + const char* sql = "SELECT 1 FROM media_files WHERE block_id=? AND node_id=? LIMIT 1"; + sqlite3_stmt* st = NULL; + if (sqlite3_prepare_v2(md->db, sql, -1, &st, NULL) == SQLITE_OK) { + sqlite3_bind_blob(st, 1, block_id, 16, SQLITE_STATIC); + sqlite3_bind_int64(st, 2, (sqlite3_int64)md->self_node_id); + if (sqlite3_step(st) == SQLITE_ROW) { nodes[nf++] = md->self_node_id; } + sqlite3_finalize(st); + } + } + for (int j = 0; j < nf && off + 30 <= (int)sizeof(resp); j++) { - uint16_t rtt = 0; /* supernode doesn't store RTT; client will probe */ + uint16_t rtt = 0; memcpy(resp + off, &nodes[j], 8); off += 8; memcpy(resp + off, &rtt, 2); off += 2; - memcpy(resp + off, bid + i * 16, 16); off += 16; + memcpy(resp + off, block_id, 16); off += 16; uint32_t ch = q->num_blocks; memcpy(resp + off, &ch, 4); off += 4; count++; @@ -347,9 +363,9 @@ static void md_handle_query(struct media_delivery_ctx* md, uint64_t from_node, uint16_t nc = (uint16_t)count; memcpy(resp + num_pos, &nc, 2); - DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: QUERY from 0x%016llx → %d entries", - MD_ID, (unsigned long long)from_node, count); - (void)md_send(md->inst, TOPO_GROUP_UTUN, from_node, resp, (size_t)off); + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: QUERY from 0x%016llx → %d entries (super=%d)", + MD_ID, (unsigned long long)from_node, count, md->is_supernode); + md_send(md->inst, TOPO_GROUP_UTUN, from_node, resp, (size_t)off); } static void md_handle_have_block(struct media_delivery_ctx* md, uint64_t from_node, @@ -695,7 +711,7 @@ static void md_super_connect(struct media_delivery_ctx* md, uint64_t peer_node_i /* find group containing this peer */ struct ll_entry* gle = md->inst->topo_groups->group_list->head; while (gle) { - struct TOPO_GROUP* g = (struct TOPO_GROUP*)gle->data; + struct TOPO_GROUP* g = (struct TOPO_GROUP*)gle; if (g->conn_mgr) { grp = g; break; } gle = gle->next; } @@ -722,7 +738,7 @@ static int md_super_start(struct media_delivery_ctx* md) { struct ll_entry* gle = md->inst->topo_groups->group_list->head; while (gle) { - struct TOPO_GROUP* g = (struct TOPO_GROUP*)gle->data; + struct TOPO_GROUP* g = (struct TOPO_GROUP*)gle; if (g->group_type != TOPO_GROUP_TYPE_CHAT) { gle = gle->next; continue; } char sql[256]; @@ -872,12 +888,10 @@ static void md_etcp_recv(struct ETCP_CONN* conn, struct ll_entry* entry) { } break; case MEDIA_SUBCMD_QUERY: - if (md->is_supernode) md_handle_query(md, from_node, data, len); - else DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: QUERY ignored — not supernode", MD_ID); + md_handle_query(md, from_node, data, len); break; case MEDIA_SUBCMD_HAVE_BLOCK: - if (md->is_supernode) md_handle_have_block(md, from_node, data, len); - else DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: HAVE_BLOCK ignored — not supernode", MD_ID); + md_handle_have_block(md, from_node, data, len); break; case MEDIA_SUBCMD_SUPER_REPL: if (md->is_supernode) md_handle_super_repl(md, from_node, data, len); @@ -892,7 +906,7 @@ static void md_etcp_recv(struct ETCP_CONN* conn, struct ll_entry* entry) { else DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: SUPER_HELLO ignored — not supernode", MD_ID); break; case MEDIA_SUBCMD_QUERY_RESP: - media_download_handle_query_resp(inst, entry->dgram, entry->len); + media_download_handle_query_resp(inst, data, len); break; case MEDIA_SUBCMD_HAVE_BLOCK_ACK: /* handled by download module */ @@ -906,10 +920,10 @@ static void md_etcp_recv(struct ETCP_CONN* conn, struct ll_entry* entry) { media_download_handle_overloaded(inst, data, len); break; case MEDIA_SUBCMD_BLOCK_CHUNK: - media_download_handle_chunk(inst, entry->dgram, entry->len); + media_download_handle_chunk(inst, data, len); break; case MEDIA_SUBCMD_BLOCK_DONE: - media_download_handle_done(inst, entry->dgram, entry->len); + media_download_handle_done(inst, data, len); break; case MEDIA_SUBCMD_CANCEL: /* TODO: handle cancel from remote */ @@ -979,7 +993,7 @@ int media_delivery_init(struct UTUN_INSTANCE* inst) { /* determine initial supernode state and subscribe to BGP callbacks */ struct ll_entry* gle = inst->topo_groups ? inst->topo_groups->group_list->head : NULL; while (gle) { - struct TOPO_GROUP* g = (struct TOPO_GROUP*)gle->data; + struct TOPO_GROUP* g = (struct TOPO_GROUP*)gle; topo_group_add_node_cbk(g, md_on_bgp_node, md); if (g->group_type == TOPO_GROUP_TYPE_CHAT && g->channel_id[0] && inst->topo_sqlite_db) { char sql[256]; @@ -1017,7 +1031,7 @@ void media_delivery_destroy(struct UTUN_INSTANCE* inst) { /* unsubscribe from BGP callbacks */ struct ll_entry* gle = inst->topo_groups ? inst->topo_groups->group_list->head : NULL; while (gle) { - struct TOPO_GROUP* g = (struct TOPO_GROUP*)gle->data; + struct TOPO_GROUP* g = (struct TOPO_GROUP*)gle; topo_group_remove_node_cbk(g, md_on_bgp_node, md); gle = gle->next; } diff --git a/src/media_delivery/media_download.c b/src/media_delivery/media_download.c index 84521e48..f088ccdd 100644 --- a/src/media_delivery/media_download.c +++ b/src/media_delivery/media_download.c @@ -16,6 +16,7 @@ #include #include #include +#include "../media_async/media_async.h" #define MDL_ID "media_download" @@ -114,7 +115,7 @@ static void md_dl_conn_cb(int result, uint64_t node_id, void* arg) { for (int bi = 0; bi < dl->peers[pi].num_blocks; bi++) { if (!dl->peers[pi].blocks[bi].started) { dl->peers[pi].blocks[bi].started = 1; - md_dl_send_block_req((struct UTUN_INSTANCE*)dl->done_arg, node_id, dl, bi); + md_dl_send_block_req(dl->inst, node_id, dl, bi); } } } @@ -127,7 +128,7 @@ static void md_dl_collect_supernodes(struct UTUN_INSTANCE* inst, struct media_do dl->super_count = 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->data; + struct TOPO_GROUP* g = (struct TOPO_GROUP*)gle; if (g->group_type != TOPO_GROUP_TYPE_CHAT || !g->channel_id[0]) { gle = gle->next; continue; } char sql[256]; snprintf(sql, sizeof(sql), @@ -144,7 +145,12 @@ static void md_dl_collect_supernodes(struct UTUN_INSTANCE* inst, struct media_do gle = gle->next; } dl->super_current = 0; - if (dl->super_count == 0) { + if (dl->super_count == 0 && dl->author_node_id && dl->author_node_id != inst->node_id) { + dl->super_nodes[0] = dl->author_node_id; + dl->super_count = 1; + DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: no supernodes, using author 0x%016llx as fallback", + MDL_ID, (unsigned long long)dl->author_node_id); + } else if (dl->super_count == 0) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: no supernodes found", MDL_ID); } } @@ -187,14 +193,22 @@ void media_download_handle_query_resp(struct UTUN_INSTANCE* inst, int ne = r->num_entries; if (ne > 100) ne = 100; - /* find the download by media_id from first entry's block_id */ - if (ne < 1 || !inst->md.downloads) return; - const uint8_t* bid = entries[0].block_id; - struct media_download* dl = md_dl_find(inst, bid); - if (!dl) { - for (int i = 0; i < ne; i++) { dl = md_dl_find(inst, entries[i].block_id); if (dl) break; } - if (!dl) { DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: QUERY_RESP — no matching download", MDL_ID); return; } + /* find the download: QUERY_RESP entries contain block_ids; + iterate all downloads checking if any block_id matches dl->block_ids */ + if (ne < 1) return; + struct media_download* dl = NULL; + if (inst->md.downloads) { + struct ll_entry* qe = inst->md.downloads->head; + while (qe) { + struct media_download* d = (struct media_download*)qe->data; + for (int i = 0; i < ne && !dl; i++) + for (int j = 0; j < d->num_blocks; j++) + if (memcmp(entries[i].block_id, d->block_ids + j * 16, 16) == 0) { dl = d; break; } + if (dl) break; + qe = qe->next; + } } + if (!dl) { DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: QUERY_RESP — no matching download", MDL_ID); return; } dl->num_peers = 0; @@ -222,16 +236,6 @@ void media_download_handle_query_resp(struct UTUN_INSTANCE* inst, DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: QUERY_RESP entries=%d peers=%d", MDL_ID, ne, dl->num_peers); - /* connect to peers */ - struct ll_entry* gle = inst->topo_groups->group_list->head; - struct TOPO_GROUP* grp = NULL; - while (gle) { - struct TOPO_GROUP* g = (struct TOPO_GROUP*)gle->data; - if (g->conn_mgr) { grp = g; break; } - gle = gle->next; - } - if (!grp || !grp->conn_mgr) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: QUERY_RESP — no conn_mgr to connect peers", MDL_ID); return; } - /* if no peers found, try querying the author directly */ if (dl->num_peers == 0 && dl->author_node_id && dl->author_node_id != inst->node_id) { DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: no holders from supernode, querying author 0x%016llx", @@ -243,8 +247,22 @@ void media_download_handle_query_resp(struct UTUN_INSTANCE* inst, return; } + /* find a group with conn_mgr for direct connection (best-effort) */ + struct ll_entry* gle = inst->topo_groups->group_list->head; + struct TOPO_GROUP* grp = NULL; + while (gle) { + struct TOPO_GROUP* g = (struct TOPO_GROUP*)gle; + if (g->conn_mgr) { grp = g; break; } + gle = gle->next; + } + for (int pi = 0; pi < dl->num_peers; pi++) { - conn_mgr_connect_node(grp->conn_mgr, dl->peers[pi].node_id, 0, md_dl_conn_cb, dl); + dl->peers[pi].connected = 1; + if (grp && grp->conn_mgr) + conn_mgr_connect_node(grp->conn_mgr, dl->peers[pi].node_id, 0, md_dl_conn_cb, dl); + else + for (int bi = 0; bi < dl->peers[pi].num_blocks; bi++) + md_dl_send_block_req(inst, dl->peers[pi].node_id, dl, bi); } } @@ -298,36 +316,46 @@ void media_download_handle_done(struct UTUN_INSTANCE* inst, } if (bi < 0) { DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "%s: BLOCK_DONE for unknown block", MDL_ID); return; } - /* verify signature */ - EVP_PKEY* pkey = NULL; - EVP_MD_CTX* vctx = EVP_MD_CTX_new(); - int sig_ok = 0; - if (vctx) { - /* read block data from temp file */ - char tmp[2048]; - snprintf(tmp, sizeof(tmp), "%s.chunk_%d", dl->dest_path, bi); - FILE* f = fopen(tmp, "rb"); - if (!f) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: fopen(%s) failed for sig check", MDL_ID, tmp); EVP_MD_CTX_free(vctx); return; } - fseeko(f, 0, SEEK_END); - off_t fsz = ftello(f); - fseeko(f, 0, SEEK_SET); - uint8_t* buf = u_malloc((size_t)fsz); - if (!buf) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: u_malloc(%lld) failed for sig check", MDL_ID, (long long)fsz); fclose(f); EVP_MD_CTX_free(vctx); return; } - size_t rd = fread(buf, 1, (size_t)fsz, f); - uint64_t node_id = inst->node_id; - uint8_t smsg[128]; size_t soff = 0; - memcpy(smsg + soff, buf, rd); soff += rd; - memcpy(smsg + soff, &node_id, 8); soff += 8; - EVP_DigestVerifyInit(vctx, NULL, EVP_sha256(), NULL, pkey); - sig_ok = (rd == (size_t)bd->total_size && rd > 0) ? 1 : 0; - if (sig_ok && rd == (size_t)bd->total_size) { - if (memcmp(bd->block_sig, dl->block_sigs + bi * 64, 64) == 0) sig_ok = 1; - else sig_ok = 0; - } - u_free(buf); + /* verify block data: Ed25519 signature from author, or size-only fallback */ + int sig_ok = 0; + uint8_t author_ed25519_pubkey[32] = {0}; + if (dl->author_node_id && inst->topo_sqlite_db) { + sqlite3_stmt* st = NULL; + sqlite3_prepare_v2(inst->topo_sqlite_db, + "SELECT ed25519_pubkey FROM nodes WHERE node_id=?", -1, &st, NULL); + if (st) { sqlite3_bind_int64(st, 1, (sqlite3_int64)dl->author_node_id); + if (sqlite3_step(st) == SQLITE_ROW) { const void* pk = sqlite3_column_blob(st, 0); int plen = sqlite3_column_bytes(st, 0); if (pk && plen == 32) memcpy(author_ed25519_pubkey, pk, 32); } + sqlite3_finalize(st); } + } + char tmp[2048]; snprintf(tmp, sizeof(tmp), "%s.chunk_%d", dl->dest_path, bi); + FILE* f = fopen(tmp, "rb"); + if (!f) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: fopen(%s) failed", MDL_ID, tmp); } + else { + fseeko(f, 0, SEEK_END); off_t fsz = ftello(f); fseeko(f, 0, SEEK_SET); + if (fsz == (off_t)bd->total_size && fsz > 0) { + if (author_ed25519_pubkey[0]) { /* Ed25519 verify against author's pubkey */ + uint8_t* buf = u_malloc((size_t)fsz); + if (buf) { + size_t rd = fread(buf, 1, (size_t)fsz, f); + if (rd == (size_t)fsz) { + EVP_PKEY* apk = EVP_PKEY_new_raw_public_key(EVP_PKEY_ED25519, NULL, author_ed25519_pubkey, 32); + EVP_MD_CTX* vctx = apk ? EVP_MD_CTX_new() : NULL; + if (vctx) { + sig_ok = (EVP_DigestVerifyInit(vctx, NULL, NULL, NULL, apk) == 1 + && EVP_DigestVerify(vctx, dl->block_sigs + bi * 64, 64, buf, rd) == 1) ? 1 : 0; + EVP_MD_CTX_free(vctx); + } + if (apk) EVP_PKEY_free(apk); + } + u_free(buf); + } + if (!sig_ok) DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: BLOCK_DONE Ed25519 sig fail block=%d", MDL_ID, bi); + } else { + sig_ok = (memcmp(bd->block_sig, dl->block_sigs + bi * 64, 64) == 0) ? 1 : 0; + } + } else { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: BLOCK_DONE size mismatch: fsz=%lld expected=%u", MDL_ID, (long long)fsz, bd->total_size); } fclose(f); } - EVP_MD_CTX_free(vctx); if (sig_ok) { dl->blocks_received++; @@ -335,6 +363,8 @@ void media_download_handle_done(struct UTUN_INSTANCE* inst, DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: BLOCK_DONE block=%d total=%d/%d", MDL_ID, bi, dl->blocks_validated, dl->num_blocks); + if (dl->progress_cb) dl->progress_cb(dl->progress_arg, dl->blocks_validated, dl->num_blocks); + /* send HAVE_BLOCK to supernode */ md_dl_send_have_block(inst, dl, bi); @@ -358,6 +388,12 @@ void media_download_handle_done(struct UTUN_INSTANCE* inst, } fclose(out); } + /* verify content_hash (skip if all-zero — test/backwards compat) */ + { int hash_zero = 1; for (int hi = 0; hi < 32; hi++) if (dl->content_hash[hi]) { hash_zero = 0; break; } + if (!hash_zero) { + uint8_t fhash[32]; ma_sha256_file(dl->dest_path, fhash); + if (memcmp(fhash, dl->content_hash, 32) != 0) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: content_hash mismatch for %s", MDL_ID, dl->dest_path); remove(dl->dest_path); dl->blocks_validated--; return; } + } } dl->assembled = 1; dl->err = 0; if (dl->done_cb) dl->done_cb(dl->done_arg, dl->err); @@ -387,7 +423,9 @@ void media_download_handle_done(struct UTUN_INSTANCE* inst, int media_download_start(struct UTUN_INSTANCE* inst, uint64_t group_id, const struct media_index_result* result, const char* dest_path, const char* media_base, - void (*done_cb)(void* arg, int err), void* done_arg) { + uint64_t author_node_id, + void (*done_cb)(void* arg, int err), void* done_arg, + void (*progress_cb)(void* arg, int blocks_done, int num_blocks), void* progress_arg) { if (!inst || !result || !dest_path || !media_base || !done_cb) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: start with NULL args inst=%p result=%p dest=%p base=%p cb=%p", MDL_ID, (void*)inst, (void*)result, (void*)dest_path, (void*)media_base, (void*)done_cb); @@ -410,6 +448,10 @@ int media_download_start(struct UTUN_INSTANCE* inst, uint64_t group_id, dl->active = 1; dl->done_cb = done_cb; dl->done_arg = done_arg; + dl->progress_cb = progress_cb; + dl->progress_arg = progress_arg; + dl->inst = inst; + dl->author_node_id = author_node_id; dl->block_ids = u_malloc((size_t)dl->num_blocks * 16); dl->block_sigs = u_malloc((size_t)dl->num_blocks * 64); diff --git a/src/media_delivery/media_download.h b/src/media_delivery/media_download.h index 9df89830..1187b884 100644 --- a/src/media_delivery/media_download.h +++ b/src/media_delivery/media_download.h @@ -60,6 +60,9 @@ struct media_download { void* timeout_timer; void (*done_cb)(void* arg, int err); void* done_arg; + void (*progress_cb)(void* arg, int blocks_done, int num_blocks); + void* progress_arg; + struct UTUN_INSTANCE* inst; }; /* состояние одного стрима на стороне блок-холдера */ @@ -80,7 +83,9 @@ struct media_block_stream { int media_download_start(struct UTUN_INSTANCE* inst, uint64_t group_id, const struct media_index_result* result, const char* dest_path, const char* media_base, - void (*done_cb)(void* arg, int err), void* done_arg); + uint64_t author_node_id, + void (*done_cb)(void* arg, int err), void* done_arg, + void (*progress_cb)(void* arg, int blocks_done, int num_blocks), void* progress_arg); int media_download_cancel(struct UTUN_INSTANCE* inst, const uint8_t* media_id, const uint8_t* block_id, uint32_t chunk); diff --git a/src/routing_layer/conn_mgr.c b/src/routing_layer/conn_mgr.c index 11517453..5932624c 100644 --- a/src/routing_layer/conn_mgr.c +++ b/src/routing_layer/conn_mgr.c @@ -262,8 +262,12 @@ int conn_mgr_connect_node(struct CONN_MGR* mgr, uint64_t node_id, uint32_t idle_ if (!target) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "node 0x%016llx not found", (unsigned long long)node_id); return CONN_MGR_ERR_NOT_FOUND; } struct CONN_MGR_ENTRY* entry = cm_find_entry(mgr, node_id); if (entry && entry->state == CONN_MGR_STATE_CONNECTED) { - DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "node 0x%016llx already connected", (unsigned long long)node_id); - return CONN_MGR_ERR_ALREADY_CONNECTED; + if (cb) { + struct cm_cb_node* cn = u_calloc(1, sizeof(struct cm_cb_node)); + if (cn) { cn->cb = cb; cn->arg = cb_arg; entry->cb_list = cn; } + cb(CONN_MGR_OK, node_id, cb_arg); + } + return CONN_MGR_OK; } if (entry && entry->state == CONN_MGR_STATE_CONNECTING) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "node 0x%016llx already connecting, adding cb", (unsigned long long)node_id); @@ -275,7 +279,7 @@ int conn_mgr_connect_node(struct CONN_MGR* mgr, uint64_t node_id, uint32_t idle_ if (!existing) existing = instance_find_conn(mgr->instance, node_id); if (existing && existing->peer_node_id == node_id && existing->links) { struct ETCP_LINK* l = existing->links; - while (l) { if (l->link_status && l->initialized) break; l = l->next; } + while (l) { if (l->link_state == 3 && l->initialized) break; l = l->next; } if (l) { entry = cm_ensure_entry(mgr, node_id); if (!entry) return CONN_MGR_ERR_INTERNAL; diff --git a/tests/Makefile.am b/tests/Makefile.am index 393412cd..e767dc0f 100644 --- a/tests/Makefile.am +++ b/tests/Makefile.am @@ -48,6 +48,7 @@ check_PROGRAMS = \ test_bgp_route_exchange \ test_bgp_triangle \ test_conn_mgr \ + test_conn_mgr_already_connected \ test_etcp_connect \ test_db_sync \ test_merkle_sync \ @@ -282,6 +283,10 @@ test_conn_mgr_SOURCES = test_conn_mgr.c test_conn_mgr_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_conn_mgr_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) +test_conn_mgr_already_connected_SOURCES = test_conn_mgr_already_connected.c +test_conn_mgr_already_connected_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib +test_conn_mgr_already_connected_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) + test_etcp_connect_SOURCES = test_etcp_connect.c test_etcp_connect_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/lib test_etcp_connect_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS) diff --git a/tests/test_conn_mgr_already_connected.c b/tests/test_conn_mgr_already_connected.c new file mode 100644 index 00000000..48f25beb --- /dev/null +++ b/tests/test_conn_mgr_already_connected.c @@ -0,0 +1,158 @@ +/** + * @file test_conn_mgr_already_connected.c + * @brief conn_mgr: подключение к уже-соединённому узлу через instance_find_conn + * + * Проверяет, что conn_mgr_connect_node находит существующее ETCP-соединение + * и сразу вызывает callback с CONN_MGR_OK, а не уходит в фазы установки нового соединения. + */ +#include +#include +#include +#include +#include "../lib/platform_compat.h" +#include "test_utils.h" +#ifndef _WIN32 +#include +#endif + +#include "etcp.h" +#include "etcp_connections.h" +#include "../src/config_parser.h" +#include "../src/config_updater.h" +#include "../src/utun_instance.h" +#include "topo_group.h" +#include "conn_mgr.h" +#include "../lib/u_async.h" +#include "../lib/debug_config.h" +#include "../lib/mem.h" + +#define TOTAL_TIMEOUT_TB 10000 /* 1 sec */ +#define POLL_MS 5 +static struct UTUN_INSTANCE *g_a = NULL, *g_b = NULL; +static struct UASYNC *ua = NULL; +static volatile int result = 0, cb_fired = 0, cb_ok = 0; +static char tdir[] = "/tmp/utun_cm2_XXXXXX"; +static char ca[256], cb[256]; +static int pa = 0, pb = 0; + +static int wf(const char* p, const char* f, ...) { + va_list ap; FILE* fp = fopen(p, "w"); if (!fp) return -1; + va_start(ap, f); vfprintf(fp, f, ap); va_end(ap); fclose(fp); return 0; +} +static char* gv(const char* p, const char* k) { + struct utun_config* c = parse_config(p); if (!c) return NULL; + char* r = (strcmp(k, "pub") == 0) ? u_strdup(c->global.my_public_key_hex) : u_strdup(c->global.my_private_key_hex); + free_config(c); return r; +} +static void fail(const char* msg) { fprintf(stderr, "FAIL: %s\n", msg); fflush(stderr); result = 2; } +static void done(void) { if (result == 0) result = 1; } +static void test_mgr_connect(void* arg); +static void test_check(void* arg); +static void test3_preconnect(void* arg); +static void test3_check(void* arg); + +static void conn_cb(int r, uint64_t id, void* arg) { + (void)arg; + cb_fired = 1; cb_ok = (r == CONN_MGR_OK); + fprintf(stderr, "conn_cb: result=%d node=0x%llx ok=%d\n", r, (unsigned long long)id, cb_ok); fflush(stderr); +} + +static void test_find_conn(void* arg) { + (void)arg; + struct ETCP_CONN* c = instance_find_conn(g_a, g_b->node_id); + if (!c || !c->links) { uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test_find_conn, "t1w"); return; } + int has = 0; struct ETCP_LINK* l = c->links; + while (l) { if (l->link_state == 3 && l->initialized) { has = 1; break; } l = l->next; } + if (!has) { uasync_set_timeout(ua, 50, NULL, (timeout_callback_t)test_find_conn, "t1w2"); return; } + fprintf(stderr, "Test 1: instance_find_conn OK (link_state=%d)\n", l->link_state); fflush(stderr); + uasync_call_soon(ua, NULL, test_mgr_connect); +} + +static void test_mgr_connect(void* arg) { + (void)arg; + cb_fired = 0; cb_ok = 0; + struct CONN_MGR* mgr = g_a->conn_mgr; + struct TOPO_NODEQ* nq = topo_node_find_by_id(mgr->group, g_b->node_id); + if (nq) { fprintf(stderr, "node already in UTUN group\n"); fflush(stderr); } + fprintf(stderr, "Test 2: conn_mgr_connect_node(B)\n"); fflush(stderr); + conn_mgr_connect_node(mgr, g_b->node_id, 0, conn_cb, NULL); + uasync_call_soon(ua, NULL, test_check); +} + +static void test_check(void* arg) { + (void)arg; + if (!cb_fired) { fail("callback not fired"); done(); return; } + if (!cb_ok) { fail("callback result != CONN_MGR_OK"); done(); return; } + uint8_t st, ty; + conn_mgr_get_status(g_a->conn_mgr, g_b->node_id, &st, &ty); + if (st != CONN_MGR_STATE_CONNECTED) { fail("state != CONNECTED"); done(); return; } + fprintf(stderr, "Test 2: OK state=%d type=%d\n", st, ty); fflush(stderr); + fprintf(stderr, "=== ALL DONE ===\n"); fflush(stderr); + done(); +} + +static void test3_preconnect(void* arg) { + /* Test 3: conn_mgr_connect_node BEFORE ETCP link is up (no instance_find_conn match) */ + (void)arg; + struct CONN_MGR* mgr = g_a->conn_mgr; + cb_fired = 0; cb_ok = 0; + struct ETCP_CONN* c = instance_find_conn(g_a, g_b->node_id); + if (c) { fprintf(stderr, "Test3: SKIP — conn already established\n"); return; } + fprintf(stderr, "Test 3: conn_mgr_connect_node BEFORE link up\n"); fflush(stderr); + conn_mgr_connect_node(mgr, g_b->node_id, 0, conn_cb, NULL); + uasync_call_soon(ua, NULL, test3_check); +} + +static void test3_check(void* arg) { + (void)arg; + if (cb_fired) { + if (cb_ok) { fprintf(stderr, "Test 3: OK — callback CONN_MGR_OK\n"); fflush(stderr); return; } + fail("callback fired with error"); done(); return; + } + /* no callback yet — check if conn_mgr gave up (state reset) */ + uint8_t st, ty; + if (conn_mgr_get_status(g_a->conn_mgr, g_b->node_id, &st, &ty) == CONN_MGR_OK && st == CONN_MGR_STATE_DISCONNECTED) { + /* entry was cleaned up — conn not possible in this test setup with this timing */ + fprintf(stderr, "Test 3: SKIP — entry cleaned up (conn not yet established in time)\n"); fflush(stderr); + return; + } + uasync_set_timeout(ua, 30, NULL, (timeout_callback_t)test3_check, "t3w"); +} + +static void to_cb(void* arg) { (void)arg; fail("timeout"); done(); } + +int main(void) { + debug_config_init(); debug_set_level(DEBUG_LEVEL_WARN); + utun_instance_set_tun_init_enabled(0); + + test_mkdtemp(tdir); + int base = 49000 + (getpid() % 10000); pa = base; pb = base + 1; + snprintf(ca, sizeof(ca), "%s/a.conf", tdir); snprintf(cb, sizeof(cb), "%s/b.conf", tdir); + wf(ca, "[global]\ntun_ip=10.97.0.1/24\ntun_ifname=tun95\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", pa); + wf(cb, "[global]\ntun_ip=10.97.0.2/24\ntun_ifname=tun94\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", pb); + config_ensure_keys_and_node_id(ca); config_ensure_keys_and_node_id(cb); + { struct utun_config* cfa = parse_config(ca); struct utun_config* cfb = parse_config(cb); + free_config(cfa); free_config(cfb); } + char *r0 = gv(ca,"priv"), *p0 = gv(ca,"pub"), *p1 = gv(cb,"pub"); + wf(ca, "[global]\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.97.0.1/24\ntun_ifname=tun95\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[client: to_b]\nkeepalive=1\npeer_public_key=%s\nlink=s1:127.0.0.1:%d\n[allowed_keys]\nallow_all=1\n", r0, p0, pa, p1, pb); + u_free(r0); u_free(p0); u_free(p1); + + ua = uasync_create(); + g_a = utun_instance_create(ua, ca); g_b = utun_instance_create(ua, cb); + if (!g_a || !g_b) { fprintf(stderr,"create failed\n"); result=2; goto done; } + utun_instance_init(g_a); utun_instance_init(g_b); + + /* Test 3: call conn_mgr BEFORE link is up */ + uasync_call_soon(ua, NULL, test3_preconnect); + /* Test 1+2: run after link is established */ + uasync_call_soon(ua, NULL, test_find_conn); + uasync_set_timeout(ua, TOTAL_TIMEOUT_TB, NULL, to_cb, "to"); + { int el = 0; while (!result && el < TOTAL_TIMEOUT_TB / 10 + 500) { uasync_poll(ua, POLL_MS); el += POLL_MS; } } + +done: + if (g_a) { g_a->running = 0; utun_instance_destroy(g_a); } + if (g_b) { g_b->running = 0; utun_instance_destroy(g_b); } + if (ua) uasync_destroy(ua, 0); + test_unlink(ca); test_unlink(cb); test_rmdir(tdir); + return (result == 1) ? 0 : 1; +} diff --git a/tools/chatgui-android/app/src/main/java/com/utun/chat/data/AudioRecorderManager.kt b/tools/chatgui-android/app/src/main/java/com/utun/chat/data/AudioRecorderManager.kt index 07776f9f..37608adf 100644 --- a/tools/chatgui-android/app/src/main/java/com/utun/chat/data/AudioRecorderManager.kt +++ b/tools/chatgui-android/app/src/main/java/com/utun/chat/data/AudioRecorderManager.kt @@ -235,6 +235,8 @@ class AudioRecorderManager { val eMs = startMs + (readTotal * FRAME_MS / frameSamples).toInt() onProgress?.invoke(total, eMs) } + val drainMs = (bufSize * 1000L) / (sampleRate.toLong() * channels * 2) + 100 + try { Thread.sleep(drainMs) } catch (_: InterruptedException) {} track.stop(); track.release() NativeLib.voiceDecodeClose(handle) } catch (e: Exception) { diff --git a/tools/chatgui-android/app/src/main/java/com/utun/chat/ui/components/MessageBubble.kt b/tools/chatgui-android/app/src/main/java/com/utun/chat/ui/components/MessageBubble.kt index b829509a..1d1609a1 100644 --- a/tools/chatgui-android/app/src/main/java/com/utun/chat/ui/components/MessageBubble.kt +++ b/tools/chatgui-android/app/src/main/java/com/utun/chat/ui/components/MessageBubble.kt @@ -97,7 +97,7 @@ private fun VoiceBubble(message: Message, textColor: Color, isActive: Boolean, val progress = ps.progress Column { - Row(verticalAlignment = Alignment.CenterVertically) { + Row(modifier = Modifier.fillMaxWidth(), verticalAlignment = Alignment.CenterVertically) { Box( modifier = Modifier.size(32.dp).background(Color.White, CircleShape) .clickable { if (message.filePath.isNotEmpty()) onPlay(message.filePath, durMs) }, @@ -120,20 +120,22 @@ private fun VoiceBubble(message: Message, textColor: Color, isActive: Boolean, } ) { Canvas(modifier = Modifier.fillMaxSize()) { - val barW = size.width / wfCount + val barCount = wfCount + val slotW = size.width / barCount + val barW = slotW * 0.75f val ph = size.height val showCursor = playing || paused - val playedIdx = if (showCursor && progress > 0f) (progress * wfCount).toInt() else 0 - for (i in 0 until wfCount) { + val playedIdx = if (showCursor && progress > 0f) (progress * barCount).toInt() else 0 + for (i in 0 until barCount) { val level = wf[i].coerceIn(0f, 1f) val barH = (2f + level * (ph - 2f)).coerceAtMost(ph) - val x = i.toFloat() * barW + val x = i * slotW val y = ph - barH val played = showCursor && i < playedIdx drawRect( color = when { played -> Color.White; else -> Color(0xFFCCCCD4).copy(alpha = 0.8f) }, topLeft = Offset(x, y), - size = Size(barW * 0.75f, barH) + size = Size(barW, barH) ) } if (showCursor && progress > 0f && progress < 1f) { diff --git a/tools/chatgui-android/app/src/main/java/com/utun/chat/viewmodel/ChatViewModel.kt b/tools/chatgui-android/app/src/main/java/com/utun/chat/viewmodel/ChatViewModel.kt index 17b60d47..29219a8f 100644 --- a/tools/chatgui-android/app/src/main/java/com/utun/chat/viewmodel/ChatViewModel.kt +++ b/tools/chatgui-android/app/src/main/java/com/utun/chat/viewmodel/ChatViewModel.kt @@ -61,6 +61,8 @@ class ChatViewModel : ViewModel() { private var recordingStartTimeMs = 0L + @Volatile private var realPlaybackElapsedMs = 0 + init { viewModelScope.launch { AppEventHandler.events.collect { (type, data) -> handleEvent(type, data) } @@ -95,6 +97,25 @@ class ChatViewModel : ViewModel() { delay(80) } } + viewModelScope.launch { + while (isActive) { + delay(50) + val s = _playbackState.value + if (s.playing && !s.paused) { + val real = realPlaybackElapsedMs + val targetMs = if (real >= s.durMs && s.durMs > 0) s.durMs else real + val newMs = (s.elapsedMs + (targetMs - s.elapsedMs) * 0.3f).toInt().coerceIn(0, s.durMs) + val newProg = if (s.durMs > 0) newMs.toFloat() / s.durMs else 0f + + if (real >= s.durMs && s.durMs > 0 && newMs >= s.durMs - 30) { + _playbackState.value = PlaybackState() + realPlaybackElapsedMs = 0 + } else { + _playbackState.value = s.copy(elapsedMs = newMs, progress = newProg) + } + } + } + } } private fun handleEvent(type: Int, data: ByteArray?) { @@ -241,18 +262,6 @@ class ChatViewModel : ViewModel() { _recordingDurationMs.value = 0 val ch = _currentChannel.value if (ch != null && duration >= 500) { - viewModelScope.launch { - while (isActive) { - val s = _playbackState.value - if (s.playing && !s.paused) { - _playbackState.value = s.copy( - elapsedMs = (s.elapsedMs + 80).coerceIn(0, s.durMs), - progress = if (s.durMs > 0) s.elapsedMs.toFloat().coerceIn(0f, s.durMs.toFloat()) / s.durMs.toFloat() else 0f - ) - } - delay(80) - } - } viewModelScope.launch { delay(500) refreshMessages(ch.id) @@ -311,38 +320,40 @@ class ChatViewModel : ViewModel() { val cur = _playbackState.value if (cur.filePath == filePath) { if (cur.paused) { + realPlaybackElapsedMs = cur.elapsedMs audioRecorder.playFromPosition(filePath, cur.elapsedMs, - onProgress = { prog, elapsed -> - _playbackState.value = _playbackState.value.copy(elapsedMs = elapsed, progress = prog) - }, - onComplete = { - _playbackState.value = PlaybackState() - }) + onProgress = { _, elapsed -> realPlaybackElapsedMs = elapsed }, + onComplete = { realPlaybackElapsedMs = durMs }) _playbackState.value = cur.copy(playing = true, paused = false) } else { audioRecorder.pausePlayback() + realPlaybackElapsedMs = 0 _playbackState.value = cur.copy(playing = false, paused = true) } return } audioRecorder.stopPlayback() + realPlaybackElapsedMs = 0 _playbackState.value = PlaybackState(filePath, playing = true, durMs = durMs) audioRecorder.playVoiceFile(filePath, - onProgress = { prog, elapsed -> - _playbackState.value = _playbackState.value.copy(elapsedMs = elapsed, progress = prog) - }, - onComplete = { _playbackState.value = PlaybackState() }) + onProgress = { _, elapsed -> realPlaybackElapsedMs = elapsed }, + onComplete = { + realPlaybackElapsedMs = 0 + _playbackState.value = PlaybackState() + }) } fun seekVoice(filePath: String, seekMs: Int, durMs: Int) { audioRecorder.stopPlayback() + realPlaybackElapsedMs = seekMs _playbackState.value = PlaybackState(filePath, playing = true, durMs = durMs, elapsedMs = seekMs, progress = if (durMs > 0) seekMs.toFloat() / durMs.toFloat() else 0f) audioRecorder.playFromPosition(filePath, seekMs, - onProgress = { prog, elapsed -> - _playbackState.value = _playbackState.value.copy(elapsedMs = elapsed, progress = prog) - }, - onComplete = { _playbackState.value = PlaybackState() }) + onProgress = { _, elapsed -> realPlaybackElapsedMs = elapsed }, + onComplete = { + realPlaybackElapsedMs = 0 + _playbackState.value = PlaybackState() + }) } fun clearState() { diff --git a/tools/chatgui/src/mainwindow.cpp b/tools/chatgui/src/mainwindow.cpp index a718e5da..c12d8d82 100644 --- a/tools/chatgui/src/mainwindow.cpp +++ b/tools/chatgui/src/mainwindow.cpp @@ -257,12 +257,15 @@ void MainWindow::setupBridgeCallbacks() { } }); gui_bridge_set_attachment_downloaded_cb([](const char* ch_id, int ch_id_len, int64_t msg_id) { - (void)msg_id; if (!s_mainWindow) return; QString chId = QString::fromUtf8(ch_id, ch_id_len); s_mainWindow->m_channelList->loadChannels(); - if (chId == s_mainWindow->m_currentChannelId) - s_mainWindow->m_messageList->refresh(); + s_mainWindow->m_messageList->onAttachmentDownloaded(chId, msg_id); + }); + gui_bridge_set_download_progress_cb([](const char* ch_id, int ch_id_len, int64_t msg_id, int blocks_done, int num_blocks) { + if (!s_mainWindow) return; + QString chId = QString::fromUtf8(ch_id, ch_id_len); + s_mainWindow->m_messageList->updateMessageProgress(chId, msg_id, blocks_done, num_blocks); }); } diff --git a/tools/chatgui/src/messagedelegate.cpp b/tools/chatgui/src/messagedelegate.cpp index ca1c1190..aaa2cca5 100644 --- a/tools/chatgui/src/messagedelegate.cpp +++ b/tools/chatgui/src/messagedelegate.cpp @@ -18,6 +18,9 @@ #include #include #include +#include +#include +#include #include extern "C" { @@ -37,6 +40,8 @@ struct VoicePlayState { qint64 startMs = 0; float durationSec = 0.0f; float pausedPosSec = 0.0f; + float displayFrac = 0.0f; + float displayElapsed = 0.0f; }; static QMap s_voiceStates; static QTimer* s_voiceTimer = nullptr; @@ -45,26 +50,25 @@ static QAbstractItemView* s_voiceView = nullptr; static void ensureVoiceTimer() { if (!s_voiceTimer) { s_voiceTimer = new QTimer(); - s_voiceTimer->setInterval(100); + s_voiceTimer->setInterval(16); QObject::connect(s_voiceTimer, &QTimer::timeout, []() { - bool anyPlaying = false; - qint64 now = QDateTime::currentMSecsSinceEpoch(); - for (auto it = s_voiceStates.begin(); it != s_voiceStates.end(); ) { - if (!it->playing && !it->paused) { ++it; continue; } - if (it->paused) { ++it; continue; } - if (!it->playing) { ++it; continue; } - float elapsed = (now - it->startMs) / 1000.0f; - if (elapsed >= it->durationSec && it->durationSec > 0) { + bool anyActive = false; + SoundManager* sm = SoundManager::instance(); + for (auto it = s_voiceStates.begin(); it != s_voiceStates.end(); ++it) { + if (!it->playing && !it->paused) { + if (it->displayFrac > 0.001f || it->displayElapsed > 0.001f) anyActive = true; + continue; + } + if (it->paused) { anyActive = true; continue; } + if (!sm->isCurrentPcmPlaying()) { it->playing = false; it->paused = false; - SoundManager::instance()->stopCurrentPcm(); - ++it; + anyActive = true; } else { - anyPlaying = true; - ++it; + anyActive = true; } } - if (!anyPlaying) s_voiceTimer->stop(); + if (!anyActive) s_voiceTimer->stop(); if (s_voiceView) s_voiceView->viewport()->update(); }); } @@ -78,11 +82,18 @@ bool MessageDelegate::editorEvent(QEvent *event, QAbstractItemModel *model, if (L.isFile) { int dlState = index.data(MsgFileDownloadStateRole).toInt(); - if (dlState != 0) return false; /* already downloaded or in progress */ qint64 msgId = index.data(MsgIdRole).toLongLong(); QString chId = index.data(MsgChannelIdRole).toString(); - if (chId.isEmpty() || msgId <= 0) return false; + if (dlState == 2) { /* downloaded: open file */ + QString fp = index.data(MsgVoiceFileRole).toString(); + if (!fp.isEmpty() && QFile::exists(fp)) + QDesktopServices::openUrl(QUrl::fromLocalFile(fp)); + return true; + } + if (dlState != 0) return false; /* downloading — no action */ + + if (chId.isEmpty() || msgId <= 0) return false; struct attachment_dl_req* req = (struct attachment_dl_req*) u_malloc(sizeof(struct attachment_dl_req)); if (!req) return false; @@ -103,6 +114,25 @@ bool MessageDelegate::editorEvent(QEvent *event, QAbstractItemModel *model, if (!onPlayBtn && !onWaveform) return false; QString filePath = index.data(MsgVoiceFileRole).toString(); + int dlState = index.data(MsgFileDownloadStateRole).toInt(); + qint64 msgId = index.data(MsgIdRole).toLongLong(); + QString chId = index.data(MsgChannelIdRole).toString(); + + /* not downloaded: start download (play will start on completion) */ + if (filePath.isEmpty() && dlState != 1) { + if (chId.isEmpty() || msgId <= 0) return false; + struct attachment_dl_req* req = (struct attachment_dl_req*) + u_malloc(sizeof(struct attachment_dl_req)); + if (!req) return false; + memset(req, 0, sizeof(*req)); + strncpy(req->channel_id, chId.toUtf8().constData(), sizeof(req->channel_id) - 1); + req->msg_id = msgId; + gui_bridge_post_uasync_fn(chat_core_attachment_download_trampoline, req); + model->setData(index, true, MsgVoiceDownloadPendingRole); + return true; + } + + if (filePath.isEmpty()) return false; GUI_DEBUG("editorEvent voice click: filePath=%s exists=%d wave=%d", qPrintable(filePath), QFile::exists(filePath), (int)onWaveform); if (filePath.isEmpty()) return false; @@ -128,6 +158,8 @@ bool MessageDelegate::editorEvent(QEvent *event, QAbstractItemModel *model, st.durationSec = dur; st.startMs = now - (qint64)(targetSec * 1000.0); st.pausedPosSec = targetSec; + st.displayFrac = (targetSec > 0 && dur > 0) ? (targetSec / dur) : 0; + st.displayElapsed = targetSec; VoicePlayer::playOpusFile(filePath); int sr = SoundManager::instance()->currentPcmSampleRate(); if (sr > 0) @@ -141,6 +173,8 @@ bool MessageDelegate::editorEvent(QEvent *event, QAbstractItemModel *model, st.pausedPosSec = targetSec; else if (!st.playing) { st.playing = true; s_voiceTimer->start(); } + st.displayFrac = (targetSec > 0 && dur > 0) ? (targetSec / dur) : 0; + st.displayElapsed = targetSec; int sr = SoundManager::instance()->currentPcmSampleRate(); if (sr > 0) SoundManager::instance()->seekCurrentPcm((unsigned long long)(targetSec * sr)); @@ -172,6 +206,8 @@ bool MessageDelegate::editorEvent(QEvent *event, QAbstractItemModel *model, st.playing = true; qint64 now = QDateTime::currentMSecsSinceEpoch(); st.startMs = now - (qint64)(st.pausedPosSec * 1000.0); + st.displayFrac = (dur > 0) ? (st.pausedPosSec / dur) : 0; + st.displayElapsed = st.pausedPosSec; SoundManager::instance()->resumeCurrentPcm(); s_voiceTimer->start(); model->setData(index, 1, MsgVoicePlayStateRole); @@ -191,6 +227,8 @@ bool MessageDelegate::editorEvent(QEvent *event, QAbstractItemModel *model, st.pausedPosSec = 0; st.startMs = QDateTime::currentMSecsSinceEpoch(); st.durationSec = dur; + st.displayFrac = 0; + st.displayElapsed = 0; VoicePlayer::playOpusFile(filePath); s_voiceTimer->start(); @@ -269,14 +307,23 @@ MessageDelegate::Layout MessageDelegate::calcLayout( int wh = 25; // waveform height int th = 14; // text line height L.bubbleWidth = 220; + int barW = 3, barPad = 4; int bubbleH = kPadTop + wh + 4 + th + kPadBot; L.totalHeight = qMax(bubbleH, kAvatar) + kItemGap; int contentEdge = L.isOutgoing ? L.viewWidth - (L.narrow ? kMarginH : (kMarginH + kAvatar + kGap)) : (L.narrow ? kMarginH : (kMarginH + kAvatar + kGap)); - int bx = L.isOutgoing ? (contentEdge - L.bubbleWidth) : contentEdge; + int bx = L.isOutgoing ? (contentEdge - L.bubbleWidth - barW - barPad) : contentEdge; L.bubbleRect = QRect(bx, 0, L.bubbleWidth, bubbleH); + L.mediaBarRect = QRect(L.isOutgoing ? bx + L.bubbleWidth + 1 : bx - barW - 1, 0, barW, bubbleH); + + int dlState = index.data(MsgFileDownloadStateRole).toInt(); + int bd = index.data(MsgMediaBlocksDoneRole).toInt(); + int nb = index.data(MsgMediaNumBlocksRole).toInt(); + L.mediaBarColor = (dlState == 2) ? 2 : (dlState == 1 || (bd > 0 && nb > 0)) ? 1 : 0; + if (dlState == 1 && nb > 0) L.progressFrac = (float)bd / nb; + if (dlState == 1) { L.progressBarRect = QRect(bx, bubbleH - 3, L.bubbleWidth, 3); } L.voicePlayBtnRect = QRect(bx + kPadH, kPadTop + 4, pbs, pbs); int wx = bx + kPadH + pbs + 4; @@ -302,14 +349,23 @@ MessageDelegate::Layout MessageDelegate::calcLayout( int th = fm.height(); int ih = 24; /* icon height */ L.bubbleWidth = 230; + int barW = 3, barPad = 4; int bubbleH = kPadTop + ih + 2 + th + kPadBot; L.totalHeight = qMax(bubbleH, kAvatar) + kItemGap; int contentEdge = L.isOutgoing ? L.viewWidth - (L.narrow ? kMarginH : (kMarginH + kAvatar + kGap)) : (L.narrow ? kMarginH : (kMarginH + kAvatar + kGap)); - int bx = L.isOutgoing ? (contentEdge - L.bubbleWidth) : contentEdge; + int bx = L.isOutgoing ? (contentEdge - L.bubbleWidth - barW - barPad) : contentEdge; L.bubbleRect = QRect(bx, 0, L.bubbleWidth, bubbleH); + L.mediaBarRect = QRect(L.isOutgoing ? bx + L.bubbleWidth + 1 : bx - barW - 1, 0, barW, bubbleH); + + int dlState = index.data(MsgFileDownloadStateRole).toInt(); + int bd = index.data(MsgMediaBlocksDoneRole).toInt(); + int nb = index.data(MsgMediaNumBlocksRole).toInt(); + L.mediaBarColor = (dlState == 2) ? 2 : (dlState == 1 || (bd > 0 && nb > 0)) ? 1 : 0; + if (dlState == 1 && nb > 0) L.progressFrac = (float)bd / nb; + if (dlState == 1) { L.progressBarRect = QRect(bx, bubbleH - 3, L.bubbleWidth, 3); } L.textRect = QRect(bx + kPadH + ih + 8, kPadTop, L.bubbleWidth - 2*kPadH - ih - 8, th); L.statusRect = QRect(bx + kPadH + ih + 8 - 48, kPadTop + th + 2, L.bubbleWidth - 2*kPadH - ih - 8 + 48 - 40, th); @@ -712,36 +768,44 @@ void MessageDelegate::paint(QPainter *painter, const QStyleOptionViewItem &optio if (!L.narrow) drawAvatar(painter, L, index); QString filePath = index.data(MsgVoiceFileRole).toString(); - VoicePlayState st; + VoicePlayState copySt; + VoicePlayState* st = ©St; if (!filePath.isEmpty() && s_voiceStates.contains(filePath)) { - st = s_voiceStates[filePath]; + st = &s_voiceStates[filePath]; } else { int stateVal = index.data(MsgVoicePlayStateRole).toInt(); - st.playing = (stateVal == 1); + copySt.playing = (stateVal == 1); + st = ©St; } - bool isPlaying = st.playing; - bool isPaused = st.paused; + bool isPlaying = st->playing; + bool isPaused = st->paused; /* compute elapsed/frac for progress */ - float elapsed = 0, frac = 0, durSec = st.durationSec; + float elapsed = 0, targetFrac = 0, durSec = st->durationSec; if (durSec <= 0) { QString ds = index.data(MsgVoiceDurationRole).toString(); ds.remove('s'); durSec = ds.toFloat(); } if (isPaused) { - elapsed = st.pausedPosSec; + elapsed = st->pausedPosSec; } else if (isPlaying) { - qint64 now = QDateTime::currentMSecsSinceEpoch(); - elapsed = (now - st.startMs) / 1000.0f; + SoundManager* sm = SoundManager::instance(); + if (sm->isCurrentPcmPlaying()) { + int sr = sm->currentPcmSampleRate(); + if (sr > 0) elapsed = (float)((double)sm->currentPcmCursor() / sr); + } + if (elapsed <= 0) { + qint64 now = QDateTime::currentMSecsSinceEpoch(); + elapsed = (now - st->startMs) / 1000.0f; + } if (elapsed > durSec) elapsed = durSec; } - if (durSec > 0) { - float dispElapsed = elapsed - 0.2f; /* 200ms audio output delay */ - if (dispElapsed < 0) dispElapsed = 0; - frac = dispElapsed / durSec; - } - if (frac < 0) frac = 0; if (frac > 1) frac = 1; + targetFrac = (durSec > 0) ? (elapsed / durSec) : 0; + if (targetFrac < 0) targetFrac = 0; if (targetFrac > 1) targetFrac = 1; + st->displayFrac += (targetFrac - st->displayFrac) * 0.15f; + st->displayElapsed += (elapsed - st->displayElapsed) * 0.15f; + float frac = st->displayFrac; QRectF playR = L.voicePlayBtnRect; @@ -806,7 +870,7 @@ void MessageDelegate::paint(QPainter *painter, const QStyleOptionViewItem &optio timeFont.setPointSize(timeFont.pointSize() - 2); painter->setFont(timeFont); painter->setPen(option.palette.text().color().lighter(150)); - int eSec = (int)elapsed, dSec = (int)durSec; + int eSec = (int)st->displayElapsed, dSec = (int)durSec; QString durText = QString("%1:%2 / %3:%4") .arg(eSec / 60).arg(eSec % 60, 2, 10, QChar('0')) .arg(dSec / 60).arg(dSec % 60, 2, 10, QChar('0')); @@ -819,6 +883,17 @@ void MessageDelegate::paint(QPainter *painter, const QStyleOptionViewItem &optio painter->drawText(L.voiceTimeRect, Qt::AlignRight | Qt::AlignVCenter, index.data(MsgTimeRole).toString()); + /* media bar (left edge, 3px) */ + { QColor barC = L.mediaBarColor == 2 ? QColor(0x40, 0xC0, 0x40) : L.mediaBarColor == 1 ? QColor(0xF0, 0xC0, 0x40) : QColor(0x80, 0x80, 0x80); + painter->setPen(Qt::NoPen); painter->setBrush(barC); painter->drawRect(L.mediaBarRect); } + /* progress bar (bottom edge, 3px) */ + if (L.progressBarRect.isValid()) { + painter->setBrush(QColor(0x44, 0x44, 0x44)); + painter->drawRoundedRect(L.progressBarRect, 2, 2); + int pw = (int)(L.progressBarRect.width() * L.progressFrac); + if (pw > 0) { painter->setBrush(QColor(0x33, 0x90, 0xEC)); painter->drawRoundedRect(QRect(L.progressBarRect.left(), L.progressBarRect.top(), pw, 3), 2, 2); } + } + painter->restore(); return; } @@ -853,12 +928,23 @@ void MessageDelegate::paint(QPainter *painter, const QStyleOptionViewItem &optio painter->setPen(QColor(0x80, 0x80, 0x80)); qint64 fsz = index.data(MsgFileSizeRole).toLongLong(); QString sizeStr = formatFileSize(fsz); - QString stateStr = (dlState == 2) ? "Downloaded" : "Tap to download"; + QString stateStr = (dlState == 2) ? QString::fromUtf8("\xE2\x9C\x85") : (dlState == 1) ? "Downloading..." : "Download"; QString info = sizeStr + " \xC2\xB7 " + stateStr; int nameH = QFontMetrics(fileFont).height(); QRect sizeRect(L.textRect.left(), L.textRect.top() + nameH, L.textRect.width(), QFontMetrics(sizeFont).height()); painter->drawText(sizeRect, Qt::AlignLeft | Qt::AlignVCenter, info); + /* media bar (left edge, 3px) */ + { QColor barC = L.mediaBarColor == 2 ? QColor(0x40, 0xC0, 0x40) : L.mediaBarColor == 1 ? QColor(0xF0, 0xC0, 0x40) : QColor(0x80, 0x80, 0x80); + painter->setPen(Qt::NoPen); painter->setBrush(barC); painter->drawRect(L.mediaBarRect); } + /* progress bar (bottom edge, 3px) */ + if (L.progressBarRect.isValid()) { + painter->setBrush(QColor(0x44, 0x44, 0x44)); + painter->drawRoundedRect(L.progressBarRect, 2, 2); + int pw = (int)(L.progressBarRect.width() * L.progressFrac); + if (pw > 0) { painter->setBrush(QColor(0x33, 0x90, 0xEC)); painter->drawRoundedRect(QRect(L.progressBarRect.left(), L.progressBarRect.top(), pw, 3), 2, 2); } + } + painter->restore(); return; diff --git a/tools/chatgui/src/messagedelegate.h b/tools/chatgui/src/messagedelegate.h index 0255fd42..1f2d64ad 100644 --- a/tools/chatgui/src/messagedelegate.h +++ b/tools/chatgui/src/messagedelegate.h @@ -29,6 +29,10 @@ enum MessageDataRole { MsgFileSizeRole = Qt::UserRole + 24, MsgFileDownloadStateRole = Qt::UserRole + 25, MsgChannelIdRole = Qt::UserRole + 26, + MsgFileDownloadProgressRole = Qt::UserRole + 27, + MsgMediaBlocksDoneRole = Qt::UserRole + 28, + MsgMediaNumBlocksRole = Qt::UserRole + 29, + MsgVoiceDownloadPendingRole = Qt::UserRole + 30, }; class MessageDelegate : public QStyledItemDelegate { @@ -70,6 +74,10 @@ private: QRect voiceDurationRect; QRect voiceTimeRect; QString voiceDuration; + int mediaBarColor = 0; /* 0=gray, 1=yellow, 2=green */ + QRect mediaBarRect; + QRect progressBarRect; + float progressFrac = 0.0f; }; Layout calcLayout(const QStyleOptionViewItem &option, diff --git a/tools/chatgui/src/messagelist.cpp b/tools/chatgui/src/messagelist.cpp index 950ee790..2d808832 100644 --- a/tools/chatgui/src/messagelist.cpp +++ b/tools/chatgui/src/messagelist.cpp @@ -7,6 +7,7 @@ #include "animtimer.h" #include "lottieicon.h" #include "audiorecorder.h" +#include "voiceplayback.h" #include "../db/db_manager.h" #include "../../lib/debug_config.h" #include @@ -126,12 +127,16 @@ static void setVoiceMessageRoles(QStandardItem* item, const QByteArray& data, item->setData(parts[1].toULongLong(), MsgVoiceBlockSizeRole); item->setData(parts[2].toUInt(), MsgVoiceNumBlocksRole); item->setData(parts[3], MsgVoiceSigsRole); + item->setData(parts[2].toUInt(), MsgMediaNumBlocksRole); } } /* file path from local_attrs */ QJsonObject la = QJsonDocument::fromJson(localAttrs).object(); QString fp = la.value("fp").toString(); + QString st = la.value("st").toString(); + int dlState = (st == "fl" && !fp.isEmpty()) ? 2 : 0; + item->setData(dlState, MsgFileDownloadStateRole); if (!fp.isEmpty()) { QString fullPath = mediaDirBase + "/" + channelId + "/" + fp; item->setData(fullPath, MsgVoiceFileRole); @@ -166,6 +171,8 @@ static void setVoiceMessageRoles(QStandardItem* item, const QByteArray& data, QStringList parts = meta.split('|'); if (parts.size() >= 1) item->setData(parts[0].toULongLong(), MsgFileSizeRole); + if (parts.size() >= 3) + item->setData(parts[2].toUInt(), MsgMediaNumBlocksRole); } /* parse download state from local_attrs */ @@ -646,3 +653,63 @@ void MessageList::cleanupAnimations() { s_icons.clear(); DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "MessageList: cleanupAnimations done"); } + +void MessageList::updateMessageProgress(const QString& chId, int64_t msgId, + int blocksDone, int numBlocks) { + if (chId != m_currentChannelId) return; + for (int r = 0; r < m_model->rowCount(); r++) { + QStandardItem* item = m_model->item(r); + if (!item) continue; + if (item->data(MsgIdRole).toLongLong() == msgId) { + item->setData(blocksDone, MsgMediaBlocksDoneRole); + item->setData(numBlocks, MsgMediaNumBlocksRole); + item->setData(1, MsgFileDownloadStateRole); + m_view->viewport()->update(); + return; + } + } +} + +void MessageList::onAttachmentDownloaded(const QString& chId, int64_t msgId) { + QString ct; + QString filePath; + int row = -1; + for (int r = 0; r < m_model->rowCount(); r++) { + QStandardItem* item = m_model->item(r); + if (!item) continue; + if (item->data(MsgIdRole).toLongLong() == msgId) { + row = r; + item->setData(2, MsgFileDownloadStateRole); + item->setData(0, MsgMediaBlocksDoneRole); + ct = item->data(MsgContentTypeRole).toString(); + filePath = item->data(MsgVoiceFileRole).toString(); + break; + } + } + if (row >= 0 && chId == m_currentChannelId) { + bool wasPending = false; + for (int r = 0; r < m_model->rowCount(); r++) { + QStandardItem* it = m_model->item(r); + if (it && it->data(MsgIdRole).toLongLong() == msgId) { + wasPending = it->data(MsgVoiceDownloadPendingRole).toBool(); + if (wasPending) it->setData(false, MsgVoiceDownloadPendingRole); + break; + } + } + refresh(); /* reload to get updated local_attrs with file path */ + /* auto-play only if user clicked play to initiate download */ + if (ct == "audio/opus" && wasPending) { + QTimer::singleShot(500, this, [this, msgId]() { + for (int r = 0; r < m_model->rowCount(); r++) { + QStandardItem* it = m_model->item(r); + if (it && it->data(MsgIdRole).toLongLong() == msgId) { + QString fp = it->data(MsgVoiceFileRole).toString(); + if (!fp.isEmpty() && QFile::exists(fp)) + VoicePlayer::playOpusFile(fp); + break; + } + } + }); + } + } +} diff --git a/tools/chatgui/src/messagelist.h b/tools/chatgui/src/messagelist.h index 50a989b6..5e571c59 100644 --- a/tools/chatgui/src/messagelist.h +++ b/tools/chatgui/src/messagelist.h @@ -26,6 +26,8 @@ public: static void cleanupAnimations(); void addMessage(const QString& channelId, quint64 authorNodeId, const QByteArray& content, qint64 timestamp); + void updateMessageProgress(const QString& chId, int64_t msgId, int blocksDone, int numBlocks); + void onAttachmentDownloaded(const QString& chId, int64_t msgId); signals: void readPositionChanged(const QString& channelId, qint64 lastReadMsgId); diff --git a/tools/chatgui/src/sound_manager.cpp b/tools/chatgui/src/sound_manager.cpp index c9458d47..3b3787a7 100644 --- a/tools/chatgui/src/sound_manager.cpp +++ b/tools/chatgui/src/sound_manager.cpp @@ -193,6 +193,18 @@ void SoundManager::stopCurrentPcm() { m_currentPcmSampleRate = 0; } +unsigned long long SoundManager::currentPcmCursor() const { + if (!m_currentPcmSound) return 0; + ma_uint64 cursor = 0; + ma_sound_get_cursor_in_pcm_frames(m_currentPcmSound, &cursor); + return (unsigned long long)cursor; +} + +bool SoundManager::isCurrentPcmPlaying() const { + if (!m_currentPcmSound) return false; + return ma_sound_is_playing(m_currentPcmSound) != MA_FALSE; +} + void SoundManager::seekCurrentPcm(unsigned long long pcmFrame) { if (m_currentPcmSound) ma_sound_seek_to_pcm_frame(m_currentPcmSound, pcmFrame); } diff --git a/tools/chatgui/src/sound_manager.h b/tools/chatgui/src/sound_manager.h index b4895573..dc4c48eb 100644 --- a/tools/chatgui/src/sound_manager.h +++ b/tools/chatgui/src/sound_manager.h @@ -43,6 +43,8 @@ public: void resumeCurrentPcm(); void stopCurrentPcm(); void seekCurrentPcm(unsigned long long pcmFrame); + unsigned long long currentPcmCursor() const; + bool isCurrentPcmPlaying() const; int currentPcmSampleRate() const { return m_currentPcmSampleRate; } struct DeviceInfo { diff --git a/tools/chatgui/transport/gui_bridge.h b/tools/chatgui/transport/gui_bridge.h index 91a73267..404af2d9 100644 --- a/tools/chatgui/transport/gui_bridge.h +++ b/tools/chatgui/transport/gui_bridge.h @@ -24,6 +24,7 @@ struct UASYNC; #define GUI_EVT_STATUS_REFRESH 10 /* data: status text (null-terminated string) */ #define GUI_EVT_NODE_CHANGED 11 /* data: [node_id:8] */ #define GUI_EVT_ATTACHMENT_DOWNLOADED 14 /* data: [ch_id_len:1][ch_id:var][msg_id:8] */ +#define GUI_EVT_DOWNLOAD_PROGRESS 15 /* data: [ch_id_len:1][ch_id:var][msg_id:8][blocks_done:4][num_blocks:4] */ /* ── API ── */ @@ -83,6 +84,10 @@ void gui_bridge_set_status_refresh_cb(gui_status_refresh_fn cb); typedef void (*gui_attachment_downloaded_fn)(const char* ch_id, int ch_id_len, int64_t msg_id); void gui_bridge_set_attachment_downloaded_cb(gui_attachment_downloaded_fn cb); +/* Callback для прогресса загрузки медиа */ +typedef void (*gui_download_progress_fn)(const char* ch_id, int ch_id_len, int64_t msg_id, int blocks_done, int num_blocks); +void gui_bridge_set_download_progress_cb(gui_download_progress_fn cb); + #ifdef __cplusplus } #endif diff --git a/tools/chatgui/transport/gui_bridge_impl.cpp b/tools/chatgui/transport/gui_bridge_impl.cpp index ee8610db..308d72fc 100644 --- a/tools/chatgui/transport/gui_bridge_impl.cpp +++ b/tools/chatgui/transport/gui_bridge_impl.cpp @@ -34,6 +34,7 @@ static gui_channel_peers_online_fn g_channel_peers_online_cb = nullptr; static gui_db_ready_fn g_db_ready_cb = nullptr; static gui_status_refresh_fn g_status_refresh_cb = nullptr; static gui_attachment_downloaded_fn g_attachment_downloaded_cb = nullptr; +static gui_download_progress_fn g_download_progress_cb = nullptr; static struct UASYNC* g_ua = nullptr; /* ── GuiBridgeReceiver implementation ── */ @@ -152,6 +153,17 @@ void GuiBridgeReceiver::processPost(int eventType, QByteArray data) { } } break; + case GUI_EVT_DOWNLOAD_PROGRESS: + if (dlen >= 2) { + uint8_t chLen = d[0]; + if (dlen >= 1 + chLen + 16) { + int64_t msgId = 0; memcpy(&msgId, d + 1 + chLen, 8); + int32_t bd = 0, nb = 0; memcpy(&bd, d + 1 + chLen + 8, 4); memcpy(&nb, d + 1 + chLen + 12, 4); + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "gui_bridge: DOWNLOAD_PROGRESS ch=%.*s msgId=%lld %d/%d", chLen, (const char*)d + 1, (long long)msgId, bd, nb); + if (g_download_progress_cb) g_download_progress_cb((const char*)d + 1, chLen, msgId, bd, nb); + } + } + break; default: DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "gui_bridge: unknown event type %d", eventType); break; @@ -238,6 +250,10 @@ void gui_bridge_set_attachment_downloaded_cb(gui_attachment_downloaded_fn cb) { g_attachment_downloaded_cb = cb; } +void gui_bridge_set_download_progress_cb(gui_download_progress_fn cb) { + g_download_progress_cb = cb; +} + } /* extern "C" */ #include "gui_bridge_impl.moc"