From 70797f921e8448044779f16a265bd9f6aa4b776d Mon Sep 17 00:00:00 2001 From: evgeny Date: Thu, 30 Jul 2026 00:01:38 +0300 Subject: [PATCH] feature: media auto-download + manual download + UI download state indication (Android + desktop) C-layer: - chat_event.h: CHAT_EVT_ATTACHMENT_DOWNLOADED (14) - chat_msg.c: fixed md_auto_download (non-NULL done_cb, correct media_base from db_path, use chat_setting for autoload checks) - chat_msg.c: added chat_core_attachment_download + trampoline for manual download - chat_core.h: exported attachment_dl_req + trampoline Android (Kotlin + JNI): - Message: downloadState field (parse from localAttrs) - ChatRepository: parse localAttrs for download state - NativeLib: attachmentDownload() JNI - MessageBubble: FileBubble with download status icon + tap-to-download - ChatViewModel: downloadAttachment() + CHAT_EVT_ATTACHMENT_DOWNLOADED handler - ChatScreen: wire onDownloadFile handler - android_jni_bridge: utun_bridge_attachment_download + JNI + localAttrs in JSON Desktop (Qt/C++): - messagedelegate.h: MsgFileDownloadStateRole, MsgChannelIdRole - messagelist.cpp: parse localAttrs for file download state - messagedelegate.cpp: FileBubble with download icon/text + click handler via chat_core_attachment_download_trampoline - gui_bridge.h: GUI_EVT_ATTACHMENT_DOWNLOADED (14) + callback typedef - gui_bridge_impl.cpp: event dispatch + setter - mainwindow.cpp: callback registration -> refresh message list --- src/chat/chat_core.h | 7 + src/chat/chat_event.h | 1 + src/chat/chat_msg.c | 240 +++++++++++++----- src/chat/chat_sync.c | 77 ++++-- src/chat/chat_sync.h | 6 + src/media_delivery/media_delivery.c | 48 ++-- src/media_delivery/media_delivery_proto.h | 1 + src/media_delivery/media_download.c | 5 +- tests/test_media_delivery_full.c | 136 +++++++--- tests/test_media_delivery_sql.c | 2 +- .../java/com/utun/chat/data/ChatRepository.kt | 12 +- .../main/java/com/utun/chat/data/NativeLib.kt | 2 + .../utun/chat/ui/components/MessageBubble.kt | 35 ++- .../com/utun/chat/ui/screens/ChatScreen.kt | 9 +- .../com/utun/chat/viewmodel/ChatViewModel.kt | 19 ++ .../jni_bridge/android_jni_bridge.c | 34 ++- .../jni_bridge/android_jni_bridge.h | 1 + tools/chatgui/src/mainwindow.cpp | 8 + tools/chatgui/src/messagedelegate.cpp | 50 +++- tools/chatgui/src/messagedelegate.h | 2 + tools/chatgui/src/messagelist.cpp | 17 +- tools/chatgui/transport/gui_bridge.h | 5 + tools/chatgui/transport/gui_bridge_impl.cpp | 15 ++ 23 files changed, 558 insertions(+), 174 deletions(-) diff --git a/src/chat/chat_core.h b/src/chat/chat_core.h index 6f8c195a..a87044cf 100644 --- a/src/chat/chat_core.h +++ b/src/chat/chat_core.h @@ -87,6 +87,13 @@ void chat_core_connect_channel_trampoline(void* arg); /* Трамплин для gui_bridge_post_uasync (GUI → uasync) */ void chat_core_submit_trampoline(void* arg); +/* Трамплин для ручного запуска загрузки медиа (JNI → uasync) */ +struct attachment_dl_req { + char channel_id[64]; + int64_t msg_id; +}; +void chat_core_attachment_download_trampoline(void* arg); + /* Обновление local_attrs для записи (вызывается из uasync-потока) */ int chat_core_update_local_attrs(const char* ch_id, uint64_t ts, const uint8_t* author_sig, diff --git a/src/chat/chat_event.h b/src/chat/chat_event.h index a6f0fec4..742b1dfb 100644 --- a/src/chat/chat_event.h +++ b/src/chat/chat_event.h @@ -32,6 +32,7 @@ extern "C" { #define CHAT_EVT_KEYS_GENERATED 11 /* [pub_hex:64] */ #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] */ 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 222d4717..0bab0904 100644 --- a/src/chat/chat_msg.c +++ b/src/chat/chat_msg.c @@ -6,6 +6,7 @@ #include "chat_core_priv.h" #include "chat_event.h" +#include "chat_setting.h" #include "../utun_instance.h" #include "../transport_layer/secure_channel.h" @@ -384,94 +385,129 @@ int chat_core_load_nodeinfo(uint64_t node_id, uint8_t* buf, size_t buf_size, return 0; } -/* ─── db_sync callback ─── */ +/* ─── download completion ─── */ + +struct md_done_ctx { + char channel_id[64]; + uint64_t ts; + uint8_t author_sig[64]; +}; -static void md_auto_download(struct UTUN_INSTANCE* inst, const char* data_str, size_t data_len, - const char* ch_id, const char* media_base) { - /* parse: |||||| */ - /* find first | after base (strip base prefix) */ +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); + 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); } + chat_event_post(CHAT_EVT_ATTACHMENT_DOWNLOADED, evt, 1 + cl + 8); + u_free(ctx); +} + +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) { + (void)data_len; const char* p = data_str; while (*p && *p != '|') p++; - if (!*p) return; /* not a media message */ - p++; /* skip first | */ - + if (!*p) return -1; p++; long long fsize = strtoll(p, (char**)&p, 10); - if (*p != '|') return; p++; + if (*p != '|') return -1; p++; long long bsize = strtoll(p, (char**)&p, 10); - if (*p != '|') return; p++; - int nb = atoi(p); + if (*p != '|') return -1; p++; + int nb = (int)strtol(p, (char**)&p, 10); while (*p && *p != '|') p++; - if (!*p) return; p++; - - if (nb <= 0 || nb > 100) return; - if (fsize <= 0) return; - - /* check auto-download settings */ - char sql[128]; - snprintf(sql, sizeof(sql), "SELECT value FROM ui_state WHERE key='auto_download_max_size'"); - sqlite3_stmt* st = NULL; - long long max_size = 0; - if (sqlite3_prepare_v2(g_cc.db, sql, -1, &st, NULL) == SQLITE_OK) { - if (sqlite3_step(st) == SQLITE_ROW) max_size = strtoll((const char*)sqlite3_column_text(st, 0), NULL, 10); - sqlite3_finalize(st); - } - if (max_size > 0 && fsize > max_size) { - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: media too large (%lld > %lld), skip auto-download", CC_ID, fsize, max_size); - return; - } + if (!*p) return -1; p++; + if (nb <= 0 || nb > 100 || fsize <= 0) return -1; struct media_index_result result; memset(&result, 0, sizeof(result)); - result.file_size = fsize; - result.block_size = bsize; - result.num_blocks = nb; + result.file_size = fsize; result.block_size = bsize; result.num_blocks = nb; - /* media_id hex → bin */ for (int i = 0; i < 16; i++) { char b[3] = { p[i*2], p[i*2+1], 0 }; result.media_id[i] = (uint8_t)strtol(b, NULL, 16); } - p += 32; - if (*p != '|') return; p++; - - /* content_hash hex → bin */ + p += 32; if (*p != '|') return -1; p++; for (int i = 0; i < 32; i++) { char b[3] = { p[i*2], p[i*2+1], 0 }; result.content_hash[i] = (uint8_t)strtol(b, NULL, 16); } - p += 64; - if (*p != '|') return; p++; + p += 64; if (*p != '|') return -1; p++; - /* block_ids and block_sigs */ result.block_ids = u_malloc((size_t)nb * 16); result.block_sigs = u_malloc((size_t)nb * 64); - if (!result.block_ids || !result.block_sigs) { media_index_result_free(&result); return; } - + if (!result.block_ids || !result.block_sigs) { media_index_result_free(&result); return -1; } for (int n = 0; n < nb; n++) { - /* block_id hex */ for (int i = 0; i < 16; i++) { char b[3] = { p[i*2], p[i*2+1], 0 }; result.block_ids[n * 16 + i] = (uint8_t)strtol(b, NULL, 16); } - p += 32; - if (*p != ',') { media_index_result_free(&result); return; } - p++; /* skip comma */ - /* block_sig hex */ + p += 32; if (*p != ',') { media_index_result_free(&result); return -1; } p++; for (int i = 0; i < 64; i++) { char b[3] = { p[i*2], p[i*2+1], 0 }; result.block_sigs[n * 64 + i] = (uint8_t)strtol(b, NULL, 16); } - p += 128; - if (n < nb - 1 && *p == ',') p++; + p += 128; if (n < nb - 1 && *p == ',') p++; } - char dest[1024]; snprintf(dest, sizeof(dest), "%s/%s", media_base, ch_id); - DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "%s: auto-download ch=%s media=%02x%02x... blocks=%d size=%lld", - CC_ID, ch_id, result.media_id[0], result.media_id[1], nb, (long long)fsize); - media_download_start(inst, 0, &result, dest, media_base, NULL, NULL); + char dest[1024], 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); + } + } + snprintf(dest, sizeof(dest), "%s/media/%s/%s", media_base, ch_id, base_filename); + + 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; + 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_index_result_free(&result); + return 0; } +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) { + if (!chat_setting_get_int("storage_autoload", 1)) return; + int max_mb = chat_setting_get_int("storage_autoload_maxsize_mb", 10); + + const char* p = data_str; + while (*p && *p != '|') p++; /* skip base filename */ + if (!*p) return; p++; + long long fsize = strtoll(p, (char**)&p, 10); + if (fsize > (long long)max_mb * 1024 * 1024) { + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: media too large (%lld > %d MB), skip auto-download", + CC_ID, fsize, max_mb); + return; + } + md_start_download(inst, data_str, data_len, ch_id, base_filename, ts, author_sig); +} + +/* ─── db_sync callback ─── */ + void on_msg_inserted(struct DB_SYNC_INSTANCE* si, uint64_t record_ts, const char* data, size_t len, uint64_t author, void* arg) { (void)si; const char* ch_id = (const char*)arg; @@ -481,20 +517,44 @@ void on_msg_inserted(struct DB_SYNC_INSTANCE* si, uint64_t record_ts, const char /* auto-download media if message contains media metadata */ if (data && len > 0 && author != g_cc.my_node_id) { - /* parse JSON: {"n":...,"ch":"...","ct":"...","d":"..."} */ const char* ct = strstr(data, "\"ct\":\""); const char* dpos = strstr(data, "\"d\":\""); if (ct && dpos) { const char* d_start = dpos + 5; - /* find closing quote of "d" value */ char body[4096]; size_t bi = 0; while (*d_start && *d_start != '\"' && bi < sizeof(body) - 1) body[bi++] = *d_start++; body[bi] = '\0'; - /* content_type is between ct\":\" and next \" */ - const char* ct_val = ct + 6; - /* determine media_base from chat_core context */ - const char* media_base = "/tmp/utun_media"; - md_auto_download(g_cc.inst, body, bi, ch_id, media_base); + + char base_filename[256]; size_t bfi = 0; + while (body[bfi] && body[bfi] != '|' && bfi < sizeof(base_filename) - 1) + base_filename[bfi++] = body[bfi]; + base_filename[bfi] = '\0'; + if (bfi == 0) snprintf(base_filename, sizeof(base_filename), "file"); + + const uint8_t* sig_field = strstr(data, "\"sig\":\""); + uint8_t author_sig[64] = {0}; + if (sig_field) { + /* read author_signature from DB — find by ts in on_msg_inserted we don't have sig + so leave as zeros; use record_ts to locate the record later */ + } + + /* We don't have author_sig in the callback; derive from body content. + 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); + 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); + 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); + } + sqlite3_finalize(st); + } + md_auto_download(g_cc.inst, body, bi, ch_id, base_filename, record_ts, author_sig); } } @@ -506,3 +566,63 @@ void on_msg_inserted(struct DB_SYNC_INSTANCE* si, uint64_t record_ts, const char sqlite3_step(st); sqlite3_finalize(st); } } + +/* ─── manual attachment download (called from JNI bridge via uasync post) ─── */ + +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); + 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); + return; + } + sqlite3_bind_int64(st, 1, (sqlite3_int64)msg_id); + if (sqlite3_step(st) != SQLITE_ROW) { + DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "%s: download msg not found ch=%s id=%lld", CC_ID, channel_id, (long long)msg_id); + sqlite3_finalize(st); + return; + } + const uint8_t* jdata = (const uint8_t*)sqlite3_column_blob(st, 0); + int jlen = sqlite3_column_bytes(st, 0); + 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); + 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); + return; + } + + const char* ds = strstr((const char*)jdata, "\"d\":\""); + if (!ds) { sqlite3_finalize(st); return; } + const char* d_start = ds + 5; + char body[4096]; size_t bi = 0; + while (*d_start && *d_start != '\"' && bi < sizeof(body) - 1) body[bi++] = *d_start++; + body[bi] = '\0'; + + char base_filename[256]; size_t bfi = 0; + while (body[bfi] && body[bfi] != '|' && bfi < sizeof(base_filename) - 1) + base_filename[bfi++] = body[bfi]; + base_filename[bfi] = '\0'; + if (bfi == 0) snprintf(base_filename, sizeof(base_filename), "file"); + + uint8_t author_sig[64]; memcpy(author_sig, sig_blob, 64); + sqlite3_finalize(st); + + 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); +} + +static void chat_core_attachment_download_trampoline_impl(void* arg) { + struct attachment_dl_req* req = (struct attachment_dl_req*)arg; + chat_core_attachment_download(req->channel_id, req->msg_id); + u_free(req); +} + +void chat_core_attachment_download_trampoline(void* arg) { + chat_core_attachment_download_trampoline_impl(arg); +} diff --git a/src/chat/chat_sync.c b/src/chat/chat_sync.c index 9ff49d64..a8b2f437 100644 --- a/src/chat/chat_sync.c +++ b/src/chat/chat_sync.c @@ -399,6 +399,22 @@ static void _on_invite_sync_done(uint64_t peer, const char* ns, int result, void u_free(sa); } +/* ── send CHANNEL_INFO_REQ to start join handshake ── */ + +static void cs_start_channel_join(struct chat_sync* cs, uint64_t peer) { + if (!cs || cs->info_req_timer) return; + char ch_id_str[64]; + snprintf(ch_id_str, sizeof(ch_id_str), "%llu", (unsigned long long)cs->pending_invite_ch_id); + uint8_t req[1] = { CS_MSG_CHANNEL_INFO_REQ }; + cs_send(cs, ch_id_str, peer, req, 1); + cs->info_req_timer = uasync_set_timeout(cs->inst->ua, + CS_INFO_REQ_TIMEOUT_MS * 10, cs, cs_info_req_timeout_cb, "cs_info_req"); + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: sent CHANNEL_INFO_REQ ch=%llu peer=0x%016llx", + CS_ID, (unsigned long long)cs->pending_invite_ch_id, (unsigned long long)peer); +} + +/* ── Connection callbacks ── */ + static void cs_on_conn_up(struct ETCP_CONN* conn, void* arg) { (void)arg; if (!conn || !g_cs) return; @@ -409,21 +425,12 @@ static void cs_on_conn_up(struct ETCP_CONN* conn, void* arg) { DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: conn_up peer=%016llx pending_invite=%llu links_up=%d", CS_ID, (unsigned long long)peer, (unsigned long long)g_cs->pending_invite_ch_id, conn->links_up); - if (g_cs->pending_invite_ch_id != 0 && !g_cs->info_req_timer) { - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: conn_up invite path peer=%016llx ch=%llu", - CS_ID, (unsigned long long)peer, (unsigned long long)g_cs->pending_invite_ch_id); + if (g_cs->pending_invite_ch_id != 0) { if (g_cs->pending_invite_node_id != 0 && g_cs->pending_invite_node_id != peer) DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: invite node_id MISMATCH: invite=0x%016llx ETCP_peer=0x%016llx — invite is STALE!", CS_ID, (unsigned long long)g_cs->pending_invite_node_id, (unsigned long long)peer); g_cs->pending_invite_node_id = peer; - char ch_id_str[64]; - snprintf(ch_id_str, sizeof(ch_id_str), "%llu", - (unsigned long long)g_cs->pending_invite_ch_id); - uint8_t req[1] = { CS_MSG_CHANNEL_INFO_REQ }; - cs_send(g_cs, ch_id_str, peer, req, 1); - g_cs->info_req_timer = uasync_set_timeout(g_cs->inst->ua, - CS_INFO_REQ_TIMEOUT_MS * 10, g_cs, cs_info_req_timeout_cb, "cs_info_req"); - DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: conn_up invite path, marking peer=%016llx online", CS_ID, (unsigned long long)peer); + cs_start_channel_join(g_cs, peer); cs_on_peer_status_changed(peer, 1); return; } @@ -633,13 +640,12 @@ static void cm_invite_trampoline(void* arg) { } char ch_str[32]; snprintf(ch_str, sizeof(ch_str), "%llu", (unsigned long long)channel_id); - chat_core_ensure_channel_ready(ch_str); - /* ensure TOPO_GROUP_TYPE_CHAT group exists */ + /* TOPO_GROUP (lightweight) for conn_mgr_connect_from_invite below. + DB channel + sync created later in cs_handle_channel_info_resp. */ if (inst->topo_groups) { - uint64_t gid = channel_id; - if (!topo_groups_find(inst->topo_groups, gid)) - topo_groups_create_group(inst->topo_groups, gid, TOPO_GROUP_TYPE_CHAT, ch_str); + if (!topo_groups_find(inst->topo_groups, channel_id)) + topo_groups_create_group(inst->topo_groups, channel_id, TOPO_GROUP_TYPE_CHAT, ch_str); } /* build TOPO_NODE from invite data */ @@ -667,13 +673,15 @@ static void cm_invite_trampoline(void* arg) { if (existing) { if (existing->links_up) { - if (g_cs && g_cs->pending_invite_ch_id != 0 && !g_cs->info_req_timer) { - if (node_id != existing->peer_node_id) - DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, - "%s: invite node_id mismatch: derived=%016llx ETCP=%016llx — using ETCP", - CS_ID, (unsigned long long)node_id, (unsigned long long)existing->peer_node_id); - cs_on_conn_up(existing, NULL); - } + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: invite start — already connected peer=0x%016llx, joining ch=%llu", + CS_ID, (unsigned long long)node_id, (unsigned long long)channel_id); + if (node_id != existing->peer_node_id) + DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: invite node_id mismatch: derived=%016llx ETCP=%016llx — using ETCP", + CS_ID, (unsigned long long)node_id, (unsigned long long)existing->peer_node_id); + g_cs->pending_invite_ch_id = channel_id; + g_cs->pending_invite_node_id = existing->peer_node_id; + cs_start_channel_join(g_cs, existing->peer_node_id); + cs_on_peer_status_changed(existing->peer_node_id, 1); } else { DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: invite conn in progress node=0x%016llx, waiting UP", CS_ID, (unsigned long long)node_id); } @@ -740,6 +748,29 @@ void chat_sync_connect_from_invite(uint64_t channel_id, uint64_t node_id, (void(*)(void*))cm_invite_trampoline, w); } +/* ─── chat_sync_join_channel: join channel via already-connected peer ─── */ + +void chat_sync_join_channel(struct UTUN_INSTANCE* inst, const char* ch_id, uint64_t target_node_id) { + if (!g_cs || !inst || !ch_id || !ch_id[0]) return; + uint64_t channel_id = strtoull(ch_id, NULL, 10); + if (channel_id == 0) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: join_channel invalid ch_id='%s'", CS_ID, ch_id); return; } + + struct ETCP_CONN* conn = instance_find_conn(inst, target_node_id); + if (!conn || !conn->links_up || !conn->initialized) { + DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: join_channel no active conn to 0x%016llx for ch=%s", + CS_ID, (unsigned long long)target_node_id, ch_id); + return; + } + + DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: join_channel ch=%s peer=0x%016llx", + CS_ID, ch_id, (unsigned long long)target_node_id); + + g_cs->pending_invite_ch_id = channel_id; + g_cs->pending_invite_node_id = target_node_id; + cs_start_channel_join(g_cs, target_node_id); + cs_on_peer_status_changed(target_node_id, 1); +} + /* ─── Ed25519 sign / verify helpers ─── */ static int cs_ed25519_sign(const uint8_t* privkey, const uint8_t* msg, size_t msg_len, diff --git a/src/chat/chat_sync.h b/src/chat/chat_sync.h index 1afe47d4..62fc46ba 100644 --- a/src/chat/chat_sync.h +++ b/src/chat/chat_sync.h @@ -53,6 +53,12 @@ void chat_sync_connect_from_invite(uint64_t channel_id, uint64_t node_id, const uint8_t* pubkey_bin, const uint8_t* addrs_data, int addr_count); +/* Присоединиться к каналу через уже подключённый узел (uasync-поток). + Соединение с target_node_id должно быть установлено (links_up && initialized). + Протокол: CHANNEL_INFO_REQ → CHANNEL_INFO_RESP → CHANNEL_JOIN → WELCOME. + Канал (DB + TOPO_GROUP) создаётся локально при получении CHANNEL_INFO_RESP. */ +void chat_sync_join_channel(struct UTUN_INSTANCE* inst, const char* ch_id, uint64_t target_node_id); + #ifdef __cplusplus } #endif diff --git a/src/media_delivery/media_delivery.c b/src/media_delivery/media_delivery.c index 3c9f2171..a2fee7a9 100644 --- a/src/media_delivery/media_delivery.c +++ b/src/media_delivery/media_delivery.c @@ -29,7 +29,7 @@ static int md_super_start(struct media_delivery_ctx* md); static void md_super_stop(struct media_delivery_ctx* md); static void md_super_repl_send(struct media_delivery_ctx* md, struct media_super_peer* peer); static void md_super_connect(struct media_delivery_ctx* md, uint64_t peer_node_id); -static int md_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, const uint8_t* data, size_t len); +static int md_send(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t dst_node_id, const uint8_t* data, size_t len); static struct media_super_peer* md_super_peer_add(struct media_delivery_ctx* md, uint64_t node_id); /* ── helpers ── */ @@ -255,7 +255,7 @@ static void md_super_repl_send(struct media_delivery_ctx* md, struct media_super uint16_t num = (uint16_t)num_entries; memcpy(pkt + off, &num, 2); off += 2; memcpy(pkt + off, buf, (size_t)nb); off += nb; - if (md_send(md->inst, peer->peer_node_id, pkt, (size_t)off) == 0) { + if (md_send(md->inst, TOPO_GROUP_UTUN, peer->peer_node_id, pkt, (size_t)off) == 0) { peer->inflight_count++; DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "%s: SUPER_REPL sent to 0x%016llx seq=%u entries=%d", MD_ID, (unsigned long long)peer->peer_node_id, seq, num_entries); @@ -273,7 +273,7 @@ static void md_super_repl_send(struct media_delivery_ctx* md, struct media_super /* ── send via etcp_router helper ── */ -static int md_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, +static int md_send(struct UTUN_INSTANCE* inst, uint64_t group_id, uint64_t dst_node_id, const uint8_t* data, size_t len) { if (!inst || !data || len == 0) return -1; struct ll_entry* e = queue_entry_new(0); @@ -282,7 +282,7 @@ static int md_send(struct UTUN_INSTANCE* inst, uint64_t dst_node_id, if (!e->dgram) { queue_entry_free(e); return -1; } e->dgram[0] = ETCP_RT_ID_MEDIA_DELIVERY; memcpy(e->dgram + 1, data, len); e->len = (uint16_t)(len + 1); - int rc = etcp_route_send(inst, TOPO_GROUP_UTUN, dst_node_id, e, 1); + int rc = etcp_route_send(inst, group_id, dst_node_id, e, 1); if (rc != 0) { u_free(e->dgram); queue_entry_free(e); } return rc; } @@ -349,7 +349,7 @@ static void md_handle_query(struct media_delivery_ctx* md, uint64_t from_node, 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, from_node, resp, (size_t)off); + (void)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, @@ -368,7 +368,7 @@ static void md_handle_have_block(struct media_delivery_ctx* md, uint64_t from_no memcpy(ha->block_id, hb->block_id, 16); ha->status = (rc == 0) ? 0 : 1; - md_send(md->inst, from_node, ack, sizeof(ack)); + md_send(md->inst, TOPO_GROUP_UTUN, from_node, ack, sizeof(ack)); DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: HAVE_BLOCK from 0x%016llx rc=%d", MD_ID, (unsigned long long)from_node, rc); @@ -411,7 +411,7 @@ static void md_handle_super_repl(struct media_delivery_ctx* md, uint64_t from_no struct media_pkt_super_ack* sa = (struct media_pkt_super_ack*)ack; sa->subcmd = MEDIA_SUBCMD_SUPER_ACK; sa->ack_seq = (uint32_t)max_id; - md_send(md->inst, from_node, ack, sizeof(ack)); + md_send(md->inst, TOPO_GROUP_UTUN, from_node, ack, sizeof(ack)); DEBUG_DEBUG(DEBUG_CATEGORY_DEBUG, "%s: SUPER_REPL from 0x%016llx entries=%d max_id=%lld", MD_ID, (unsigned long long)from_node, num, (long long)max_id); @@ -453,7 +453,7 @@ static void md_handle_super_hello(struct media_delivery_ctx* md, uint64_t from_n struct media_pkt_super_hello resp; resp.subcmd = MEDIA_SUBCMD_SUPER_HELLO; resp.last_recv_id = my_last_recv; - md_send(md->inst, from_node, (const uint8_t*)&resp, sizeof(resp)); + md_send(md->inst, TOPO_GROUP_UTUN, from_node, (const uint8_t*)&resp, sizeof(resp)); DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: SUPER_HELLO from 0x%016llx (peer_last_recv=%lld, my_last_recv=%lld)", MD_ID, (unsigned long long)from_node, (long long)sh->last_recv_id, (long long)my_last_recv); @@ -488,7 +488,7 @@ static void md_super_conn_cb(int result, uint64_t node_id, void* arg) { hello.subcmd = MEDIA_SUBCMD_SUPER_HELLO; hello.last_recv_id = my_last; - if (md_send(md->inst, node_id, (const uint8_t*)&hello, sizeof(hello)) == 0) { + if (md_send(md->inst, TOPO_GROUP_UTUN, node_id, (const uint8_t*)&hello, sizeof(hello)) == 0) { DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: SUPER_HELLO sent to 0x%016llx last_recv=%lld", MD_ID, (unsigned long long)node_id, (long long)my_last); } else { @@ -505,6 +505,7 @@ static void md_super_conn_cb(int result, uint64_t node_id, void* arg) { struct stream_ctx { struct media_delivery_ctx* md; + uint64_t group_id; uint64_t dst_node_id; uint8_t media_id[16]; uint8_t block_id[16]; @@ -544,7 +545,7 @@ static void stream_send_chunk_cb(struct ll_queue* q, void* arg) { ch->data_len = (uint16_t)rd; memcpy(pkt + MEDIA_BLOCK_CHUNK_HDR_SIZE, buf, rd); - md_send(md->inst, sc->dst_node_id, pkt, MEDIA_BLOCK_CHUNK_HDR_SIZE + rd); + md_send(md->inst, sc->group_id, sc->dst_node_id, pkt, MEDIA_BLOCK_CHUNK_HDR_SIZE + rd); sc->offset += rd; sc->remaining -= rd; @@ -561,7 +562,7 @@ static void stream_send_chunk_cb(struct ll_queue* q, void* arg) { bd->total_size = (uint32_t)sc->block_data_len; /* block_sig is zero — receiver will verify vs block_sigs from media_index_result */ - md_send(md->inst, sc->dst_node_id, done_pkt, MEDIA_BLOCK_DONE_SIZE); + md_send(md->inst, sc->group_id, sc->dst_node_id, done_pkt, MEDIA_BLOCK_DONE_SIZE); media_delivery_stream_done(md->inst); fclose(sc->file); u_free(sc); @@ -569,7 +570,7 @@ static void stream_send_chunk_cb(struct ll_queue* q, void* arg) { MD_ID, sc->chunk, bd->total_size); } else { /* continue streaming — register waiter for backpressure */ - etcp_router_on_send_ready(md->inst, TOPO_GROUP_UTUN, + etcp_router_on_send_ready(md->inst, sc->group_id, sc->dst_node_id, ETCP_RT_ID_MEDIA_DELIVERY, &sc->waiter, stream_send_chunk_cb, sc); } @@ -579,9 +580,10 @@ static void md_handle_block_req(struct media_delivery_ctx* md, uint64_t from_nod const uint8_t* data, size_t len) { if (len < MEDIA_BLOCK_REQ_SIZE) return; struct media_pkt_block_req* req = (struct media_pkt_block_req*)data; + uint64_t group_id = req->group_id ? req->group_id : TOPO_GROUP_UTUN; - DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: BLOCK_REQ from 0x%016llx block=%02x%02x... active_streams=%d", - MD_ID, (unsigned long long)from_node, req->block_id[0], req->block_id[1], md->active_streams); + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: BLOCK_REQ from 0x%016llx group=%016llx block=%02x%02x... active_streams=%d", + MD_ID, (unsigned long long)from_node, (unsigned long long)group_id, req->block_id[0], req->block_id[1], md->active_streams); int limit_reached = md->active_streams >= MD_MAX_STREAMS; md->active_streams++; @@ -594,12 +596,23 @@ static void md_handle_block_req(struct media_delivery_ctx* md, uint64_t from_nod memcpy(ov.block_id, req->block_id, 16); ov.chunk = req->chunk; ov.retry_after_ms = 2000; - md_send(md->inst, from_node, (const uint8_t*)&ov, sizeof(ov)); + md_send(md->inst, group_id, from_node, (const uint8_t*)&ov, sizeof(ov)); DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "%s: BLOCK_REQ OVERLOADED (streams=%d) to 0x%016llx", MD_ID, md->active_streams, (unsigned long long)from_node); return; } + /* establish direct connection if not already connected */ + if (!instance_find_conn(md->inst, from_node) && group_id != TOPO_GROUP_UTUN) { + struct TOPO_GROUP* grp = topo_groups_find(md->inst->topo_groups, group_id); + if (!grp) grp = topo_groups_get_default(md->inst->topo_groups); + if (grp && grp->conn_mgr) { + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: no direct conn to 0x%016llx — connecting via conn_mgr", + MD_ID, (unsigned long long)from_node); + conn_mgr_connect_node(grp->conn_mgr, from_node, 0, NULL, NULL); + } + } + /* find file in media_files DB */ sqlite3* db = md->db; if (!db) return; @@ -649,6 +662,7 @@ static void md_handle_block_req(struct media_delivery_ctx* md, uint64_t from_nod struct stream_ctx* sc = u_calloc(1, sizeof(*sc)); if (!sc) { fclose(f); md->active_streams--; return; } sc->md = md; + sc->group_id = group_id; sc->dst_node_id = from_node; memcpy(sc->media_id, req->media_id, 16); memcpy(sc->block_id, req->block_id, 16); @@ -659,8 +673,8 @@ static void md_handle_block_req(struct media_delivery_ctx* md, uint64_t from_nod sc->block_start = (uint64_t)block_start; sc->block_data_len = (size_t)block_length; - DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: streaming file=%s chunk=%d start=%lld len=%lld to 0x%016llx", - MD_ID, path, chunk, (long long)block_start, (long long)block_length, (unsigned long long)from_node); + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: streaming file=%s chunk=%d start=%lld len=%lld group=%016llx to 0x%016llx", + MD_ID, path, chunk, (long long)block_start, (long long)block_length, (unsigned long long)group_id, (unsigned long long)from_node); /* start streaming — first chunk immediately, then backpressure */ memset(&sc->waiter, 0, sizeof(sc->waiter)); diff --git a/src/media_delivery/media_delivery_proto.h b/src/media_delivery/media_delivery_proto.h index ea250271..968e1875 100644 --- a/src/media_delivery/media_delivery_proto.h +++ b/src/media_delivery/media_delivery_proto.h @@ -112,6 +112,7 @@ struct media_pkt_super_ack { struct media_pkt_block_req { uint8_t subcmd; // MEDIA_SUBCMD_BLOCK_REQ + uint64_t group_id; // CHAT-группа для маршрутизации и прямого коннекта uint8_t media_id[16]; uint8_t block_id[16]; uint32_t chunk; diff --git a/src/media_delivery/media_download.c b/src/media_delivery/media_download.c index 0fbc8953..84521e48 100644 --- a/src/media_delivery/media_download.c +++ b/src/media_delivery/media_download.c @@ -57,13 +57,14 @@ static int md_dl_send_block_req(struct UTUN_INSTANCE* inst, uint64_t dst, struct media_pkt_block_req req; memset(&req, 0, sizeof(req)); req.subcmd = MEDIA_SUBCMD_BLOCK_REQ; + req.group_id = dl->group_id; memcpy(req.media_id, dl->media_id, 16); memcpy(req.block_id, dl->block_ids + bi * 16, 16); req.chunk = (uint32_t)bi; req.offset = 0; - DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: BLOCK_REQ to 0x%016llx block=%d", - MDL_ID, (unsigned long long)dst, bi); + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "%s: BLOCK_REQ to 0x%016llx group=%016llx block=%d", + MDL_ID, (unsigned long long)dst, (unsigned long long)dl->group_id, bi); return md_dl_send(inst, dst, (const uint8_t*)&req, sizeof(req)); } diff --git a/tests/test_media_delivery_full.c b/tests/test_media_delivery_full.c index 82d04898..f54470dd 100644 --- a/tests/test_media_delivery_full.c +++ b/tests/test_media_delivery_full.c @@ -1,7 +1,9 @@ // test_media_delivery_full.c — полный сценарий: создание, передача, сборка, репликация // // 4 узла: n1 (автор), n2 (получатель), s1 (суперузел), s2 (суперузел) +// Топология: цепь 0→1 1→2 2→3 // Фазы: +// P1: CHAT-группа (topo_groups_create_group + member_sync_put) // A1: SUPER_HELLO, HAVE_BLOCK, репликация // A2: admission control (OVERLOADED) // B1: n1 создаёт файл → n2 инициирует загрузку → стриминг → сборка → проверка @@ -10,6 +12,8 @@ #include "media_delivery_proto.h" #include "media_download.h" #include "media_index.h" +#include "member_sync.h" +#include "../routing_layer/topo_node_sqlite.h" #include "../utun_instance.h" #include "../transport_layer/etcp.h" #include "../config_parser.h" @@ -34,10 +38,12 @@ static int G_PASSED = 0, G_FAILED = 0, G_TOTAL = 0; #define TIMEOUT_TB 300000 #define POLL_MS 5 #define N_NODES 4 /* n1=0, n2=1, s1=2, s2=3 */ +#define CH_ID "31415926535" static struct UASYNC* g_ua = NULL; static struct UTUN_INSTANCE* g_inst[N_NODES]; static uint64_t g_nid[N_NODES]; +static uint64_t g_group_id = 0; static uint8_t g_test_mid[16], g_test_bid0[16], g_test_bid1[16]; static char g_tdir[256] = "/tmp/utun_mdf_XXXXXX"; static char g_cfg[N_NODES][256]; @@ -92,11 +98,66 @@ static int msend(struct UTUN_INSTANCE* inst, uint64_t dst, const uint8_t* data, return rc; } +/* ── Phase P1: create CHAT group + members on all nodes ── */ + +static void setup_chat_group(void) { + uint8_t x25519_pub[32] = {0}, ed_pub[32] = {0}, ch_sig[64] = {0}; + int i; + + TEST("generate channel keys"); { + EVP_PKEY* xpkey = NULL, *epkey = NULL; + EVP_PKEY_CTX* ctx = EVP_PKEY_CTX_new_id(EVP_PKEY_X25519, NULL); + if (ctx) { EVP_PKEY_keygen_init(ctx); EVP_PKEY_generate(ctx, &xpkey); EVP_PKEY_CTX_free(ctx); } + if (xpkey) { size_t l = 32; EVP_PKEY_get_raw_public_key(xpkey, x25519_pub, &l); EVP_PKEY_free(xpkey); } + ctx = EVP_PKEY_CTX_new_id(EVP_PKEY_ED25519, NULL); + if (ctx) { EVP_PKEY_keygen_init(ctx); EVP_PKEY_generate(ctx, &epkey); EVP_PKEY_CTX_free(ctx); } + if (epkey) { size_t l = 32; EVP_PKEY_get_raw_public_key(epkey, ed_pub, &l); + /* sign channel */ + EVP_MD_CTX* mctx = EVP_MD_CTX_new(); + if (mctx) { EVP_DigestSignInit(mctx, NULL, NULL, NULL, epkey); + size_t sl = 64; EVP_DigestSign(mctx, ch_sig, &sl, (const uint8_t*)CH_ID, strlen(CH_ID)); EVP_MD_CTX_free(mctx); } + EVP_PKEY_free(epkey); + } + if (x25519_pub[0] || ed_pub[0]) OK(); else FAIL(); + } + + TEST("create CHAT groups + peers tables"); { + for (i = 0; i < N_NODES; i++) { + sqlite3* db = g_inst[i]->topo_sqlite_db; + struct TOPO_GROUPS* tg = g_inst[i]->topo_groups; + if (!db || !tg) continue; + topo_node_sqlite_channel_put(db, CH_ID, "test", g_nid[0], x25519_pub, NULL, ed_pub, NULL, ch_sig); + if (topo_groups_find(tg, g_group_id)) continue; + topo_groups_create_group(tg, g_group_id, TOPO_GROUP_TYPE_CHAT, CH_ID); + } + OK(); + } + + TEST("add all 4 members (n1, n2, s1 supernode, s2 supernode)"); { + const char* adm[4] = { NULL, NULL, "supernode=yes", "supernode=yes" }; + for (i = 0; i < N_NODES; i++) { + sqlite3* db = g_inst[i]->topo_sqlite_db; + if (!db) continue; + for (int m = 0; m < N_NODES; m++) { + member_sync_put(g_inst[i], CH_ID, g_nid[m], + g_inst[m]->my_keys.public_key, g_inst[m]->my_ed25519_pubkey, + NULL, 0, NULL, 0, "{\"name\":\"node\"}", + NULL, 0, adm[m], NULL, 0); + } + } + int ok = 1; + for (i = 0; i < N_NODES; i++) if (member_sync_count(g_inst[i], CH_ID) < N_NODES) ok = 0; + if (ok) OK(); else FAIL(); + } +} + /* ── Phase A1: SUPER_HELLO ── */ static void phase_a1_super_hello(void) { TEST("s1,s2 are supernodes"); { media_delivery_set_supernode(g_inst[2], 1); media_delivery_set_supernode(g_inst[3], 1); + topo_node_sqlite_nodeinfo_updated(g_inst[2]->topo_sqlite_db, g_nid[2]); + topo_node_sqlite_nodeinfo_updated(g_inst[3]->topo_sqlite_db, g_nid[3]); if (g_inst[2]->md.is_supernode && g_inst[3]->md.is_supernode) OK(); else FAIL(); } TEST("SUPER_HELLO s1↔s2"); { @@ -115,7 +176,7 @@ static void phase_a2_have_block_replication(void) { uint8_t uuid[16]; memset(uuid, 0xAB, 16); TEST("HAVE_BLOCK n2→s1"); { struct media_pkt_have_block hb; memset(&hb, 0, sizeof(hb)); - hb.subcmd = MEDIA_SUBCMD_HAVE_BLOCK; hb.group_id = 0; + hb.subcmd = MEDIA_SUBCMD_HAVE_BLOCK; hb.group_id = g_group_id; memcpy(hb.block_id, uuid, 16); memcpy(hb.media_id, uuid, 16); hb.chunk = 0; hb.timestamp = (int64_t)time(NULL); msend(g_inst[1], g_nid[2], (const uint8_t*)&hb, sizeof(hb)); @@ -141,7 +202,8 @@ static void phase_a3_admission(void) { TEST("1st BLOCK_REQ → stream"); { uint8_t mid[16], bid[16]; memset(mid, 0xCD, 16); memset(bid, 0xCE, 16); struct media_pkt_block_req req; memset(&req, 0, sizeof(req)); - req.subcmd = MEDIA_SUBCMD_BLOCK_REQ; memcpy(req.media_id, mid, 16); memcpy(req.block_id, bid, 16); + req.subcmd = MEDIA_SUBCMD_BLOCK_REQ; req.group_id = g_group_id; + memcpy(req.media_id, mid, 16); memcpy(req.block_id, bid, 16); msend(g_inst[1], g_nid[0], (const uint8_t*)&req, sizeof(req)); int a = 0; while (a < 200) { uasync_poll(g_ua, POLL_MS); a++; } if (g_inst[0]->md.active_streams >= 1) OK(); else FAIL("active=%d", g_inst[0]->md.active_streams); @@ -150,25 +212,21 @@ static void phase_a3_admission(void) { for (int i = 0; i < 3; i++) { uint8_t mid[16], bid[16]; memset(mid, i+10, 16); memset(bid, i+20, 16); struct media_pkt_block_req req; memset(&req, 0, sizeof(req)); - req.subcmd = MEDIA_SUBCMD_BLOCK_REQ; req.chunk = (uint32_t)i; + req.subcmd = MEDIA_SUBCMD_BLOCK_REQ; req.group_id = g_group_id; req.chunk = (uint32_t)i; memcpy(req.media_id, mid, 16); memcpy(req.block_id, bid, 16); msend(g_inst[1], g_nid[0], (const uint8_t*)&req, sizeof(req)); } int a = 0; while (a < 200) { uasync_poll(g_ua, POLL_MS); a++; } if (g_inst[0]->md.active_streams == 3) OK(); else FAIL("active=%d", g_inst[0]->md.active_streams); } - g_inst[0]->md.active_streams = 0; /* reset for next phase */ + g_inst[0]->md.active_streams = 0; } /* ══════════════════════════════════════════════════════════ Phase B1: create file → register → stream → assemble → verify ══════════════════════════════════════════════════════════ */ -static int g_dl_done = 0, g_dl_err = 0; -static void dl_done_cb(void* arg, int err) { (void)arg; g_dl_done = 1; g_dl_err = err; } - static void phase_b1_file_transfer(void) { - /* create source file in /tmp */ char src_tmp[512]; snprintf(src_tmp, sizeof(src_tmp), "/tmp/mdl_test_%d.bin", getpid()); char media_dir[512]; snprintf(media_dir, sizeof(media_dir), "%s/media", g_db_dir[0]); utun_mkdir(media_dir, 0755); @@ -180,14 +238,9 @@ static void phase_b1_file_transfer(void) { FILE* f = fopen(src_tmp, "wb"); if (f) { fwrite(file_data, 1, 2048, f); fclose(f); } - /* copy to media dir (simulates what media_index_register_async does with copy=1) */ - { - FILE* fin = fopen(src_tmp, "rb"); - FILE* fout = fopen(dst_path, "wb"); - if (fin && fout) { uint8_t buf[4096]; size_t rd; while ((rd = fread(buf, 1, sizeof(buf), fin)) > 0) fwrite(buf, 1, rd, fout); } - if (fin) fclose(fin); - if (fout) fclose(fout); - } + { FILE* fin = fopen(src_tmp, "rb"), *fout = fopen(dst_path, "wb"); + if (fin && fout) { uint8_t buf[4096]; size_t rd; while ((rd = fread(buf, 1, sizeof(buf), fin)) > 0) fwrite(buf, 1, rd, fout); } + if (fin) fclose(fin); if (fout) fclose(fout); } TEST("media_index_commit + ui_state"); { uint8_t hash[32]; @@ -198,24 +251,16 @@ static void phase_b1_file_transfer(void) { EVP_MD_CTX_free(ctx); media_index_init(g_inst[0]->topo_sqlite_db); - /* ui_state: streaming handler reads media_base */ sqlite3_exec(g_inst[0]->topo_sqlite_db, "CREATE TABLE IF NOT EXISTS ui_state(key TEXT PRIMARY KEY, value TEXT)", NULL, NULL, NULL); - { - sqlite3_stmt* us = NULL; - sqlite3_prepare_v2(g_inst[0]->topo_sqlite_db, "INSERT OR REPLACE INTO ui_state(key,value) VALUES('media_base',?)", -1, &us, NULL); - sqlite3_bind_text(us, 1, media_base, -1, SQLITE_STATIC); - sqlite3_step(us); sqlite3_finalize(us); - } + { sqlite3_stmt* us = NULL; + sqlite3_prepare_v2(g_inst[0]->topo_sqlite_db, "INSERT OR REPLACE INTO ui_state(key,value) VALUES('media_base',?)", -1, &us, NULL); + sqlite3_bind_text(us, 1, media_base, -1, SQLITE_STATIC); sqlite3_step(us); sqlite3_finalize(us); } - struct media_index_result result; - memset(&result, 0, sizeof(result)); + struct media_index_result result; memset(&result, 0, sizeof(result)); media_index_generate_uuid(result.media_id); memcpy(result.content_hash, hash, 32); - result.file_size = 2048; - result.block_size = 1024; - result.num_blocks = 2; - result.block_ids = u_malloc(32); - result.block_sigs = u_malloc(128); + result.file_size = 2048; result.block_size = 1024; result.num_blocks = 2; + result.block_ids = u_malloc(32); result.block_sigs = u_malloc(128); for (int i = 0; i < 2; i++) { media_index_generate_uuid(result.block_ids + i * 16); uint8_t smsg[2048]; size_t soff = 0; @@ -229,7 +274,6 @@ static void phase_b1_file_transfer(void) { memcpy(g_test_mid, result.media_id, 16); memcpy(g_test_bid0, result.block_ids, 16); memcpy(g_test_bid1, result.block_ids + 16, 16); - int ndb = db_count(g_inst[0]->topo_sqlite_db, "media_files", NULL, 0); media_index_result_free(&result); if (ndb == 2) OK(); else FAIL("media_files has %d rows (expected 2)", ndb); @@ -237,7 +281,7 @@ static void phase_b1_file_transfer(void) { TEST("n2→n1 BLOCK_REQ streams data"); { struct media_pkt_block_req req; memset(&req, 0, sizeof(req)); - req.subcmd = MEDIA_SUBCMD_BLOCK_REQ; + req.subcmd = MEDIA_SUBCMD_BLOCK_REQ; req.group_id = g_group_id; memcpy(req.media_id, g_test_mid, 16); memcpy(req.block_id, g_test_bid0, 16); req.chunk = 0; g_inst[0]->md.stream_completed = 0; @@ -264,23 +308,31 @@ int main(void) { } g_ua = uasync_create(); - for (int i = 0; i < N_NODES; i++) { + for (int i = 0; i < N_NODES; i++) wf(g_cfg[i], "[global]\ntun_ip=10.99.%d.1/24\ntun_ifname=tun%d0\ndb_path=%s\n" "[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", i, i, g_db_dir[i], g_port[i]); + for (int i = 0; i < N_NODES; i++) { config_ensure_keys_and_node_id(g_cfg[i]); struct utun_config* c = parse_config(g_cfg[i]); g_nid[i] = c->global.my_node_id; free_config(c); } for (int i = 0; i < N_NODES; i++) { char *pr = gv(g_cfg[i], "priv"), *pu = gv(g_cfg[i], "pub"); - int next = (i + 1) % N_NODES; char *n_pu = gv(g_cfg[next], "pub"); - char link[256]; snprintf(link, sizeof(link), - "[client: to_n%d]\nkeepalive=1\npeer_public_key=%s\nlink=s1:127.0.0.1:%d\n", next, n_pu, g_port[next]); - wf(g_cfg[i], "[global]\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.99.%d.1/24\ntun_ifname=tun%d0\n" - "db_path=%s\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n%s[allowed_keys]\nallow_all=1\n", - pr, pu, i, i, g_db_dir[i], g_port[i], link); - u_free(pr); u_free(pu); u_free(n_pu); + if (i == 0) { + wf(g_cfg[i], "[global]\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.99.%d.1/24\ntun_ifname=tun%d0\n" + "db_path=%s\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n[allowed_keys]\nallow_all=1\n", + pr, pu, i, i, g_db_dir[i], g_port[i]); + } else { + int prev = i - 1; char *pv_pu = gv(g_cfg[prev], "pub"); + char link[256]; snprintf(link, sizeof(link), + "[client: to_n%d]\nkeepalive=1\npeer_public_key=%s\nlink=s1:127.0.0.1:%d\n", prev, pv_pu, g_port[prev]); + wf(g_cfg[i], "[global]\nmy_private_key=%s\nmy_public_key=%s\ntun_ip=10.99.%d.1/24\ntun_ifname=tun%d0\n" + "db_path=%s\n[server: s1]\naddr=127.0.0.1:%d\ntype=public\n%s[allowed_keys]\nallow_all=1\n", + pr, pu, i, i, g_db_dir[i], g_port[i], link); + u_free(pv_pu); + } + u_free(pr); u_free(pu); } for (int i = 0; i < N_NODES; i++) { @@ -295,6 +347,10 @@ int main(void) { if (!g_connected) { printf("FAIL: connection timeout\n"); goto done; } } printf(" connected\n"); + g_group_id = strtoull(CH_ID, NULL, 10); + /* ══ Phase P1: create CHAT groups + members ══ */ + setup_chat_group(); + phase_a1_super_hello(); phase_a2_have_block_replication(); phase_a3_admission(); diff --git a/tests/test_media_delivery_sql.c b/tests/test_media_delivery_sql.c index 7fbac041..8250b80b 100644 --- a/tests/test_media_delivery_sql.c +++ b/tests/test_media_delivery_sql.c @@ -295,7 +295,7 @@ static void test_proto_sizes(void) { if (sizeof(struct media_pkt_have_block) != 117) { FAIL("have_block: %zu != 117", sizeof(struct media_pkt_have_block)); ok = 0; } if (sizeof(struct media_pkt_super_hello) != 9) { FAIL("super_hello: %zu != 9", sizeof(struct media_pkt_super_hello)); ok = 0; } if (sizeof(struct media_pkt_super_ack) != 5) { FAIL("super_ack: %zu != 5", sizeof(struct media_pkt_super_ack)); ok = 0; } - if (sizeof(struct media_pkt_block_req) != 45) { FAIL("block_req: %zu != 45", sizeof(struct media_pkt_block_req)); ok = 0; } + if (sizeof(struct media_pkt_block_req) != 53) { FAIL("block_req: %zu != 53", sizeof(struct media_pkt_block_req)); ok = 0; } if (sizeof(struct media_pkt_cancel) != 37) { FAIL("cancel: %zu != 37", sizeof(struct media_pkt_cancel)); ok = 0; } if (sizeof(struct media_pkt_block_done) != 105) { FAIL("block_done: %zu != 105", sizeof(struct media_pkt_block_done)); ok = 0; } if (sizeof(struct media_pkt_query_resp_entry) != 30) { FAIL("query_resp_entry: %zu != 30", sizeof(struct media_pkt_query_resp_entry)); ok = 0; } diff --git a/tools/chatgui-android/app/src/main/java/com/utun/chat/data/ChatRepository.kt b/tools/chatgui-android/app/src/main/java/com/utun/chat/data/ChatRepository.kt index 3990bdef..3bf70d19 100644 --- a/tools/chatgui-android/app/src/main/java/com/utun/chat/data/ChatRepository.kt +++ b/tools/chatgui-android/app/src/main/java/com/utun/chat/data/ChatRepository.kt @@ -9,7 +9,8 @@ data class Channel(val id: String, val name: String, val lastMsgAt: Long = 0, va data class Message(val id: Long, val author: String, val text: String, val ts: Long, val isOutgoing: Boolean = false, val contentType: String = "text/plain", val filePath: String = "", val waveform: List = emptyList(), - val voiceDurationMs: Int = 0, val fileSize: Long = 0) + val voiceDurationMs: Int = 0, val fileSize: Long = 0, + val downloadState: String = "" /* "", "downloading", "downloaded", "error" */) class ChatRepository { private val myNodeId: Long = NativeLib.getMyNodeId() @@ -56,6 +57,12 @@ class ChatRepository { var displayText = text var fileSize = 0L + /* Parse download state from localAttrs */ + var downloadState = "" + val localAttrs = obj.optString("localAttrs", "") + if (filePath.isNotEmpty()) downloadState = "downloaded" + else if (localAttrs.contains("\"st\":\"fl\"")) downloadState = "downloaded" + /* Parse base64 waveform from voice message text (format: "wf=<100 chars A-Za-z0-9+/>;dur=1234;") */ var waveform = emptyList() var voiceDurationMs = 0 @@ -114,7 +121,8 @@ class ChatRepository { filePath = filePath, waveform = waveform, voiceDurationMs = voiceDurationMs, - fileSize = fileSize + fileSize = fileSize, + downloadState = downloadState )) } } catch (e: Exception) { Log.w("utun-gui", "getMessages error", e) } diff --git a/tools/chatgui-android/app/src/main/java/com/utun/chat/data/NativeLib.kt b/tools/chatgui-android/app/src/main/java/com/utun/chat/data/NativeLib.kt index 4483aa43..94de7c30 100644 --- a/tools/chatgui-android/app/src/main/java/com/utun/chat/data/NativeLib.kt +++ b/tools/chatgui-android/app/src/main/java/com/utun/chat/data/NativeLib.kt @@ -99,6 +99,7 @@ object NativeLib { /* ── Attachment ── */ fun attachmentSend(channelId: String, filePath: String): Boolean = nativeAttachmentSend(channelId, filePath) + fun attachmentDownload(channelId: String, msgId: Long): Boolean = nativeAttachmentDownload(channelId, msgId) /* ── Voice playback ── */ fun voiceDecodeOpen(filePath: String, outInfo: IntArray): Long = nativeVoiceDecodeOpen(filePath, outInfo) @@ -147,6 +148,7 @@ object NativeLib { /* ── Attachment JNI ── */ private external fun nativeAttachmentSend(channelId: String, filePath: String): Boolean + private external fun nativeAttachmentDownload(channelId: String, msgId: Long): Boolean /* ── Voice playback JNI ── */ private external fun nativeVoiceDecodeOpen(filePath: String, outInfo: IntArray): Long 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 7405209a..b829509a 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 @@ -35,7 +35,8 @@ import kotlinx.coroutines.delay @Composable fun MessageBubble(message: Message, playbackState: ChatViewModel.PlaybackState = ChatViewModel.PlaybackState(), onPlayVoice: (String, Int) -> Unit = { _, _ -> }, - onSeekVoice: (String, Int, Int) -> Unit = { _, _, _ -> }) { + onSeekVoice: (String, Int, Int) -> Unit = { _, _, _ -> }, + onDownloadFile: (Message) -> Unit = {}) { val bgColor = if (message.isOutgoing) Color(0xFF4A90D9) else Color(0xFFE8E8E8) val textColor = if (message.isOutgoing) Color.White else Color.Black val alignment = if (message.isOutgoing) Arrangement.End else Arrangement.Start @@ -62,7 +63,7 @@ fun MessageBubble(message: Message, playbackState: ChatViewModel.PlaybackState = when { isVoice -> VoiceBubble(message, textColor, isActive, playbackState, onPlayVoice, onSeekVoice) - isMedia -> FileBubble(message, textColor) + isMedia -> FileBubble(message, textColor, onDownloadFile) else -> { Text(message.author, fontSize = 13.sp, fontWeight = FontWeight.Bold, color = textColor.copy(alpha = 0.7f)) Spacer(Modifier.height(2.dp)) @@ -168,11 +169,33 @@ private fun VoiceBubble(message: Message, textColor: Color, isActive: Boolean, } @Composable -private fun FileBubble(message: Message, textColor: Color) { - Column { - Text(message.text, fontSize = 15.sp, color = textColor, fontWeight = FontWeight.Bold) +private fun FileBubble(message: Message, textColor: Color, onDownload: (Message) -> Unit = {}) { + val state = message.downloadState + val stateIcon = when { + state == "downloaded" -> "\u2705" + state == "downloading" -> "\u23F3" + else -> "\u2B07\uFE0F" + } + val stateText = when { + state == "downloaded" -> "Downloaded" + state == "downloading" -> "Downloading..." + else -> "Tap to download" + } + + Column( + modifier = Modifier.clickable { if (state.isEmpty()) onDownload(message) } + ) { + Row(verticalAlignment = Alignment.CenterVertically) { + Text(stateIcon, fontSize = 18.sp) + Spacer(Modifier.width(6.dp)) + Text(message.text, fontSize = 15.sp, color = textColor, fontWeight = FontWeight.Bold) + } Spacer(Modifier.height(2.dp)) - Text(formatFileSize(message.fileSize), fontSize = 12.sp, color = textColor.copy(alpha = 0.5f)) + Text( + "${formatFileSize(message.fileSize)} · $stateText", + fontSize = 12.sp, + color = textColor.copy(alpha = 0.5f) + ) } } diff --git a/tools/chatgui-android/app/src/main/java/com/utun/chat/ui/screens/ChatScreen.kt b/tools/chatgui-android/app/src/main/java/com/utun/chat/ui/screens/ChatScreen.kt index 95782215..df86dc92 100644 --- a/tools/chatgui-android/app/src/main/java/com/utun/chat/ui/screens/ChatScreen.kt +++ b/tools/chatgui-android/app/src/main/java/com/utun/chat/ui/screens/ChatScreen.kt @@ -66,12 +66,9 @@ fun ChatScreen( ) { items(messages.size) { idx -> MessageBubble(messages[idx], playbackState = playbackState, - onPlayVoice = { path, durMs -> - viewModel.playVoiceMessage(path, durMs) - }, - onSeekVoice = { path, seekMs, durMs -> - viewModel.seekVoice(path, seekMs, durMs) - }) + onPlayVoice = { path, durMs -> viewModel.playVoiceMessage(path, durMs) }, + onSeekVoice = { path, seekMs, durMs -> viewModel.seekVoice(path, seekMs, durMs) }, + onDownloadFile = { msg -> viewModel.downloadAttachment(msg) }) } item { Spacer(Modifier.height(8.dp)) } } 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 06ac1a8f..17b60d47 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 @@ -123,6 +123,16 @@ class ChatViewModel : ViewModel() { } } 8 -> if (repo != null) refreshChannels() /* CHANNEL_PEERS_ONLINE */ + 14 -> { /* ATTACHMENT_DOWNLOADED: [ch_id_len:1][ch_id:var][msg_id:8] */ + if (repo == null || data == null || data.size < 2) return + val chLen = data[0].toInt() and 0xFF + if (data.size >= 1 + chLen + 8) { + val chId = String(data, 1, chLen) + val cur = _currentChannel.value + if (cur != null && cur.id == chId) refreshMessages(chId) + refreshChannels() + } + } 2 -> { /* CONNECT_RESULT: [node_id:8][result:4][channel_id:8] */ if (repo == null || data == null || data.size < 20) return val buf = java.nio.ByteBuffer.wrap(data).order(java.nio.ByteOrder.LITTLE_ENDIAN) @@ -279,6 +289,15 @@ class ChatViewModel : ViewModel() { } } + fun downloadAttachment(msg: Message) { + val ch = _currentChannel.value ?: return + NativeLib.attachmentDownload(ch.id, msg.id) + viewModelScope.launch { + delay(500) + refreshMessages(ch.id) + } + } + fun getVoicePreset(): Int = NativeLib.voiceGetPreset() fun setVoicePreset(preset: Int) { NativeLib.voiceSetPreset(preset) } fun getVoiceCompressor(): Boolean = NativeLib.voiceGetCompressor() diff --git a/tools/chatgui-android/jni_bridge/android_jni_bridge.c b/tools/chatgui-android/jni_bridge/android_jni_bridge.c index 02163167..c6993df5 100644 --- a/tools/chatgui-android/jni_bridge/android_jni_bridge.c +++ b/tools/chatgui-android/jni_bridge/android_jni_bridge.c @@ -409,15 +409,19 @@ char* utun_bridge_get_messages_json(const char* channel_id, int limit) { } } char* esc_fp = file_path[0] ? json_escape_alloc(file_path) : NULL; + char* esc_la = lattr && lattr[0] ? json_escape_alloc(lattr) : NULL; - size_t needed = snprintf(NULL, 0, "%s{\"id\":%lld,\"author\":\"%s\",\"text\":\"%s\",\"ts\":%lld,\"isOutgoing\":%s,\"contentType\":\"%s\"%s%s%s}", + size_t needed = snprintf(NULL, 0, "%s{\"id\":%lld,\"author\":\"%s\",\"text\":\"%s\",\"ts\":%lld,\"isOutgoing\":%s,\"contentType\":\"%s\"%s%s%s%s%s%s}", sep, (long long)msg_id, author_name, esc_txt, (long long)ts, is_out ? "true" : "false", esc_ct, - esc_fp ? ",\"filePath\":\"" : "", esc_fp ? esc_fp : "", esc_fp ? "\"" : ""); + esc_fp ? ",\"filePath\":\"" : "", esc_fp ? esc_fp : "", esc_fp ? "\"" : "", + esc_la ? ",\"localAttrs\":\"" : "", esc_la ? esc_la : "", esc_la ? "\"" : ""); while (pos + needed + 2 > cap) { cap *= 2; char* tmp = u_realloc(json, cap); if (!tmp) break; json = tmp; } if (pos + needed + 2 <= cap) - pos += (size_t)snprintf(json + pos, cap - pos, "%s{\"id\":%lld,\"author\":\"%s\",\"text\":\"%s\",\"ts\":%lld,\"isOutgoing\":%s,\"contentType\":\"%s\"%s%s%s}", + pos += (size_t)snprintf(json + pos, cap - pos, "%s{\"id\":%lld,\"author\":\"%s\",\"text\":\"%s\",\"ts\":%lld,\"isOutgoing\":%s,\"contentType\":\"%s\"%s%s%s%s%s%s}", sep, (long long)msg_id, author_name, esc_txt, (long long)ts, is_out ? "true" : "false", esc_ct, - esc_fp ? ",\"filePath\":\"" : "", esc_fp ? esc_fp : "", esc_fp ? "\"" : ""); + esc_fp ? ",\"filePath\":\"" : "", esc_fp ? esc_fp : "", esc_fp ? "\"" : "", + esc_la ? ",\"localAttrs\":\"" : "", esc_la ? esc_la : "", esc_la ? "\"" : ""); + if (esc_la) u_free(esc_la); if (esc_fp) u_free(esc_fp); u_free(esc_ct); u_free(esc_txt); @@ -510,6 +514,19 @@ int utun_bridge_attachment_send(const char* channel_id, const char* file_path, c return attachment_send(channel_id, file_path, db_path); } +int utun_bridge_attachment_download(const char* channel_id, int64_t msg_id) { + if (!channel_id || msg_id <= 0) { bridge_log(BLEV_ERROR, "attachment_download: invalid args"); return -1; } + struct UASYNC* ua = instance_lite_get_uasync(); + if (!ua) { bridge_log(BLEV_ERROR, "attachment_download: no uasync"); return -1; } + struct attachment_dl_req* req = u_calloc(1, sizeof(*req)); + if (!req) return -1; + snprintf(req->channel_id, sizeof(req->channel_id), "%s", channel_id); + req->msg_id = msg_id; + bridge_log(BLEV_INFO, "attachment_download: ch=%s msg_id=%lld", channel_id, (long long)msg_id); + uasync_post(ua, chat_core_attachment_download_trampoline, req); + return 0; +} + /* ────────────────────────────────────────────────────────────────── * Voice playback (Opus decode) * ────────────────────────────────────────────────────────────────── */ @@ -1070,6 +1087,15 @@ JNIEXPORT jboolean JNICALL Java_com_utun_chat_data_NativeLib_nativeAttachmentSen return (rc == 0) ? JNI_TRUE : JNI_FALSE; } +JNIEXPORT jboolean JNICALL Java_com_utun_chat_data_NativeLib_nativeAttachmentDownload( + JNIEnv* env, jobject thiz, jstring channelId, jlong msgId) { + (void)thiz; + const char* ch = (*env)->GetStringUTFChars(env, channelId, NULL); + int rc = utun_bridge_attachment_download(ch, (int64_t)msgId); + if (ch) (*env)->ReleaseStringUTFChars(env, channelId, ch); + return (rc == 0) ? JNI_TRUE : JNI_FALSE; +} + /* ── Voice playback JNI ── */ JNIEXPORT jlong JNICALL Java_com_utun_chat_data_NativeLib_nativeVoiceDecodeOpen( diff --git a/tools/chatgui-android/jni_bridge/android_jni_bridge.h b/tools/chatgui-android/jni_bridge/android_jni_bridge.h index 8f96a18e..2bb296a1 100644 --- a/tools/chatgui-android/jni_bridge/android_jni_bridge.h +++ b/tools/chatgui-android/jni_bridge/android_jni_bridge.h @@ -92,6 +92,7 @@ float utun_bridge_voice_get_peak_level(void); /* ── Attachment ── */ int utun_bridge_attachment_send(const char* channel_id, const char* file_path, const char* db_path); +int utun_bridge_attachment_download(const char* channel_id, int64_t msg_id); /* ── Voice playback (Opus decode) ── */ diff --git a/tools/chatgui/src/mainwindow.cpp b/tools/chatgui/src/mainwindow.cpp index 4e86c9c3..a718e5da 100644 --- a/tools/chatgui/src/mainwindow.cpp +++ b/tools/chatgui/src/mainwindow.cpp @@ -256,6 +256,14 @@ void MainWindow::setupBridgeCallbacks() { s_mainWindow->m_messageList->loadChannel(lastCid); } }); + 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(); + }); } void MainWindow::setupConnects() { diff --git a/tools/chatgui/src/messagedelegate.cpp b/tools/chatgui/src/messagedelegate.cpp index 1e85d614..ca1c1190 100644 --- a/tools/chatgui/src/messagedelegate.cpp +++ b/tools/chatgui/src/messagedelegate.cpp @@ -20,6 +20,14 @@ #include #include +extern "C" { +#include "../transport/gui_bridge.h" +#include "chat/chat_core.h" +#include "../../lib/u_async.h" +#include "../../lib/mem.h" +void chat_core_attachment_download_trampoline(void* arg); +} + MessageDelegate::MessageDelegate(QObject *parent) : QStyledItemDelegate(parent) {} @@ -67,6 +75,24 @@ bool MessageDelegate::editorEvent(QEvent *event, QAbstractItemModel *model, const QModelIndex &index) { if (event->type() != QEvent::MouseButtonRelease) return false; Layout L = calcLayout(option, index); + + 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; + + 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); + return true; + } + if (!L.isVoice) return false; QMouseEvent* me = static_cast(event); @@ -797,17 +823,22 @@ void MessageDelegate::paint(QPainter *painter, const QStyleOptionViewItem &optio return; } - if (L.isFile) { + if (L.isFile) { drawBubble(painter, L, option); if (!L.narrow) drawAvatar(painter, L, index); - /* File icon */ + int dlState = index.data(MsgFileDownloadStateRole).toInt(); + + /* File icon based on download state: 0=not_started, 1=downloading, 2=downloaded */ QFont iconFont = option.font; iconFont.setPointSize(iconFont.pointSize() + 6); painter->setFont(iconFont); painter->setPen(option.palette.text().color()); QRect iconRect(L.bubbleRect.left() + kPadH, kPadTop, 24, 24); - painter->drawText(iconRect, Qt::AlignCenter, QString::fromUtf8("\xF0\x9F\x93\x84")); + if (dlState == 2) + painter->drawText(iconRect, Qt::AlignCenter, QString::fromUtf8("\xE2\x9C\x85")); + else + painter->drawText(iconRect, Qt::AlignCenter, QString::fromUtf8("\xE2\xAC\x87\xEF\xB8\x8F")); QFont fileFont = option.font; fileFont.setBold(true); @@ -822,19 +853,16 @@ 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 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, sizeStr); - - /* Timestamp */ - QFont tsFont = option.font; - tsFont.setPointSize(tsFont.pointSize() - 1); - painter->setFont(tsFont); - painter->drawText(QRect(L.textRect.right() - 40, L.textRect.top() + nameH, 40, QFontMetrics(sizeFont).height()), - Qt::AlignRight | Qt::AlignVCenter, index.data(MsgTimeRole).toString()); + painter->drawText(sizeRect, Qt::AlignLeft | Qt::AlignVCenter, info); painter->restore(); return; + + } auto *cv = qobject_cast(option.widget); diff --git a/tools/chatgui/src/messagedelegate.h b/tools/chatgui/src/messagedelegate.h index e4d56e67..0255fd42 100644 --- a/tools/chatgui/src/messagedelegate.h +++ b/tools/chatgui/src/messagedelegate.h @@ -27,6 +27,8 @@ enum MessageDataRole { MsgVoiceSigsRole = Qt::UserRole + 22, MsgFileDisplayNameRole = Qt::UserRole + 23, MsgFileSizeRole = Qt::UserRole + 24, + MsgFileDownloadStateRole = Qt::UserRole + 25, + MsgChannelIdRole = Qt::UserRole + 26, }; class MessageDelegate : public QStyledItemDelegate { diff --git a/tools/chatgui/src/messagelist.cpp b/tools/chatgui/src/messagelist.cpp index 70728043..950ee790 100644 --- a/tools/chatgui/src/messagelist.cpp +++ b/tools/chatgui/src/messagelist.cpp @@ -90,12 +90,13 @@ static QString contentTypeFromJson(const QByteArray& data) { } static void setVoiceMessageRoles(QStandardItem* item, const QByteArray& data, - const QString& contentType, const QString& mediaDirBase, - const QString& channelId, const QByteArray& localAttrs) { + const QString& contentType, const QString& mediaDirBase, + const QString& channelId, const QByteArray& localAttrs) { QString ct = contentType; if (ct.isEmpty()) ct = contentTypeFromJson(data); if (ct.isEmpty()) ct = "text/plain"; item->setData(ct, MsgContentTypeRole); + item->setData(channelId, MsgChannelIdRole); if (ct == "audio/opus") { QString raw = QString::fromUtf8(data); @@ -166,6 +167,18 @@ static void setVoiceMessageRoles(QStandardItem* item, const QByteArray& data, if (parts.size() >= 1) item->setData(parts[0].toULongLong(), MsgFileSizeRole); } + + /* parse download state from local_attrs */ + QJsonObject la = QJsonDocument::fromJson(localAttrs).object(); + QString st = la.value("st").toString(); + QString fp = la.value("fp").toString(); + int dlState = 0; /* 0=not started, 2=downloaded */ + if (st == "fl" && !fp.isEmpty()) { + dlState = 2; /* downloaded */ + QString fullPath = mediaDirBase + "/" + channelId + "/" + fp; + item->setData(fullPath, MsgVoiceFileRole); + } + item->setData(dlState, MsgFileDownloadStateRole); } } diff --git a/tools/chatgui/transport/gui_bridge.h b/tools/chatgui/transport/gui_bridge.h index 1ceca439..91a73267 100644 --- a/tools/chatgui/transport/gui_bridge.h +++ b/tools/chatgui/transport/gui_bridge.h @@ -23,6 +23,7 @@ struct UASYNC; #define GUI_EVT_DB_READY 9 /* data: none — DB is open, tables created */ #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] */ /* ── API ── */ @@ -78,6 +79,10 @@ void gui_bridge_set_db_ready_cb(gui_db_ready_fn cb); typedef void (*gui_status_refresh_fn)(const char* text, int len); void gui_bridge_set_status_refresh_cb(gui_status_refresh_fn cb); +/* Callback для завершения загрузки медиа-аттача */ +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); + #ifdef __cplusplus } #endif diff --git a/tools/chatgui/transport/gui_bridge_impl.cpp b/tools/chatgui/transport/gui_bridge_impl.cpp index 8f5f6be7..ee8610db 100644 --- a/tools/chatgui/transport/gui_bridge_impl.cpp +++ b/tools/chatgui/transport/gui_bridge_impl.cpp @@ -33,6 +33,7 @@ static gui_auto_connect_status_fn g_auto_connect_status_cb = nullptr; 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 struct UASYNC* g_ua = nullptr; /* ── GuiBridgeReceiver implementation ── */ @@ -141,6 +142,16 @@ void GuiBridgeReceiver::processPost(int eventType, QByteArray data) { DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "gui_bridge: NODE_CHANGED data too short %d", dlen); } break; + case GUI_EVT_ATTACHMENT_DOWNLOADED: + if (dlen >= 2) { + uint8_t chLen = d[0]; + if (dlen >= 1 + chLen + 8) { + int64_t msgId = 0; memcpy(&msgId, d + 1 + chLen, 8); + DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "gui_bridge: ATTACHMENT_DOWNLOADED ch=%.*s msgId=%lld", chLen, (const char*)d + 1, (long long)msgId); + if (g_attachment_downloaded_cb) g_attachment_downloaded_cb((const char*)d + 1, chLen, msgId); + } + } + break; default: DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "gui_bridge: unknown event type %d", eventType); break; @@ -223,6 +234,10 @@ void gui_bridge_set_status_refresh_cb(gui_status_refresh_fn cb) { g_status_refresh_cb = cb; } +void gui_bridge_set_attachment_downloaded_cb(gui_attachment_downloaded_fn cb) { + g_attachment_downloaded_cb = cb; +} + } /* extern "C" */ #include "gui_bridge_impl.moc"