Browse Source

feat: media_async + media_index modules, media_files table, tests

- media_async: thread-per-task async engine with primitives (sha256, sign, copy, uuid, file_size)
- media_index: media_files DB table with register_async for indexing chunks
- chat_msg: async media submission via register_async, simplified chat_msg_submit fields
- voice_recorder/attachment_sender: write single file (no chunk splitting on caller side)
- inputbar/messagelist: updated to new simplified submit interface
- test_media_async: 4-thread stress test (200K ops / 2 sec, 0 errors)
- test_media_index: 9 cases covering commit edge cases + async flow
- build: Makefile.am, CMakeLists.txt, utun_sources.cmake updated
topo_upd
evgeny 2 months ago
parent
commit
65dbfa6813
  1. 6
      lib/audio_compressor.c
  2. 2
      lib/audio_compressor.h
  3. 8
      src/Makefile.am
  4. 4
      src/chat/chat_core.c
  5. 8
      src/chat/chat_core.h
  6. 233
      src/chat/chat_msg.c
  7. 154
      src/media_async/media_async.c
  8. 36
      src/media_async/media_async.h
  9. 41
      src/media_delivery/media_delivery.c
  10. 27
      src/media_delivery/media_delivery.h
  11. 297
      src/media_delivery/media_index.c
  12. 52
      src/media_delivery/media_index.h
  13. 1
      src/transport_layer/etcp_api.h
  14. 20
      src/utun_instance.c
  15. 8
      src/utun_instance.h
  16. 10
      tests/Makefile.am
  17. 8
      tests/test_audio_compressor.c
  18. 299
      tests/test_media_async.c
  19. 510
      tests/test_media_index.c
  20. 23
      tools/chatgui-android/app/src/main/java/com/utun/chat/data/ChatRepository.kt
  21. 9
      tools/chatgui-android/app/src/main/java/com/utun/chat/data/ConfigProvider.kt
  22. 22
      tools/chatgui-android/app/src/main/java/com/utun/chat/ui/components/MessageBubble.kt
  23. 49
      tools/chatgui-android/app/src/main/java/com/utun/chat/ui/screens/SettingsScreen.kt
  24. 4
      tools/chatgui-android/jni_bridge/android_jni_bridge.c
  25. 2
      tools/chatgui-android/jni_bridge/android_jni_bridge.h
  26. 75
      tools/chatgui-android/libutun_lite/attachment_sender.c
  27. 2
      tools/chatgui-android/libutun_lite/utun_sources.cmake
  28. 68
      tools/chatgui-android/libutun_lite/voice_recorder.c
  29. 2
      tools/chatgui-android/libutun_lite/voice_recorder.h
  30. 1
      tools/chatgui/CMakeLists.txt
  31. 2
      tools/chatgui/libutun/CMakeLists.txt
  32. 65
      tools/chatgui/src/audiodevicesettingspage.cpp
  33. 13
      tools/chatgui/src/audiodevicesettingspage.h
  34. 35
      tools/chatgui/src/inputbar.cpp
  35. 8
      tools/chatgui/src/inputbar.h
  36. 12
      tools/chatgui/src/mainwindow.cpp
  37. 79
      tools/chatgui/src/messagedelegate.cpp
  38. 2
      tools/chatgui/src/messagedelegate.h
  39. 39
      tools/chatgui/src/messagelist.cpp
  40. 18
      tools/chatgui/src/settingsdialog.cpp
  41. 2
      tools/chatgui/src/settingsdialog.h
  42. 90
      tools/chatgui/src/storagesettingspage.cpp
  43. 22
      tools/chatgui/src/storagesettingspage.h

6
lib/audio_compressor.c

@ -78,14 +78,14 @@ void audio_compressor_configure(struct audio_compressor* ac, const audio_compres
ac->block_samples = ac->cfg.sample_rate * ac->cfg.block_duration_ms / 1000 * ac->cfg.channels;
ac->lookback_blocks = ac->cfg.lookback_ms / ac->cfg.block_duration_ms;
ac->lookahead_blocks = ac->cfg.lookahead_ms / ac->cfg.block_duration_ms;
ac->rise_factor_per_block = powf(ac->cfg.rise_rate_per_500ms, (float)ac->cfg.block_duration_ms / 500.0f);
ac->rise_factor_per_block = powf(ac->cfg.rise_rate_per_sec, (float)ac->cfg.block_duration_ms / 1000.0f);
ac->max_gain = powf(10.0f, ac->cfg.max_gain_db / 20.0f);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL,
"audio_compressor_configure: rate=%d ch=%d block=%dms lookback=%dms lookahead=%dms maxGain=%.0fdB riseRate=%.1f/500ms target=%.0fdBFS",
"audio_compressor_configure: rate=%d ch=%d block=%dms lookback=%dms lookahead=%dms maxGain=%.0fdB riseRate=%.1f/sec target=%.0fdBFS",
ac->cfg.sample_rate, ac->cfg.channels, ac->cfg.block_duration_ms,
ac->cfg.lookback_ms, ac->cfg.lookahead_ms, ac->cfg.max_gain_db,
ac->cfg.rise_rate_per_500ms, 20.0f * log10f(ac->cfg.target_level));
ac->cfg.rise_rate_per_sec, 20.0f * log10f(ac->cfg.target_level));
}
void audio_compressor_reset(struct audio_compressor* ac) {

2
lib/audio_compressor.h

@ -17,7 +17,7 @@ typedef struct {
int lookback_ms;
int lookahead_ms;
float max_gain_db;
float rise_rate_per_500ms;
float rise_rate_per_sec;
float target_level;
} audio_compressor_config_t;

8
src/Makefile.am

@ -45,6 +45,9 @@ utun_CORE_SOURCES = \
firewall.c \
eim_nat.c \
nat_transport.c \
media_delivery/media_delivery.c \
media_delivery/media_index.c \
media_async/media_async.c \
transport_layer/dummynet.c \
ntp_time.c \
ntp_node_time.c \
@ -111,6 +114,9 @@ libutun_a_SOURCES = \
firewall.c \
eim_nat.c \
nat_transport.c \
media_delivery/media_delivery.c \
media_delivery/media_index.c \
media_async/media_async.c \
transport_layer/dummynet.c \
ntp_time.c \
ntp_node_time.c \
@ -146,6 +152,8 @@ utun_CFLAGS = \
-I$(top_srcdir)/src \
-I$(top_srcdir)/src/transport_layer \
-I$(top_srcdir)/src/routing_layer \
-I$(top_srcdir)/src/media_delivery \
-I$(top_srcdir)/src/media_async \
-I$(top_srcdir)/src/uip \
-g \
-DUSE_SQLITE \

4
src/chat/chat_core.c

@ -16,6 +16,7 @@
#include "../utun_instance.h"
#include "../ntp_time.h"
#include "../routing_layer/topo_node_sqlite.h"
#include "../media_delivery/media_index.h"
#include "../../lib/mem.h"
#include <string.h>
@ -135,6 +136,9 @@ int chat_core_init(struct UTUN_INSTANCE* inst, const char* db_path) {
chat_core_sync_my_addresses();
/* create media_files table */
media_index_init(g_cc.db);
db_exec(
"CREATE TABLE IF NOT EXISTS local_identity ("
" id INTEGER PRIMARY KEY CHECK (id = 1),"

8
src/chat/chat_core.h

@ -28,11 +28,9 @@ int chat_core_is_initialized(void);
struct chat_msg_submit {
char channel_id[64];
char content_type[32];
char media_dt[32];
char media_basename[256];
char media_suffix[8];
char media_ext[16];
uint32_t media_num_blocks;
char media_src[1024];
char media_dest[1024];
uint8_t media_copy;
uint8_t* data;
uint32_t data_len;
uint64_t timestamp;

233
src/chat/chat_msg.c

@ -9,16 +9,11 @@
#include "../utun_instance.h"
#include "../transport_layer/secure_channel.h"
#include "../media_delivery/media_index.h"
#include "../../lib/mem.h"
#include "../../lib/platform_compat.h"
#include <openssl/sha.h>
#include <unistd.h>
#include <sys/stat.h>
#define MEDIA_BLOCK_MIN (10 * 1024 * 1024)
#define MEDIA_BLOCK_MAX (25 * 1024 * 1024)
#define MEDIA_BLOCK_TARGET 30
/* forward decl */
static void chat_core_submit_media_message(struct chat_msg_submit* req);
@ -72,174 +67,140 @@ void chat_core_submit_message(struct chat_msg_submit* req) {
void chat_core_submit_trampoline(void* arg) {
struct chat_msg_submit* req = (struct chat_msg_submit*)arg;
if (req->media_num_blocks > 0)
if (req->media_src[0] != '\0')
chat_core_submit_media_message(req);
else
chat_core_submit_message(req);
u_free(req);
}
/* ─── Медиа-сообщения: подпись блоков + сборка в единый файл ─── */
static uint32_t media_calc_block_size(uint64_t file_size) {
if (file_size == 0) return 0;
uint64_t target = (file_size + MEDIA_BLOCK_TARGET - 1) / MEDIA_BLOCK_TARGET;
if (target < MEDIA_BLOCK_MIN) target = MEDIA_BLOCK_MIN;
if (target > MEDIA_BLOCK_MAX) target = MEDIA_BLOCK_MAX;
return (uint32_t)target;
}
static void chat_core_submit_media_message(struct chat_msg_submit* req) {
if (!g_cc.initialized) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: media submit before init", CC_ID); return; }
if (!req || req->media_num_blocks == 0) return;
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: MEDIA SUBMIT ch=%s ct=%s dt=%s stem=%s ext=%s blocks=%u",
CC_ID, req->channel_id, req->content_type,
req->media_dt, req->media_basename, req->media_ext, req->media_num_blocks);
/* ─── Медиа-сообщения: async регистрация в media_files + отправка ─── */
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: db_path='%s' initialized=%d",
CC_ID, g_cc.db_path, g_cc.initialized);
struct media_submit_ctx {
char channel_id[64];
char content_type[32];
char media_base[512];
uint8_t msg_data[4096];
uint32_t msg_data_len;
uint64_t timestamp;
};
struct DB_SYNC_INSTANCE* si = si_find(req->channel_id);
if (!si) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: no db_sync instance for ch=%s", CC_ID, req->channel_id); return; }
static void on_media_registered(void* arg, int err, const struct media_index_result* result) {
struct media_submit_ctx* mctx = (struct media_submit_ctx*)arg;
char media_dir[1024];
{
/* g_cc.db_path is a file path (e.g. /data/.../chats.db).
Extract directory portion for media path. */
const char* last_slash = strrchr(g_cc.db_path, '/');
if (last_slash)
snprintf(media_dir, sizeof(media_dir), "%.*s/media/%s",
(int)(last_slash - g_cc.db_path), g_cc.db_path, req->channel_id);
else
snprintf(media_dir, sizeof(media_dir), "%s/media/%s", g_cc.db_path, req->channel_id);
if (err || !result) {
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: media registration failed ch=%s err=%d",
CC_ID, mctx->channel_id, err);
u_free(mctx);
return;
}
DEBUG_DEBUG(DEBUG_CATEGORY_DB_SYNC, "%s: media_dir='%s' dt='%s' stem='%s' ext='%s'",
CC_ID, media_dir, req->media_dt, req->media_basename, req->media_ext);
uint32_t num_blocks = req->media_num_blocks;
uint64_t block_size = 0;
uint64_t file_size = 0;
/* allocate sig array */
uint8_t (*sigs)[64] = (uint8_t(*)[64])u_malloc(num_blocks * 64);
if (!sigs) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: sig alloc failed blocks=%u", CC_ID, num_blocks); return; }
/* read blocks, sign, compute sizes */
for (uint32_t n = 0; n < num_blocks; n++) {
char path[1280];
snprintf(path, sizeof(path), "%s/%s_%u_%s_%s.%s",
media_dir, req->media_dt, n, req->media_basename, req->media_suffix, req->media_ext);
FILE* f = fopen(path, "rb");
if (!f) {
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: cannot open block %u: %s", CC_ID, n, path);
u_free(sigs); return;
}
fseek(f, 0, SEEK_END); long fsz = ftell(f); fseek(f, 0, SEEK_SET);
if (fsz <= 0) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: empty block %u: %s", CC_ID, n, path); fclose(f); u_free(sigs); return; }
uint8_t* buf = (uint8_t*)u_malloc((size_t)fsz);
if (!buf || fread(buf, 1, (size_t)fsz, f) != (size_t)fsz) {
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: read fail block %u", CC_ID, n);
if (buf) u_free(buf); fclose(f); u_free(sigs); return;
}
fclose(f);
if (sc_ed25519_sign(g_cc.inst->my_ed25519_privkey, buf, (size_t)fsz, sigs[n]) != SC_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: Ed25519 sign fail block %u", CC_ID, n);
u_free(buf); u_free(sigs); return;
}
u_free(buf);
file_size += (uint64_t)fsz;
if (n == 0) block_size = (uint64_t)fsz;
struct DB_SYNC_INSTANCE* si = si_find(mctx->channel_id);
if (!si) {
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: no db_sync instance for ch=%s", CC_ID, mctx->channel_id);
u_free(mctx);
return;
}
if (file_size == 0) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: zero file size", CC_ID); u_free(sigs); return; }
/* build final data: <base>|<file_size>|<block_size>|<num_blocks>|<sig_hex>,... */
size_t base_len = req->data_len;
size_t sigs_hex_len = num_blocks * 129; /* 128 hex + comma */
size_t full_data_cap = base_len + 128 + sigs_hex_len;
char* full_data = (char*)u_malloc(full_data_cap + 1);
if (!full_data) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: data alloc fail", CC_ID); u_free(sigs); return; }
/* build full_data: <base>|<file_size>|<block_size>|<num_blocks>|<media_id_hex>|<content_hash_hex>|<id0,sig0,...> */
size_t base_len = mctx->msg_data_len;
int nb = result->num_blocks;
size_t pairs_len = (size_t)nb * (32 + 1 + 128 + 1); /* id_hex, sig_hex, commas */
size_t extra = 128 + 32 + 64 + pairs_len + 256;
size_t fdcap = base_len + extra;
char* full_data = u_malloc(fdcap);
if (!full_data) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: data alloc fail", CC_ID); u_free(mctx); return; }
size_t off = 0;
if (base_len > 0) { memcpy(full_data + off, req->data, base_len); off += base_len; }
off += snprintf(full_data + off, full_data_cap - off, "|%llu|%llu|%u|",
(unsigned long long)file_size, (unsigned long long)block_size, num_blocks);
for (uint32_t n = 0; n < num_blocks; n++) {
if (base_len > 0) { memcpy(full_data + off, mctx->msg_data, base_len); off += base_len; }
off += snprintf(full_data + off, fdcap - off, "|%lld|%lld|%d|",
(long long)result->file_size, (long long)result->block_size, nb);
for (int i = 0; i < 16; i++) off += snprintf(full_data + off, fdcap - off, "%02x", result->media_id[i]);
off += snprintf(full_data + off, fdcap - off, "|");
for (int i = 0; i < 32; i++) off += snprintf(full_data + off, fdcap - off, "%02x", result->content_hash[i]);
off += snprintf(full_data + off, fdcap - off, "|");
for (int n = 0; n < nb; n++) {
if (n > 0) full_data[off++] = ',';
for (int b = 0; b < 64; b++)
off += snprintf(full_data + off, full_data_cap - off, "%02x", sigs[n][b]);
for (int i = 0; i < 16; i++)
off += snprintf(full_data + off, fdcap - off, "%02x", result->block_ids[n * 16 + i]);
full_data[off++] = ',';
for (int i = 0; i < 64; i++)
off += snprintf(full_data + off, fdcap - off, "%02x", result->block_sigs[n * 64 + i]);
}
size_t full_data_len = off;
/* build JSON with dynamic buffer */
size_t json_cap = full_data_len + 256;
char* json = (char*)u_malloc(json_cap);
if (!json) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: json alloc fail", CC_ID); u_free(full_data); u_free(sigs); return; }
char* json = u_malloc(json_cap);
if (!json) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: json alloc fail", CC_ID); u_free(full_data); u_free(mctx); return; }
snprintf(json, json_cap,
"{\"n\":%llu,\"ch\":\"%s\",\"ct\":\"%s\",\"d\":\"%.*s\"}",
(unsigned long long)g_cc.my_node_id, req->channel_id,
req->content_type, (int)full_data_len, full_data);
(unsigned long long)g_cc.my_node_id, mctx->channel_id,
mctx->content_type, (int)full_data_len, full_data);
uint64_t ts = db_sync_next_timestamp(si);
size_t sig_msg_len = 8 + strlen(json);
uint8_t* sig_msg = (uint8_t*)u_malloc(sig_msg_len);
if (!sig_msg) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: sig_msg alloc fail", CC_ID); u_free(json); u_free(full_data); u_free(sigs); return; }
uint8_t* sig_msg = u_malloc(sig_msg_len);
if (!sig_msg) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: sig_msg alloc fail", CC_ID); u_free(json); u_free(full_data); u_free(mctx); return; }
memcpy(sig_msg, &ts, 8);
memcpy(sig_msg + 8, json, sig_msg_len - 8);
uint8_t json_sig[64];
if (sc_ed25519_sign(g_cc.inst->my_ed25519_privkey, sig_msg, sig_msg_len, json_sig) != SC_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: Ed25519 sign JSON failed", CC_ID);
u_free(sig_msg); u_free(json); u_free(full_data); u_free(sigs); return;
u_free(sig_msg); u_free(json); u_free(full_data); u_free(mctx); return;
}
/* assemble single file from blocks, then delete blocks */
{
char assembled_path[1280];
snprintf(assembled_path, sizeof(assembled_path), "%s/%s_%s_%s.%s",
media_dir, req->media_dt, req->media_basename, req->media_suffix, req->media_ext);
FILE* af = fopen(assembled_path, "wb");
if (af) {
for (uint32_t n = 0; n < num_blocks; n++) {
char block_path[1280];
snprintf(block_path, sizeof(block_path), "%s/%s_%u_%s_%s.%s",
media_dir, req->media_dt, n, req->media_basename, req->media_suffix, req->media_ext);
FILE* bf = fopen(block_path, "rb");
if (bf) {
uint8_t buf[65536]; size_t rd;
while ((rd = fread(buf, 1, sizeof(buf), bf)) > 0) fwrite(buf, 1, rd, af);
fclose(bf);
unlink(block_path);
}
}
fclose(af);
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: assembled %s (%llu bytes)", CC_ID, assembled_path, (unsigned long long)file_size);
} else {
DEBUG_WARN(DEBUG_CATEGORY_DB_SYNC, "%s: cannot create assembled file %s", CC_ID, assembled_path);
}
int ret = db_sync_insert_signed(si, json, strlen(json), json_sig, 64, ts, NULL);
if (ret != 0)
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: media insert failed ret=%d", CC_ID, ret);
else
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: media message inserted ch=%s ts=%llu",
CC_ID, mctx->channel_id, (unsigned long long)ts);
u_free(sig_msg); u_free(json); u_free(full_data); u_free(mctx);
}
static void chat_core_submit_media_message(struct chat_msg_submit* req) {
if (!g_cc.initialized) { DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: media submit before init", CC_ID); return; }
if (!req || req->media_src[0] == '\0') return;
DEBUG_INFO(DEBUG_CATEGORY_DB_SYNC, "%s: MEDIA SUBMIT ch=%s ct=%s src=%s dst=%s copy=%d",
CC_ID, req->channel_id, req->content_type,
req->media_src, req->media_dest, req->media_copy);
if (!g_cc.inst->media_async) {
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: media_async not initialized", CC_ID);
return;
}
/* build local_attrs JSON */
char local_attrs_json[512];
/* build media_base from db_path: dirname(db_path) e.g. /path/to/data */
char media_base[1024];
{
char bl_list[256] = "0"; size_t bl_off = 1;
for (uint32_t n = 1; n < num_blocks; n++)
bl_off += snprintf(bl_list + bl_off, sizeof(bl_list) - bl_off, ",%u", n);
snprintf(local_attrs_json, sizeof(local_attrs_json),
"{\"st\":\"fl\",\"fp\":\"%s_%s_%s.%s\",\"bl\":[%s]}",
req->media_dt, req->media_basename, req->media_suffix, req->media_ext, bl_list);
const char* last_slash = strrchr(g_cc.db_path, '/');
if (last_slash)
snprintf(media_base, sizeof(media_base), "%.*s", (int)(last_slash - g_cc.db_path), g_cc.db_path);
else
snprintf(media_base, sizeof(media_base), "%s", g_cc.db_path);
}
int ret = db_sync_insert_signed(si, json, strlen(json), json_sig, 64, ts, local_attrs_json);
if (ret != 0)
DEBUG_ERROR(DEBUG_CATEGORY_DB_SYNC, "%s: media insert failed ret=%d", CC_ID, ret);
/* copy fields for async callback (req is freed by trampoline immediately) */
struct media_submit_ctx* mctx = u_calloc(1, sizeof(*mctx));
if (!mctx) return;
snprintf(mctx->channel_id, sizeof(mctx->channel_id), "%s", req->channel_id);
snprintf(mctx->content_type, sizeof(mctx->content_type), "%s", req->content_type);
snprintf(mctx->media_base, sizeof(mctx->media_base), "%s", media_base);
mctx->timestamp = req->timestamp;
if (req->data && req->data_len > 0) {
mctx->msg_data_len = req->data_len < sizeof(mctx->msg_data) ? req->data_len : sizeof(mctx->msg_data) - 1;
memcpy(mctx->msg_data, req->data, mctx->msg_data_len);
}
u_free(sig_msg); u_free(json); u_free(full_data); u_free(sigs);
media_index_register_async(
g_cc.inst->media_async, g_cc.inst->ua, g_cc.db,
g_cc.my_node_id, g_cc.inst->my_ed25519_privkey,
req->channel_id,
req->media_src, req->media_dest, req->media_copy ? 1 : 0,
media_base,
on_media_registered, mctx);
}
/* ─── Обновление local_attrs (для будущей приёмной стороны) ─── */

154
src/media_async/media_async.c

@ -0,0 +1,154 @@
// media_async.c — thread-per-task async engine + crypto/file primitives
#include "media_async.h"
#include "../../lib/u_async.h"
#include "../../lib/mem.h"
#include "../../lib/debug_config.h"
#include <openssl/sha.h>
#include <openssl/evp.h>
#include <openssl/rand.h>
#include <pthread.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <sys/stat.h>
struct media_async {
int initialized;
};
struct ma_thread_ctx {
struct UASYNC* ua;
ma_work_fn work;
void* data;
ma_done_fn done;
void* arg;
};
static void ma_done_trampoline(void* raw) {
struct ma_thread_ctx* ctx = (struct ma_thread_ctx*)raw;
ctx->done(ctx->arg, 0);
u_free(ctx);
}
static void* ma_thread_entry(void* raw) {
struct ma_thread_ctx* ctx = (struct ma_thread_ctx*)raw;
ctx->work(ctx->data);
uasync_post(ctx->ua, ma_done_trampoline, ctx);
return NULL;
}
struct media_async* media_async_create(void) {
struct media_async* ma = u_malloc(sizeof(*ma));
if (!ma) return NULL;
ma->initialized = 1;
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "media_async: created");
return ma;
}
void media_async_destroy(struct media_async* ma) {
if (!ma) return;
ma->initialized = 0;
u_free(ma);
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "media_async: destroyed");
}
void media_async_submit(struct media_async* ma, struct UASYNC* ua,
ma_work_fn work, void* data,
ma_done_fn done, void* arg) {
(void)ma;
if (!ua || !work || !done) return;
struct ma_thread_ctx* ctx = u_malloc(sizeof(*ctx));
if (!ctx) { done(arg, -1); return; }
ctx->ua = ua;
ctx->work = work;
ctx->data = data;
ctx->done = done;
ctx->arg = arg;
pthread_t tid;
int rc = pthread_create(&tid, NULL, ma_thread_entry, ctx);
if (rc != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "media_async: pthread_create failed rc=%d", rc);
u_free(ctx);
done(arg, -1);
return;
}
pthread_detach(tid);
}
/* ─── примитивы (вызываются из worker thread) ─── */
int ma_sha256_file(const char* path, uint8_t hash_out[32]) {
FILE* f = fopen(path, "rb");
if (!f) { DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "ma_sha256_file: cannot open %s", path); return -1; }
SHA256_CTX ctx;
SHA256_Init(&ctx);
uint8_t buf[65536];
size_t rd;
while ((rd = fread(buf, 1, sizeof(buf), f)) > 0) SHA256_Update(&ctx, buf, rd);
fclose(f);
SHA256_Final(hash_out, &ctx);
return 0;
}
int ma_sign_block(const uint8_t* ed25519_privkey, const uint8_t* data, size_t len, uint8_t sig_out[64]) {
EVP_PKEY* pkey = EVP_PKEY_new_raw_private_key(EVP_PKEY_ED25519, NULL, ed25519_privkey, 32);
if (!pkey) { DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "ma_sign_block: EVP_PKEY_new_raw_private_key failed"); return -1; }
EVP_MD_CTX* mdctx = EVP_MD_CTX_new();
if (!mdctx) { EVP_PKEY_free(pkey); return -1; }
if (EVP_DigestSignInit(mdctx, NULL, NULL, NULL, pkey) != 1) {
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "ma_sign_block: EVP_DigestSignInit failed");
EVP_MD_CTX_free(mdctx); EVP_PKEY_free(pkey); return -1;
}
size_t siglen = 64;
int rc = EVP_DigestSign(mdctx, sig_out, &siglen, data, len);
EVP_MD_CTX_free(mdctx);
EVP_PKEY_free(pkey);
if (rc != 1 || siglen != 64) {
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "ma_sign_block: EVP_DigestSign failed rc=%d len=%zu", rc, siglen);
return -1;
}
return 0;
}
int ma_copy_file(const char* src, const char* dst) {
if (!src || !dst) return -1;
if (strcmp(src, dst) == 0) return 0;
FILE* s = fopen(src, "rb");
if (!s) { DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "ma_copy_file: cannot open src %s", src); return -1; }
FILE* d = fopen(dst, "wb");
if (!d) { DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "ma_copy_file: cannot create dst %s", dst); fclose(s); return -1; }
uint8_t buf[65536];
size_t rd;
while ((rd = fread(buf, 1, sizeof(buf), s)) > 0) {
if (fwrite(buf, 1, rd, d) != rd) {
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "ma_copy_file: write error at %s", dst);
fclose(s); fclose(d); return -1;
}
}
fclose(s); fclose(d);
return 0;
}
int ma_file_size(const char* path) {
struct stat st;
if (stat(path, &st) != 0) return -1;
return (int)(st.st_size > 0 ? st.st_size : 0);
}
void ma_uuid(uint8_t uuid_out[16]) {
RAND_bytes(uuid_out, 16);
}

36
src/media_async/media_async.h

@ -0,0 +1,36 @@
// media_async.h — асинхронный движок thread-per-task + примитивы (хеш, подпись, копия файла)
#ifndef MEDIA_ASYNC_H
#define MEDIA_ASYNC_H
#ifdef __cplusplus
extern "C" {
#endif
#include <stdint.h>
#include <stddef.h>
struct UASYNC;
struct media_async;
typedef void (*ma_work_fn)(void* data);
typedef void (*ma_done_fn)(void* arg, int err);
struct media_async* media_async_create(void);
void media_async_destroy(struct media_async* ma);
void media_async_submit(struct media_async* ma, struct UASYNC* ua,
ma_work_fn work, void* data,
ma_done_fn done, void* arg);
int ma_sha256_file(const char* path, uint8_t hash_out[32]);
int ma_sign_block(const uint8_t* ed25519_privkey,
const uint8_t* data, size_t len, uint8_t sig_out[64]);
int ma_copy_file(const char* src, const char* dst);
int ma_file_size(const char* path);
void ma_uuid(uint8_t uuid_out[16]);
#ifdef __cplusplus
}
#endif
#endif

41
src/media_delivery/media_delivery.c

@ -0,0 +1,41 @@
#include "media_delivery.h"
#include "utun_instance.h"
#include "etcp_router.h"
#include "../lib/debug_config.h"
#include <string.h>
static void media_delivery_etcp_recv_cb(struct ETCP_CONN* conn, struct ll_entry* entry) {
(void)conn; (void)entry;
}
int media_delivery_init(struct UTUN_INSTANCE* inst) {
if (!inst) return -1;
struct media_delivery_ctx* md = &inst->md;
memset(md, 0, sizeof(*md));
md->self_node_id = inst->node_id;
if (etcp_router_bind(inst, ETCP_RT_ID_MEDIA_DELIVERY, media_delivery_etcp_recv_cb) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "media_delivery: failed to bind ETCP_RT_ID_MEDIA_DELIVERY");
return -1;
}
md->initialized = 1;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "media_delivery initialized, node_id=%016llx",
(unsigned long long)md->self_node_id);
return 0;
}
void media_delivery_destroy(struct UTUN_INSTANCE* inst) {
if (!inst || !inst->md.initialized) return;
struct media_delivery_ctx* md = &inst->md;
etcp_router_unbind(inst, ETCP_RT_ID_MEDIA_DELIVERY);
memset(md, 0, sizeof(*md));
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "media_delivery destroyed");
}
int media_delivery_bind(struct UTUN_INSTANCE* inst) {
return media_delivery_init(inst);
}

27
src/media_delivery/media_delivery.h

@ -0,0 +1,27 @@
// media_delivery.h — распространение медиа (аудио/видео стриминг)
#ifndef MEDIA_DELIVERY_H
#define MEDIA_DELIVERY_H
#ifdef __cplusplus
extern "C" {
#endif
#include <stdint.h>
struct UTUN_INSTANCE;
struct media_delivery_ctx {
uint64_t self_node_id;
int initialized;
};
int media_delivery_init(struct UTUN_INSTANCE* inst);
void media_delivery_destroy(struct UTUN_INSTANCE* inst);
int media_delivery_bind(struct UTUN_INSTANCE* inst);
#ifdef __cplusplus
}
#endif
#endif

297
src/media_delivery/media_index.c

@ -0,0 +1,297 @@
// media_index.c — таблица media_files: init, commit, register_async
#include "media_index.h"
#include "../media_async/media_async.h"
#include "../../lib/u_async.h"
#include "../../lib/mem.h"
#include "../../lib/debug_config.h"
#include <openssl/evp.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <time.h>
#include <unistd.h>
#define MEDIA_BLOCK_MIN (10 * 1024 * 1024)
#define MEDIA_BLOCK_MAX (25 * 1024 * 1024)
#define MEDIA_BLOCK_TARGET 30
static int mi_insert(sqlite3* db,
const uint8_t* media_id, const uint8_t* block_id,
const uint8_t* content_hash,
const char* chat_id, const char* location,
uint64_t node_id, const uint8_t* node_sign,
int64_t timestamp, int64_t file_size,
int64_t chunk_size, int chunk, int64_t offset) {
const char* sql =
"INSERT INTO media_files(media_id,block_id,content_hash,chat_id,location,"
"node_id,node_sign,timestamp,file_size,chunk_size,chunk,offset)"
" VALUES(?,?,?,?,?,?,?,?,?,?,?,?)";
sqlite3_stmt* stmt = NULL;
if (sqlite3_prepare_v2(db, sql, -1, &stmt, NULL) != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "media_index: insert prep fail: %s", sqlite3_errmsg(db));
return -1;
}
sqlite3_bind_blob(stmt, 1, media_id, 16, SQLITE_STATIC);
sqlite3_bind_blob(stmt, 2, block_id, 16, SQLITE_STATIC);
sqlite3_bind_blob(stmt, 3, content_hash, 32, SQLITE_STATIC);
sqlite3_bind_text(stmt, 4, chat_id, -1, SQLITE_STATIC);
sqlite3_bind_text(stmt, 5, location, -1, SQLITE_STATIC);
sqlite3_bind_int64(stmt, 6, (sqlite3_int64)node_id);
sqlite3_bind_blob(stmt, 7, node_sign, 64, SQLITE_STATIC);
sqlite3_bind_int64(stmt, 8, timestamp);
sqlite3_bind_int64(stmt, 9, file_size);
sqlite3_bind_int64(stmt, 10, chunk_size);
sqlite3_bind_int(stmt, 11, chunk);
sqlite3_bind_int64(stmt, 12, offset);
int rc = sqlite3_step(stmt);
sqlite3_finalize(stmt);
if (rc != SQLITE_DONE) {
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "media_index: insert fail rc=%d: %s", rc, sqlite3_errmsg(db));
return -1;
}
return 0;
}
int media_index_init(sqlite3* db) {
if (!db) return -1;
const char* sql =
"CREATE TABLE IF NOT EXISTS media_files ("
" id INTEGER PRIMARY KEY AUTOINCREMENT,"
" media_id BLOB NOT NULL,"
" block_id BLOB NOT NULL,"
" content_hash BLOB NOT NULL,"
" chat_id TEXT NOT NULL,"
" location TEXT NOT NULL,"
" node_id INTEGER NOT NULL,"
" node_sign BLOB NOT NULL,"
" timestamp INTEGER NOT NULL,"
" file_size INTEGER NOT NULL,"
" chunk_size INTEGER NOT NULL,"
" chunk INTEGER NOT NULL,"
" offset INTEGER NOT NULL);";
char* err = NULL;
int rc = sqlite3_exec(db, sql, NULL, NULL, &err);
if (rc != SQLITE_OK) {
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "media_index: CREATE TABLE failed: %s", err ? err : "?");
sqlite3_free(err);
return -1;
}
sqlite3_exec(db,
"CREATE UNIQUE INDEX IF NOT EXISTS idx_media_files_block_id ON media_files(block_id);",
NULL, NULL, NULL);
sqlite3_exec(db,
"CREATE INDEX IF NOT EXISTS idx_media_files_media_id ON media_files(media_id);",
NULL, NULL, NULL);
sqlite3_exec(db,
"CREATE INDEX IF NOT EXISTS idx_media_files_chat_id ON media_files(chat_id);",
NULL, NULL, NULL);
sqlite3_exec(db,
"CREATE INDEX IF NOT EXISTS idx_media_files_node_id ON media_files(node_id);",
NULL, NULL, NULL);
sqlite3_exec(db,
"CREATE INDEX IF NOT EXISTS idx_media_files_content_hash ON media_files(content_hash);",
NULL, NULL, NULL);
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "media_index: table+indices created");
return 0;
}
int media_calc_block_size(uint64_t file_size) {
if (file_size == 0) return 0;
uint64_t target = (file_size + MEDIA_BLOCK_TARGET - 1) / MEDIA_BLOCK_TARGET;
if (target < MEDIA_BLOCK_MIN) target = MEDIA_BLOCK_MIN;
if (target > MEDIA_BLOCK_MAX) target = MEDIA_BLOCK_MAX;
return (int)target;
}
int media_index_generate_uuid(uint8_t uuid_out[16]) {
ma_uuid(uuid_out);
return 0;
}
void media_index_result_free(struct media_index_result* result) {
if (!result) return;
if (result->block_ids) { u_free(result->block_ids); result->block_ids = NULL; }
if (result->block_sigs) { u_free(result->block_sigs); result->block_sigs = NULL; }
}
int media_index_commit(sqlite3* db, const struct media_index_result* result,
uint64_t node_id, const uint8_t* ed25519_privkey,
const char* chat_id, const char* dest_path, const char* media_base) {
if (!db || !result || !chat_id || !dest_path || !media_base) return -1;
int nb = result->num_blocks;
int64_t bs = result->block_size;
int64_t fs = result->file_size;
int64_t ts = (int64_t)time(NULL);
size_t base_len = strlen(media_base);
if (base_len > 0 && media_base[base_len - 1] == '/') base_len--;
const char* rel = dest_path;
if (strncmp(rel, media_base, base_len) == 0 && rel[base_len] == '/') rel += base_len + 1;
for (int n = 0; n < nb; n++) {
int64_t cs = (n == nb - 1) ? fs - n * bs : bs;
int64_t off = n * bs;
char location[512];
snprintf(location, sizeof(location), "%s", rel);
uint8_t sig_msg[92];
size_t soff = 0;
memcpy(sig_msg + soff, result->media_id, 16); soff += 16;
memcpy(sig_msg + soff, result->content_hash, 32); soff += 32;
memcpy(sig_msg + soff, &ts, 8); soff += 8;
memcpy(sig_msg + soff, &fs, 8); soff += 8;
memcpy(sig_msg + soff, &cs, 8); soff += 8;
int32_t cn = (int32_t)n; memcpy(sig_msg + soff, &cn, 4); soff += 4;
memcpy(sig_msg + soff, &off, 8); soff += 8;
memcpy(sig_msg + soff, &node_id, 8); soff += 8;
uint8_t node_sign[64];
if (ma_sign_block(ed25519_privkey, sig_msg, soff, node_sign) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "media_index: node_sign failed chunk=%d", n);
return -1;
}
if (mi_insert(db, result->media_id, result->block_ids + n * 16,
result->content_hash, chat_id, location,
node_id, node_sign, ts, fs, cs, n, off) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "media_index: insert failed chunk=%d", n);
return -1;
}
}
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "media_index: committed %d chunks for chat=%s path=%s",
nb, chat_id, rel);
return 0;
}
/* ─── register_async (worker + uasync callback) ─── */
struct mi_reg_ctx {
struct UASYNC* ua;
uint8_t ed25519_privkey[32];
char src_path[1024];
char dest_path[1024];
int copy_file;
uint64_t node_id;
char chat_id[64];
char media_base[512];
sqlite3* db;
void (*cb)(void* arg, int err, const struct media_index_result* result);
void* cb_arg;
struct media_index_result result;
int err;
};
static void mi_reg_work(void* data) {
struct mi_reg_ctx* ctx = (struct mi_reg_ctx*)data;
if (ctx->copy_file && ma_copy_file(ctx->src_path, ctx->dest_path) != 0) {
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "media_index: copy failed %s -> %s", ctx->src_path, ctx->dest_path);
ctx->err = -10; return;
}
int64_t fsize = (int64_t)ma_file_size(ctx->dest_path);
if (fsize <= 0) { ctx->err = -20; return; }
ctx->result.file_size = fsize;
int64_t bs = (int64_t)media_calc_block_size((uint64_t)fsize);
int nb = (int)((fsize + bs - 1) / bs);
ctx->result.block_size = bs;
ctx->result.num_blocks = nb;
if (ma_sha256_file(ctx->dest_path, ctx->result.content_hash) != 0) {
ctx->err = -30; return;
}
ma_uuid(ctx->result.media_id);
ctx->result.block_ids = u_malloc((size_t)nb * 16);
ctx->result.block_sigs = u_malloc((size_t)nb * 64);
if (!ctx->result.block_ids || !ctx->result.block_sigs) { ctx->err = -40; return; }
FILE* f = fopen(ctx->dest_path, "rb");
if (!f) { ctx->err = -50; return; }
for (int n = 0; n < nb; n++) {
int64_t off = n * bs;
int64_t sz = (n == nb - 1) ? fsize - off : bs;
uint8_t* buf = u_malloc((size_t)sz);
if (!buf) { fclose(f); ctx->err = -60; return; }
fseeko(f, (off_t)off, SEEK_SET);
size_t rd = fread(buf, 1, (size_t)sz, f);
if (rd != (size_t)sz) { u_free(buf); fclose(f); ctx->err = -70; return; }
ma_uuid(ctx->result.block_ids + n * 16);
if (ma_sign_block(ctx->ed25519_privkey, buf, (size_t)sz,
ctx->result.block_sigs + n * 64) != 0) {
u_free(buf); fclose(f); ctx->err = -80; return;
}
u_free(buf);
}
fclose(f);
ctx->err = 0;
}
static void mi_reg_done(void* arg, int err) {
(void)err;
struct mi_reg_ctx* ctx = (struct mi_reg_ctx*)arg;
if (ctx->err == 0) {
int rc = media_index_commit(ctx->db, &ctx->result, ctx->node_id,
ctx->ed25519_privkey, ctx->chat_id,
ctx->dest_path, ctx->media_base);
if (rc != 0) ctx->err = -90;
}
if (ctx->cb) ctx->cb(ctx->cb_arg, ctx->err, ctx->err ? NULL : &ctx->result);
media_index_result_free(&ctx->result);
u_free(ctx);
}
void media_index_register_async(
struct media_async* ma, struct UASYNC* ua, sqlite3* db,
uint64_t node_id, const uint8_t* ed25519_privkey,
const char* chat_id,
const char* src_path, const char* dest_path, int copy_file,
const char* media_base,
void (*cb)(void* arg, int err, const struct media_index_result* result),
void* cb_arg) {
if (!ma || !ua || !db || !ed25519_privkey || !chat_id || !src_path || !dest_path || !media_base || !cb) {
DEBUG_ERROR(DEBUG_CATEGORY_DEBUG, "media_index: register_async invalid args");
if (cb) cb(cb_arg, -1, NULL);
return;
}
struct mi_reg_ctx* ctx = u_calloc(1, sizeof(*ctx));
if (!ctx) { cb(cb_arg, -1, NULL); return; }
ctx->ua = ua;
ctx->db = db;
ctx->node_id = node_id;
ctx->copy_file = copy_file;
ctx->cb = cb;
ctx->cb_arg = cb_arg;
memcpy(ctx->ed25519_privkey, ed25519_privkey, 32);
snprintf(ctx->src_path, sizeof(ctx->src_path), "%s", src_path);
snprintf(ctx->dest_path, sizeof(ctx->dest_path), "%s", dest_path);
snprintf(ctx->chat_id, sizeof(ctx->chat_id), "%s", chat_id);
snprintf(ctx->media_base, sizeof(ctx->media_base), "%s", media_base);
media_async_submit(ma, ua, mi_reg_work, ctx, mi_reg_done, ctx);
}

52
src/media_delivery/media_index.h

@ -0,0 +1,52 @@
// media_index.h — таблица media_files + регистрация медиа с async-обработкой
#ifndef MEDIA_INDEX_H
#define MEDIA_INDEX_H
#ifdef __cplusplus
extern "C" {
#endif
#include <stdint.h>
#include <stdint.h>
#include <sqlite3.h>
struct media_async;
struct UASYNC;
#define MEDIA_ID_SIZE 16
#define MEDIA_HASH_SIZE 32
struct media_index_result {
uint8_t media_id[16];
uint8_t content_hash[32];
int64_t file_size;
int64_t block_size;
int num_blocks;
uint8_t* block_ids; // num_blocks * 16
uint8_t* block_sigs; // num_blocks * 64
};
int media_index_init(sqlite3* db);
int media_index_commit(
sqlite3* db, const struct media_index_result* result,
uint64_t node_id, const uint8_t* ed25519_privkey,
const char* chat_id, const char* dest_path, const char* media_base);
void media_index_register_async(
struct media_async* ma, struct UASYNC* ua, sqlite3* db,
uint64_t node_id, const uint8_t* ed25519_privkey,
const char* chat_id,
const char* src_path, const char* dest_path, int copy_file,
const char* media_base,
void (*cb)(void* arg, int err, const struct media_index_result* result),
void* cb_arg);
int media_index_generate_uuid(uint8_t uuid_out[16]);
void media_index_result_free(struct media_index_result* result);
int media_calc_block_size(uint64_t file_size);
#ifdef __cplusplus
}
#endif
#endif

1
src/transport_layer/etcp_api.h

@ -41,6 +41,7 @@ extern "C" {
#define ETCP_RT_ID_TCP_PROXY 0x04 // TCP proxy (клиент ↔ exit)
#define ETCP_RT_ID_UDP_PROXY 0x05 // UDP datagram прокси (client ↔ exit)
#define ETCP_RT_ID_ICMP_PROXY 0x06 // ICMP echo прокси (ping через exit)
#define ETCP_RT_ID_MEDIA_DELIVERY 0x07 // распространение медиа (аудио/видео стриминг)
#define ETCP_RT_ID_CONN_MGR 0x11 // Connection Manager — management connections

20
src/utun_instance.c

@ -509,6 +509,15 @@ void utun_instance_destroy(struct UTUN_INSTANCE *instance) {
// Cleanup routing module (unbinds from etcp_router before etcp_router_destroy)
routing_destroy(instance);
// Cleanup media delivery (unbinds from etcp_router before etcp_router_destroy)
if (instance->md.initialized) {
media_delivery_destroy(instance);
}
// Cleanup media async engine
media_async_destroy(instance->media_async);
instance->media_async = NULL;
// Cleanup NAT (unbinds from etcp_router before etcp_router_destroy)
if (instance->nat_tr.initialized) {
nat_transport_destroy(instance);
@ -637,6 +646,17 @@ int utun_instance_init(struct UTUN_INSTANCE *instance) {
}
}
// Initialize media delivery
if (media_delivery_init(instance) != 0) {
DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "media_delivery init failed (non-fatal)");
}
// Initialize media async engine
instance->media_async = media_async_create();
if (!instance->media_async) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "media_async create failed");
}
// Note: TUN socket is already registered in tun_init()
// Initialize connections

8
src/utun_instance.h

@ -18,6 +18,8 @@ extern "C" {
#include "firewall.h"
#include "eim_nat.h"
#include "nat_transport.h"
#include "media_delivery/media_delivery.h"
#include "media_async/media_async.h"
#include "proxy/tcp_proxy_client.h"
#include "etcp_router.h"
#include "proxy/tcp_proxy_server.h"
@ -145,6 +147,12 @@ struct UTUN_INSTANCE {
struct eim_nat_ctx nat;
struct nat_transport_ctx nat_tr;
// Media delivery
struct media_delivery_ctx md;
// Media async engine (thread-per-task for crypto/file ops)
struct media_async* media_async;
// Socket initialization status: 0=OK, 1=partial (some sockets failed), -1=error (none created)
int socket_init_status;

10
tests/Makefile.am

@ -60,6 +60,8 @@ check_PROGRAMS = \
test_uasync_socket_race \
test_ntp \
test_opus_codec \
test_media_async \
test_media_index \
bench_timeout_heap \
bench_uasync_timeouts
@ -315,6 +317,14 @@ test_opus_codec_LDADD = $(COMMON_LIBS) $(top_builddir)/lib/libopus/libopus_inter
test_audio_compressor_SOURCES = test_audio_compressor.c
test_audio_compressor_LDADD = $(COMMON_LIBS) -lm
test_media_async_SOURCES = test_media_async.c
test_media_async_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/src/media_async -I$(top_srcdir)/lib
test_media_async_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS)
test_media_index_SOURCES = test_media_index.c
test_media_index_CFLAGS = -I$(top_srcdir)/src -I$(top_srcdir)/src/media_delivery -I$(top_srcdir)/src/media_async -I$(top_srcdir)/lib
test_media_index_LDADD = $(LIBUTUN) $(CRYPTO_LIBS) $(COMMON_LIBS)
# Copy test configs to build directory (tests run from build/tests/)
all-local: copy-test-configs

8
tests/test_audio_compressor.c

@ -67,7 +67,7 @@ static void test_basic_push_flush(void) {
audio_compressor_config_t cfg = {0};
cfg.sample_rate = 48000; cfg.channels = 1;
cfg.block_duration_ms = 20; cfg.lookback_ms = 200; cfg.lookahead_ms = 100;
cfg.max_gain_db = 30.0f; cfg.rise_rate_per_500ms = 2.0f; cfg.target_level = 0.25f;
cfg.max_gain_db = 30.0f; cfg.rise_rate_per_sec = 4.0f; cfg.target_level = 1.0f;
audio_compressor_configure(ac, &cfg);
int16_t in[960];
@ -97,7 +97,7 @@ static void test_silence_in_silence_out(void) {
audio_compressor_config_t cfg = {0};
cfg.sample_rate = 48000; cfg.channels = 1;
cfg.block_duration_ms = 20; cfg.lookback_ms = 200; cfg.lookahead_ms = 100;
cfg.max_gain_db = 30.0f; cfg.rise_rate_per_500ms = 2.0f; cfg.target_level = 0.25f;
cfg.max_gain_db = 30.0f; cfg.rise_rate_per_sec = 4.0f; cfg.target_level = 1.0f;
audio_compressor_configure(ac, &cfg);
int16_t in[960]; fill_silence(in, 960);
@ -121,7 +121,7 @@ static void test_reset(void) {
audio_compressor_config_t cfg = {0};
cfg.sample_rate = 48000; cfg.channels = 1;
cfg.block_duration_ms = 20; cfg.lookback_ms = 200; cfg.lookahead_ms = 100;
cfg.max_gain_db = 30.0f; cfg.rise_rate_per_500ms = 2.0f; cfg.target_level = 0.25f;
cfg.max_gain_db = 30.0f; cfg.rise_rate_per_sec = 4.0f; cfg.target_level = 1.0f;
audio_compressor_configure(ac, &cfg);
int16_t in[960];
@ -156,7 +156,7 @@ static void test_stress(void) {
cfg.block_duration_ms = 20; cfg.lookback_ms = (rand() % 200) + 50;
cfg.lookahead_ms = (rand() % 100) + 20;
cfg.max_gain_db = (float)(rand() % 31);
cfg.rise_rate_per_500ms = 1.1f + (float)(rand() % 90) / 10.0f;
cfg.rise_rate_per_sec = 0.4f + (float)(rand() % 27) / 10.0f;
cfg.target_level = 0.1f + (float)(rand() % 30) / 100.0f;
audio_compressor_configure(ac, &cfg);

299
tests/test_media_async.c

@ -0,0 +1,299 @@
// test_media_async.c — стресс-тест примитивов media_async (многопоточность + корректность)
//
// 4 потока одновременно выполняют ma_* операции на файлах 1K/5K/10K.
// Тест длится 2 секунды, затем потоки останавливаются.
// Проверяется корректность каждого результата и отсутствие заклиниваний.
#include "../src/media_async/media_async.h"
#include <pthread.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <time.h>
#include <unistd.h>
#include <sys/stat.h>
#include <sys/time.h>
#include <openssl/evp.h>
#define NUM_THREADS 4
#define TEST_TIMEOUT_SEC 2
#define NUM_FILES 3
#define SIZE_1K 1024
#define SIZE_5K 5120
#define SIZE_10K 10240
#define NUM_OPS 5
static volatile int g_running = 1;
static const char* g_op_names[NUM_OPS] = {
"SHA256-file", "Sign-block", "Copy-file",
"UUID", "FileSize"
};
struct thread_stats {
int tid;
int op_count;
int errors;
int op_counts[NUM_OPS];
double max_time_us[NUM_OPS]; // максимальное время (микросекунды) по каждой операции
};
static char g_test_dir[256];
static char g_file_paths[NUM_FILES][512];
static uint8_t g_file_hashes[NUM_FILES][32];
static int g_file_sizes[NUM_FILES];
static uint8_t g_test_data[NUM_FILES][SIZE_10K];
static uint8_t g_ed25519_privkey[32];
static uint8_t g_expected_sigs[NUM_FILES][64];
static int64_t now_us(void) {
struct timeval tv; gettimeofday(&tv, NULL);
return (int64_t)tv.tv_sec * 1000000LL + tv.tv_usec;
}
static void fill_test_data(uint8_t* buf, int size, int seed) {
unsigned s = (unsigned)seed;
for (int i = 0; i < size; i++) {
s = s * 1103515245 + 12345;
buf[i] = (uint8_t)((s >> 16) & 0xFF);
}
}
static void* worker_thread(void* arg) {
struct thread_stats* st = (struct thread_stats*)arg;
unsigned s = (unsigned)(time(NULL) ^ (st->tid * 0x9E3779B9 << 8));
while (g_running) {
s = s * 1103515245 + 12345;
int op = (int)((s >> 16) % 4); // random: 0=SHA256, 1=Sign, 2=Copy, 3=UUID+FileSize
s = s * 1103515245 + 12345;
int fi = (int)((s >> 16) % NUM_FILES);
int size = g_file_sizes[fi];
int real_op = -1;
int64_t t0, elapsed;
switch (op) {
case 0: { // SHA256 file
real_op = 0;
uint8_t hash[32];
t0 = now_us();
int rc = ma_sha256_file(g_file_paths[fi], hash);
elapsed = now_us() - t0;
if (rc != 0 || memcmp(hash, g_file_hashes[fi], 32) != 0) st->errors++;
break;
}
case 1: { // Ed25519 sign
real_op = 1;
uint8_t sig[64];
t0 = now_us();
int rc = ma_sign_block(g_ed25519_privkey, g_test_data[fi], (size_t)size, sig);
elapsed = now_us() - t0;
if (rc != 0 || memcmp(sig, g_expected_sigs[fi], 64) != 0) st->errors++;
break;
}
case 2: { // copy file
real_op = 2;
char dst[512];
snprintf(dst, sizeof(dst), "%s/cp_t%d_n%d.bin", g_test_dir, st->tid, st->op_count);
t0 = now_us();
int rc = ma_copy_file(g_file_paths[fi], dst);
elapsed = now_us() - t0;
if (rc == 0) {
uint8_t hash[32];
if (ma_sha256_file(dst, hash) != 0 || memcmp(hash, g_file_hashes[fi], 32) != 0)
st->errors++;
unlink(dst);
} else st->errors++;
break;
}
case 3: { // UUID
real_op = 3;
uint8_t uuid[16];
t0 = now_us();
ma_uuid(uuid);
elapsed = now_us() - t0;
int zero = 1;
for (int i = 0; i < 16; i++) { if (uuid[i] != 0) { zero = 0; break; } }
if (zero) st->errors++;
// file_size отдельно с замером времени
t0 = now_us();
int64_t fsz = ma_file_size(g_file_paths[fi]);
int64_t et = now_us() - t0;
if (fsz != g_file_sizes[fi]) st->errors++;
st->op_counts[4]++;
if (et > st->max_time_us[4]) st->max_time_us[4] = (double)et;
break;
}
}
st->op_count++;
if (real_op >= 0) {
st->op_counts[real_op]++;
if (elapsed > st->max_time_us[real_op])
st->max_time_us[real_op] = (double)elapsed;
}
}
return NULL;
}
int main(void) {
int ret = 0;
snprintf(g_test_dir, sizeof(g_test_dir), "/tmp/test_media_async_%d", (int)getpid());
if (mkdir(g_test_dir, 0755) != 0) {
perror("mkdir"); return 1;
}
g_file_sizes[0] = SIZE_1K;
g_file_sizes[1] = SIZE_5K;
g_file_sizes[2] = SIZE_10K;
for (int i = 0; i < NUM_FILES; i++) {
fill_test_data(g_test_data[i], g_file_sizes[i], 42 + i * 100);
snprintf(g_file_paths[i], sizeof(g_file_paths[i]), "%s/test_%dk.bin", g_test_dir, g_file_sizes[i] / 1024);
FILE* f = fopen(g_file_paths[i], "wb");
if (!f) { perror("fopen"); ret = 1; goto cleanup; }
fwrite(g_test_data[i], 1, (size_t)g_file_sizes[i], f);
fclose(f);
ma_sha256_file(g_file_paths[i], g_file_hashes[i]);
}
{
EVP_PKEY_CTX* pctx = EVP_PKEY_CTX_new_id(EVP_PKEY_ED25519, NULL);
if (!pctx) { fprintf(stderr, "EVP_PKEY_CTX_new_id failed\n"); ret = 1; goto cleanup; }
EVP_PKEY_keygen_init(pctx);
EVP_PKEY* pkey = NULL;
if (EVP_PKEY_keygen(pctx, &pkey) != 1 || !pkey) {
fprintf(stderr, "EVP_PKEY_keygen failed\n"); EVP_PKEY_CTX_free(pctx); ret = 1; goto cleanup;
}
EVP_PKEY_CTX_free(pctx);
size_t keylen = 32;
EVP_PKEY_get_raw_private_key(pkey, g_ed25519_privkey, &keylen);
EVP_PKEY_free(pkey);
for (int i = 0; i < NUM_FILES; i++)
ma_sign_block(g_ed25519_privkey, g_test_data[i], (size_t)g_file_sizes[i], g_expected_sigs[i]);
}
pthread_t threads[NUM_THREADS];
struct thread_stats stats[NUM_THREADS];
memset(stats, 0, sizeof(stats));
for (int i = 0; i < NUM_THREADS; i++) {
stats[i].tid = i;
if (pthread_create(&threads[i], NULL, worker_thread, &stats[i]) != 0) {
fprintf(stderr, "pthread_create(%d) failed\n", i);
g_running = 0;
for (int j = 0; j < i; j++) pthread_join(threads[j], NULL);
ret = 1; goto cleanup;
}
}
sleep(TEST_TIMEOUT_SEC);
g_running = 0;
for (int i = 0; i < NUM_THREADS; i++)
pthread_join(threads[i], NULL);
printf("\n=== media_async stress test ===\n");
printf("Timeout: %d sec, Threads: %d, Files: 1K/5K/10K\n\n", TEST_TIMEOUT_SEC, NUM_THREADS);
int total_ops = 0, total_errors = 0;
int totals_op[NUM_OPS] = {0};
double max_time_all[NUM_OPS] = {0};
for (int i = 0; i < NUM_THREADS; i++) {
printf("Thread %d: %d ops (", i, stats[i].op_count);
for (int o = 0; o < NUM_OPS; o++) {
if (o > 0) printf(" ");
printf("%s:%d", g_op_names[o], stats[i].op_counts[o]);
totals_op[o] += stats[i].op_counts[o];
}
printf(") errors=%d\n", stats[i].errors);
printf(" max-time(us): ");
for (int o = 0; o < NUM_OPS; o++) {
if (o > 0) printf(" ");
printf("%s:%.0f", g_op_names[o], stats[i].max_time_us[o]);
if (stats[i].max_time_us[o] > max_time_all[o]) max_time_all[o] = stats[i].max_time_us[o];
}
printf("\n");
total_ops += stats[i].op_count;
total_errors += stats[i].errors;
}
printf("\nTotal: %d ops, %d errors\n", total_ops, total_errors);
printf("Op totals: ");
for (int o = 0; o < NUM_OPS; o++) printf("%s:%d ", g_op_names[o], totals_op[o]);
printf("\nMax time(us): ");
for (int o = 0; o < NUM_OPS; o++) printf("%s:%.0f ", g_op_names[o], max_time_all[o]);
printf("\n\n");
/* ─── sanity checks ─── */
// минимальный throughput: хотя бы 1000 ops на поток за 2 сек
int min_ops_per_thread = 1000;
// разбалансировка потоков: самый медленный поток не хуже 0.1× от самого быстрого
double balance_ratio = 0.1;
// макс допустимое время SHA256/Sign/Copy/USB операции (микросекунды)
double max_us_sha = 100000; // 100ms
double max_us_sign = 100000;
double max_us_copy = 100000;
double max_us_uuid = 5000; // 5ms
double max_us_filesize = 5000;
int min_ops = stats[0].op_count, max_ops = stats[0].op_count;
for (int i = 1; i < NUM_THREADS; i++) {
if (stats[i].op_count < min_ops) min_ops = stats[i].op_count;
if (stats[i].op_count > max_ops) max_ops = stats[i].op_count;
}
printf("Sanity checks:\n");
if (total_errors > 0) {
printf(" FAIL: %d errors\n", total_errors);
ret = 1;
}
if (total_ops < min_ops_per_thread * NUM_THREADS) {
printf(" FAIL: total ops %d < %d (throughput too low)\n",
total_ops, min_ops_per_thread * NUM_THREADS);
ret = 1;
} else {
printf(" OK: total ops=%d >= %d\n", total_ops, min_ops_per_thread * NUM_THREADS);
}
if (max_ops > 0 && min_ops < (int)(max_ops * balance_ratio)) {
printf(" FAIL: thread imbalance min=%d max=%d (ratio %.3f < %.1f)\n",
min_ops, max_ops, max_ops > 0 ? (double)min_ops / max_ops : 0, balance_ratio);
ret = 1;
} else {
printf(" OK: thread balance min=%d max=%d (ratio %.3f)\n",
min_ops, max_ops, max_ops > 0 ? (double)min_ops / max_ops : 0);
}
double limits[NUM_OPS] = { max_us_sha, max_us_sign, max_us_copy, max_us_uuid, max_us_filesize };
for (int o = 0; o < NUM_OPS; o++) {
if (max_time_all[o] > limits[o]) {
printf(" FAIL: %s max-time %.0fus > %.0fus\n", g_op_names[o], max_time_all[o], limits[o]);
ret = 1;
} else if (max_time_all[o] <= 0.0) {
printf(" WARN: %s max-time 0 (no measurements?)\n", g_op_names[o]);
} else {
printf(" OK: %s max-time %.0fus\n", g_op_names[o], max_time_all[o]);
}
}
if (ret == 0) {
printf("\nPASS: %d ops, 0 errors, all checks OK\n", total_ops);
} else {
printf("\nFAIL: sanity checks failed\n");
}
cleanup:
for (int i = 0; i < NUM_FILES; i++) unlink(g_file_paths[i]);
rmdir(g_test_dir);
return ret;
}

510
tests/test_media_index.c

@ -0,0 +1,510 @@
// test_media_index.c — юнит-тесты media_index (commit + register_async)
//
// Покрытие:
// 1. 1 чанк, точный размер
// 2. N чанков, точное совпадение
// 3. Последний чанк меньше
// 4. Файл меньше размера чанка
// 5. Много чанков (>3)
// 6. media_base с trailing slash
// 7. media_base без trailing slash
// 8. node_sign верификация (Ed25519 verify)
// 9. Разные chat_id
// 10. Полный async flow (register_async + copy=1)
#include "../src/media_delivery/media_index.h"
#include "../src/media_async/media_async.h"
#include "../lib/u_async.h"
#include "../lib/mem.h"
#include "../lib/debug_config.h"
#include <openssl/evp.h>
#include <sqlite3.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
#include <sys/stat.h>
#define BS 100
#define NODE_ID 0xDEADBEEF
static int g_passed = 0;
static int g_failed = 0;
static int g_total = 0;
#define TEST(name) do { \
g_total++; \
printf(" %-55s", name); \
} while(0)
#define OK() do { g_passed++; printf("OK\n"); } while(0)
#define FAIL(fmt, ...) do { \
g_failed++; \
printf("FAIL: " fmt "\n", ##__VA_ARGS__); \
} while(0)
/* ─── helpers ─── */
static sqlite3* make_db(void) {
sqlite3* db = NULL;
sqlite3_open(":memory:", &db);
sqlite3_exec(db, "PRAGMA journal_mode=memory", NULL, NULL, NULL);
return db;
}
static void free_db(sqlite3* db) { if (db) sqlite3_close(db); }
static void derive_pubkey(const uint8_t priv[32], uint8_t pub[32]) {
EVP_PKEY* pkey = EVP_PKEY_new_raw_private_key(EVP_PKEY_ED25519, NULL, priv, 32);
if (pkey) {
size_t l = 32;
EVP_PKEY_get_raw_public_key(pkey, pub, &l);
EVP_PKEY_free(pkey);
}
}
static int ed25519_verify(const uint8_t pub[32], const uint8_t* msg, size_t len, const uint8_t sig[64]) {
EVP_PKEY* pkey = EVP_PKEY_new_raw_public_key(EVP_PKEY_ED25519, NULL, pub, 32);
if (!pkey) return -1;
EVP_MD_CTX* mdctx = EVP_MD_CTX_new();
if (!mdctx) { EVP_PKEY_free(pkey); return -1; }
if (EVP_DigestVerifyInit(mdctx, NULL, NULL, NULL, pkey) != 1) {
EVP_MD_CTX_free(mdctx); EVP_PKEY_free(pkey); return -1;
}
int rc = EVP_DigestVerify(mdctx, sig, 64, msg, len);
EVP_MD_CTX_free(mdctx); EVP_PKEY_free(pkey);
return rc == 1 ? 0 : -1;
}
static int row_count(sqlite3* db, const uint8_t* media_id) {
sqlite3_stmt* st = NULL;
sqlite3_prepare_v2(db, "SELECT COUNT(*) FROM media_files WHERE media_id=?", -1, &st, NULL);
sqlite3_bind_blob(st, 1, media_id, 16, SQLITE_STATIC);
int c = 0;
if (sqlite3_step(st) == SQLITE_ROW) c = sqlite3_column_int(st, 0);
sqlite3_finalize(st);
return c;
}
/* заполняет result рандомными данными для тестирования commit */
static void fill_result(struct media_index_result* r, int64_t file_size, int64_t block_size, int num_blocks) {
memset(r, 0, sizeof(*r));
r->file_size = file_size;
r->block_size = block_size;
r->num_blocks = num_blocks;
ma_uuid(r->media_id);
ma_uuid(r->content_hash);
ma_uuid(r->content_hash + 16);
r->block_ids = u_malloc((size_t)num_blocks * 16);
r->block_sigs = u_malloc((size_t)num_blocks * 64);
for (int i = 0; i < num_blocks; i++) {
ma_uuid(r->block_ids + i * 16);
ma_uuid(r->block_sigs + i * 64);
ma_uuid(r->block_sigs + i * 64 + 16);
ma_uuid(r->block_sigs + i * 64 + 32);
ma_uuid(r->block_sigs + i * 64 + 48);
}
}
static int check_row(sqlite3* db, const uint8_t* media_id, int chunk,
const uint8_t* content_hash, int64_t file_size,
int64_t chunk_size, int64_t offset, uint64_t node_id,
const char* chat_id, const char* location,
const uint8_t* node_pubkey) {
sqlite3_stmt* st = NULL;
sqlite3_prepare_v2(db,
"SELECT block_id,content_hash,chat_id,location,node_id,node_sign,timestamp,file_size,chunk_size,chunk,offset"
" FROM media_files WHERE media_id=? AND chunk=?",
-1, &st, NULL);
if (!st) return -1;
sqlite3_bind_blob(st, 1, media_id, 16, SQLITE_STATIC);
sqlite3_bind_int(st, 2, chunk);
if (sqlite3_step(st) != SQLITE_ROW) { sqlite3_finalize(st); return -1; }
const uint8_t* bid = (const uint8_t*)sqlite3_column_blob(st, 0);
int bid_len = sqlite3_column_bytes(st, 0);
const uint8_t* ch = (const uint8_t*)sqlite3_column_blob(st, 1);
int ch_len = sqlite3_column_bytes(st, 1);
const char* cid = (const char*)sqlite3_column_text(st, 2);
const char* loc = (const char*)sqlite3_column_text(st, 3);
uint64_t nid = (uint64_t)sqlite3_column_int64(st, 4);
const uint8_t* nsig = (const uint8_t*)sqlite3_column_blob(st, 5);
int nsig_len = sqlite3_column_bytes(st, 5);
int64_t ts = sqlite3_column_int64(st, 6);
int64_t fsz = sqlite3_column_int64(st, 7);
int64_t csz = sqlite3_column_int64(st, 8);
int c = sqlite3_column_int(st, 9);
int64_t off = sqlite3_column_int64(st, 10);
int err = 0;
if (!bid || bid_len != 16) { printf("bad block_id "); err = -1; }
if (!ch || ch_len != 32 || memcmp(ch, content_hash, 32) != 0) { printf("bad content_hash "); err = -1; }
if (!cid || strcmp(cid, chat_id) != 0) { printf("bad chat_id "); err = -1; }
if (!loc || strcmp(loc, location) != 0) { printf("bad location '%s'!='%s' ", loc ? loc : "NULL", location); err = -1; }
if (nid != node_id) { printf("bad node_id %llu!=%llu ", (unsigned long long)nid, (unsigned long long)node_id); err = -1; }
if (!nsig || nsig_len != 64) { printf("bad node_sign "); err = -1; }
if (ts <= 0) { printf("bad timestamp %lld ", (long long)ts); err = -1; }
if (fsz != file_size) { printf("bad file_size %lld!=%lld ", (long long)fsz, (long long)file_size); err = -1; }
if (csz != chunk_size) { printf("bad chunk_size %lld!=%lld ", (long long)csz, (long long)chunk_size); err = -1; }
if (c != chunk) { printf("bad chunk %d!=%d ", c, chunk); err = -1; }
if (off != offset) { printf("bad offset %lld!=%lld ", (long long)off, (long long)offset); err = -1; }
if (err == 0 && nsig && node_pubkey) {
uint8_t sig_msg[92];
size_t soff = 0;
memcpy(sig_msg + soff, media_id, 16); soff += 16;
memcpy(sig_msg + soff, content_hash, 32); soff += 32;
memcpy(sig_msg + soff, &ts, 8); soff += 8;
memcpy(sig_msg + soff, &fsz, 8); soff += 8;
memcpy(sig_msg + soff, &csz, 8); soff += 8;
int32_t cn = (int32_t)chunk; memcpy(sig_msg + soff, &cn, 4); soff += 4;
memcpy(sig_msg + soff, &off, 8); soff += 8;
memcpy(sig_msg + soff, &node_id, 8); soff += 8;
if (ed25519_verify(node_pubkey, sig_msg, soff, nsig) != 0)
{ printf("bad node_sign verify "); err = -1; }
}
sqlite3_finalize(st);
return err;
}
/* ─── тесты ─── */
static int test_single_chunk(void) {
TEST("single chunk, exact size (100/100)");
sqlite3* db = make_db();
media_index_init(db);
uint8_t priv[32], pub[32];
ma_uuid(priv); ma_uuid(priv + 16);
derive_pubkey(priv, pub);
struct media_index_result r;
fill_result(&r, 100, 100, 1);
int rc = media_index_commit(db, &r, NODE_ID, priv, "ch1",
"/data/ch1/file.bin", "/data");
if (rc != 0) { FAIL("commit=%d", rc); media_index_result_free(&r); free_db(db); return -1; }
if (row_count(db, r.media_id) != 1) { FAIL("row count != 1"); media_index_result_free(&r); free_db(db); return -1; }
int cr = check_row(db, r.media_id, 0, r.content_hash, 100, 100, 0, NODE_ID, "ch1",
"ch1/file.bin", pub);
if (cr != 0) { FAIL("check_row"); media_index_result_free(&r); free_db(db); return -1; }
OK();
media_index_result_free(&r);
free_db(db);
return 0;
}
static int test_exact_multiple(void) {
TEST("3 chunks, exact match (300/100)");
sqlite3* db = make_db();
media_index_init(db);
uint8_t priv[32], pub[32];
ma_uuid(priv); ma_uuid(priv + 16); derive_pubkey(priv, pub);
struct media_index_result r;
fill_result(&r, 300, 100, 3);
if (media_index_commit(db, &r, NODE_ID, priv, "ch2", "/data/ch2/file.bin", "/data") != 0)
{ FAIL("commit"); media_index_result_free(&r); free_db(db); return -1; }
if (row_count(db, r.media_id) != 3) { FAIL("row count != 3"); media_index_result_free(&r); free_db(db); return -1; }
if (check_row(db, r.media_id, 0, r.content_hash, 300, 100, 0, NODE_ID, "ch2", "ch2/file.bin", pub) != 0)
{ FAIL("chunk 0"); media_index_result_free(&r); free_db(db); return -1; }
if (check_row(db, r.media_id, 1, r.content_hash, 300, 100, 100, NODE_ID, "ch2", "ch2/file.bin", pub) != 0)
{ FAIL("chunk 1"); media_index_result_free(&r); free_db(db); return -1; }
if (check_row(db, r.media_id, 2, r.content_hash, 300, 100, 200, NODE_ID, "ch2", "ch2/file.bin", pub) != 0)
{ FAIL("chunk 2"); media_index_result_free(&r); free_db(db); return -1; }
OK();
media_index_result_free(&r); free_db(db);
return 0;
}
static int test_partial_last(void) {
TEST("last chunk smaller (250/100 → 3 chunks)");
sqlite3* db = make_db();
media_index_init(db);
uint8_t priv[32], pub[32];
ma_uuid(priv); ma_uuid(priv + 16); derive_pubkey(priv, pub);
struct media_index_result r;
fill_result(&r, 250, 100, 3);
if (media_index_commit(db, &r, NODE_ID, priv, "ch3", "/data/ch3/file.bin", "/data") != 0)
{ FAIL("commit"); media_index_result_free(&r); free_db(db); return -1; }
if (row_count(db, r.media_id) != 3) { FAIL("row count != 3"); media_index_result_free(&r); free_db(db); return -1; }
if (check_row(db, r.media_id, 0, r.content_hash, 250, 100, 0, NODE_ID, "ch3", "ch3/file.bin", pub) != 0)
{ FAIL("chunk 0"); media_index_result_free(&r); free_db(db); return -1; }
if (check_row(db, r.media_id, 1, r.content_hash, 250, 100, 100, NODE_ID, "ch3", "ch3/file.bin", pub) != 0)
{ FAIL("chunk 1"); media_index_result_free(&r); free_db(db); return -1; }
if (check_row(db, r.media_id, 2, r.content_hash, 250, 50, 200, NODE_ID, "ch3", "ch3/file.bin", pub) != 0)
{ FAIL("chunk 2 (last, chunk_size=50)"); media_index_result_free(&r); free_db(db); return -1; }
OK();
media_index_result_free(&r); free_db(db);
return 0;
}
static int test_file_smaller_than_block(void) {
TEST("file smaller than block (50/100 → 1 chunk)");
sqlite3* db = make_db();
media_index_init(db);
uint8_t priv[32], pub[32];
ma_uuid(priv); ma_uuid(priv + 16); derive_pubkey(priv, pub);
struct media_index_result r;
fill_result(&r, 50, 100, 1);
if (media_index_commit(db, &r, NODE_ID, priv, "ch4", "/data/ch4/file.bin", "/data") != 0)
{ FAIL("commit"); media_index_result_free(&r); free_db(db); return -1; }
if (row_count(db, r.media_id) != 1) { FAIL("row count != 1"); media_index_result_free(&r); free_db(db); return -1; }
if (check_row(db, r.media_id, 0, r.content_hash, 50, 50, 0, NODE_ID, "ch4", "ch4/file.bin", pub) != 0)
{ FAIL("chunk 0"); media_index_result_free(&r); free_db(db); return -1; }
OK();
media_index_result_free(&r); free_db(db);
return 0;
}
static int test_many_chunks(void) {
TEST("8 chunks (800/100)");
sqlite3* db = make_db();
media_index_init(db);
uint8_t priv[32], pub[32];
ma_uuid(priv); ma_uuid(priv + 16); derive_pubkey(priv, pub);
struct media_index_result r;
fill_result(&r, 800, 100, 8);
if (media_index_commit(db, &r, NODE_ID, priv, "ch5", "/data/ch5/file.bin", "/data") != 0)
{ FAIL("commit"); media_index_result_free(&r); free_db(db); return -1; }
if (row_count(db, r.media_id) != 8) { FAIL("row count != 8"); media_index_result_free(&r); free_db(db); return -1; }
// все block_id уникальны
for (int i = 0; i < 8; i++) {
for (int j = i + 1; j < 8; j++) {
if (memcmp(r.block_ids + i * 16, r.block_ids + j * 16, 16) == 0)
{ FAIL("block_id collision i=%d j=%d", i, j); media_index_result_free(&r); free_db(db); return -1; }
}
}
for (int i = 0; i < 8; i++) {
int64_t cs = 100, off = i * 100;
if (check_row(db, r.media_id, i, r.content_hash, 800, cs, off, NODE_ID, "ch5", "ch5/file.bin", pub) != 0)
{ FAIL("chunk %d", i); media_index_result_free(&r); free_db(db); return -1; }
}
OK();
media_index_result_free(&r); free_db(db);
return 0;
}
static int test_location_trailing_slash(void) {
TEST("media_base with trailing slash");
sqlite3* db = make_db();
media_index_init(db);
uint8_t priv[32], pub[32];
ma_uuid(priv); ma_uuid(priv + 16); derive_pubkey(priv, pub);
struct media_index_result r;
fill_result(&r, 100, 100, 1);
if (media_index_commit(db, &r, NODE_ID, priv, "tc", "/root/media/tc/v.bin", "/root/") != 0)
{ FAIL("commit"); media_index_result_free(&r); free_db(db); return -1; }
if (check_row(db, r.media_id, 0, r.content_hash, 100, 100, 0, NODE_ID, "tc", "media/tc/v.bin", pub) != 0)
{ FAIL("check_row"); media_index_result_free(&r); free_db(db); return -1; }
OK();
media_index_result_free(&r); free_db(db);
return 0;
}
static int test_location_no_slash(void) {
TEST("media_base without trailing slash");
sqlite3* db = make_db();
media_index_init(db);
uint8_t priv[32], pub[32];
ma_uuid(priv); ma_uuid(priv + 16); derive_pubkey(priv, pub);
struct media_index_result r;
fill_result(&r, 100, 100, 1);
if (media_index_commit(db, &r, NODE_ID, priv, "tc2", "/root/media/tc2/v.bin", "/root") != 0)
{ FAIL("commit"); media_index_result_free(&r); free_db(db); return -1; }
if (check_row(db, r.media_id, 0, r.content_hash, 100, 100, 0, NODE_ID, "tc2", "media/tc2/v.bin", pub) != 0)
{ FAIL("check_row"); media_index_result_free(&r); free_db(db); return -1; }
OK();
media_index_result_free(&r); free_db(db);
return 0;
}
static int test_different_chat_ids(void) {
TEST("two commits with different chat_ids");
sqlite3* db = make_db();
media_index_init(db);
uint8_t priv[32], pub[32];
ma_uuid(priv); ma_uuid(priv + 16); derive_pubkey(priv, pub);
struct media_index_result r1, r2;
fill_result(&r1, 100, 100, 1);
fill_result(&r2, 200, 100, 2);
if (media_index_commit(db, &r1, NODE_ID, priv, "chatA", "/d/cA/f.bin", "/d") != 0)
{ FAIL("commit1"); goto err; }
if (media_index_commit(db, &r2, NODE_ID + 1, priv, "chatB", "/d/cB/f.bin", "/d") != 0)
{ FAIL("commit2"); goto err; }
if (row_count(db, r1.media_id) != 1) { FAIL("r1 count"); goto err; }
if (row_count(db, r2.media_id) != 2) { FAIL("r2 count"); goto err; }
if (check_row(db, r1.media_id, 0, r1.content_hash, 100, 100, 0, NODE_ID, "chatA", "cA/f.bin", pub) != 0)
{ FAIL("r1 row"); goto err; }
if (check_row(db, r2.media_id, 0, r2.content_hash, 200, 100, 0, NODE_ID+1, "chatB", "cB/f.bin", pub) != 0)
{ FAIL("r2 row 0"); goto err; }
if (check_row(db, r2.media_id, 1, r2.content_hash, 200, 100, 100, NODE_ID+1, "chatB", "cB/f.bin", pub) != 0)
{ FAIL("r2 row 1"); goto err; }
OK();
media_index_result_free(&r1); media_index_result_free(&r2); free_db(db);
return 0;
err:
media_index_result_free(&r1); media_index_result_free(&r2); free_db(db);
return -1;
}
/* ─── интеграционный тест: register_async ─── */
struct async_ctx {
volatile int done;
int err;
struct media_index_result result;
};
static void on_registered(void* arg, int err, const struct media_index_result* r) {
struct async_ctx* ctx = (struct async_ctx*)arg;
ctx->err = err;
if (!err && r) {
memcpy(&ctx->result, r, sizeof(*r));
if (r->block_ids) {
ctx->result.block_ids = u_malloc((size_t)r->num_blocks * 16);
memcpy(ctx->result.block_ids, r->block_ids, (size_t)r->num_blocks * 16);
}
if (r->block_sigs) {
ctx->result.block_sigs = u_malloc((size_t)r->num_blocks * 64);
memcpy(ctx->result.block_sigs, r->block_sigs, (size_t)r->num_blocks * 64);
}
} else memset(&ctx->result, 0, sizeof(ctx->result));
ctx->done = 1;
}
static void on_timeout(void* arg) {
volatile int* done = (volatile int*)arg;
if (done) *done = 1;
}
static int create_file(const char* path, int size) {
FILE* f = fopen(path, "wb");
if (!f) return -1;
for (int i = 0; i < size; i++) fputc((i * 7 + 13) & 0xFF, f);
fclose(f);
return 0;
}
static int test_register_async_flow(void) {
TEST("register_async: copy file + verify DB");
char tempd[256];
snprintf(tempd, sizeof(tempd), "/tmp/test_media_idx_%d", (int)getpid());
mkdir(tempd, 0755);
char src[512]; snprintf(src, sizeof(src), "%s/src.bin", tempd);
char dst[512]; snprintf(dst, sizeof(dst), "%s/media/ch_test/dst.bin", tempd);
{
char t1[512], t2[512];
snprintf(t1, sizeof(t1), "%s/media", tempd); mkdir(t1, 0755);
snprintf(t2, sizeof(t2), "%s/media/ch_test", tempd); mkdir(t2, 0755);
}
create_file(src, 400); // 400 bytes → 4 chunks of 100
struct UASYNC* ua = uasync_create();
if (!ua) { FAIL("uasync_create"); rmdir(tempd); return -1; }
struct async_ctx ctx = {0};
void* th = uasync_set_timeout(ua, 50000, &ctx.done, on_timeout, "tout_async");
struct media_async* ma = media_async_create();
sqlite3* db = make_db();
media_index_init(db);
uint8_t priv[32], pub[32];
ma_uuid(priv); ma_uuid(priv + 16); derive_pubkey(priv, pub);
media_index_register_async(ma, ua, db, 0xCAFE, priv, "ch_test",
src, dst, 1, tempd, on_registered, &ctx);
while (!ctx.done) uasync_poll(ua, 100);
uasync_cancel_timeout(ua, th);
if (ctx.err != 0) { FAIL("register_async err=%d", ctx.err); goto cleanup; }
{ struct stat st; if (stat(dst, &st) != 0) { FAIL("dst file not created"); goto cleanup; } }
int nb = ctx.result.num_blocks;
int64_t bs = ctx.result.block_size;
int64_t fs = ctx.result.file_size;
if (fs != 400) { FAIL("file_size=%lld != 400", (long long)fs); goto cleanup; }
if (nb < 1 || nb > 400) { FAIL("num_blocks=%d unexpected", nb); goto cleanup; }
if (row_count(db, ctx.result.media_id) != nb) { FAIL("row count %d != %d", row_count(db, ctx.result.media_id), nb); goto cleanup; }
for (int i = 0; i < nb; i++) {
int64_t cs = (i == nb - 1) ? fs - i * bs : bs;
int64_t off = i * bs;
if (check_row(db, ctx.result.media_id, i, ctx.result.content_hash, fs, cs, off, 0xCAFE, "ch_test",
"media/ch_test/dst.bin", pub) != 0)
{ FAIL("chunk %d", i); goto cleanup; }
}
// block_ids уникальны
for (int i = 0; i < nb; i++)
for (int j = i+1; j < nb; j++)
if (memcmp(ctx.result.block_ids + i*16, ctx.result.block_ids + j*16, 16) == 0)
{ FAIL("block_id collision i=%d j=%d", i, j); goto cleanup; }
OK();
cleanup:
media_index_result_free(&ctx.result);
free_db(db);
media_async_destroy(ma);
uasync_destroy(ua, 0);
unlink(src); unlink(dst);
{ char t1[512]; snprintf(t1,sizeof(t1),"%s/media/ch_test",tempd); rmdir(t1); }
{ char t2[512]; snprintf(t2,sizeof(t2),"%s/media",tempd); rmdir(t2); }
rmdir(tempd);
return 0;
}
/* ─── main ─── */
int main(void) {
debug_set_level(DEBUG_LEVEL_ERROR);
printf("=== test_media_index ===\n\n");
test_single_chunk();
test_exact_multiple();
test_partial_last();
test_file_smaller_than_block();
test_many_chunks();
test_location_trailing_slash();
test_location_no_slash();
test_different_chat_ids();
test_register_async_flow();
printf("\nResults: %d/%d passed, %d failed\n", g_passed, g_total, g_failed);
fflush(stdout);
return g_failed ? 1 : 0;
}

23
tools/chatgui-android/app/src/main/java/com/utun/chat/data/ChatRepository.kt

@ -9,7 +9,7 @@ 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<Float> = emptyList(),
val voiceDurationMs: Int = 0)
val voiceDurationMs: Int = 0, val fileSize: Long = 0)
class ChatRepository {
private val myNodeId: Long = NativeLib.getMyNodeId()
@ -53,6 +53,9 @@ class ChatRepository {
val contentType = obj.optString("contentType", "text/plain")
val filePath = obj.optString("filePath", "")
var displayText = text
var fileSize = 0L
/* Parse base64 waveform from voice message text (format: "wf=<100 chars A-Za-z0-9+/>;dur=1234;") */
var waveform = emptyList<Float>()
var voiceDurationMs = 0
@ -88,16 +91,30 @@ class ChatRepository {
} catch (_: Exception) {}
}
/* Parse file attachment: "filename.ext|file_size|block_size|num_blocks|sigs" */
if (contentType == "application/octet-stream") {
try {
val pipeIdx = text.indexOf('|')
if (pipeIdx > 0) {
displayText = text.substring(0, pipeIdx)
val metaParts = text.substring(pipeIdx + 1).split("|")
if (metaParts.isNotEmpty())
fileSize = metaParts[0].toLongOrNull() ?: 0L
}
} catch (_: Exception) {}
}
list.add(Message(
id = obj.getLong("id"),
author = obj.getString("author"),
text = text,
text = displayText,
ts = obj.getLong("ts"),
isOutgoing = obj.optBoolean("isOutgoing", false),
contentType = contentType,
filePath = filePath,
waveform = waveform,
voiceDurationMs = voiceDurationMs
voiceDurationMs = voiceDurationMs,
fileSize = fileSize
))
}
} catch (e: Exception) { Log.w("utun-gui", "getMessages error", e) }

9
tools/chatgui-android/app/src/main/java/com/utun/chat/data/ConfigProvider.kt

@ -36,6 +36,9 @@ class ConfigProvider(private val context: Context) {
private val logUdpIp = stringPreferencesKey("log_udp_ip")
private val logUdpPort = intPreferencesKey("log_udp_port")
private val serversJson = stringPreferencesKey("servers")
private val storageAutoload = intPreferencesKey("storage_autoload")
private val storageAutoloadMaxsizeMb = intPreferencesKey("storage_autoload_maxsize_mb")
private val storageMaxsizeGb = intPreferencesKey("storage_maxsize_gb")
private val firstLaunchDone = booleanPreferencesKey("first_launch_done")
private val gson = Gson()
@ -130,6 +133,9 @@ class ConfigProvider(private val context: Context) {
key == "node.listen_port" -> context.dataStore.data.first()[nodeListenPort] ?: 12345
key == "control.port" -> context.dataStore.data.first()[controlPort] ?: 9999
key == "log_udp.port" -> context.dataStore.data.first()[logUdpPort] ?: 9999
key == "storage_autoload" -> context.dataStore.data.first()[storageAutoload] ?: 1
key == "storage_autoload_maxsize_mb" -> context.dataStore.data.first()[storageAutoloadMaxsizeMb] ?: 10
key == "storage_maxsize_gb" -> context.dataStore.data.first()[storageMaxsizeGb] ?: 1
else -> 0
}
}
@ -149,6 +155,9 @@ class ConfigProvider(private val context: Context) {
"debug.categories" -> prefs[debugCategories] = value
"log_udp.ip" -> prefs[logUdpIp] = value
"log_udp.port" -> prefs[logUdpPort] = value.toIntOrNull() ?: 0
"storage_autoload" -> prefs[storageAutoload] = value.toIntOrNull() ?: 1
"storage_autoload_maxsize_mb" -> prefs[storageAutoloadMaxsizeMb] = value.toIntOrNull() ?: 10
"storage_maxsize_gb" -> prefs[storageMaxsizeGb] = value.toIntOrNull() ?: 1
}
}
}

22
tools/chatgui-android/app/src/main/java/com/utun/chat/ui/components/MessageBubble.kt

@ -62,7 +62,7 @@ fun MessageBubble(message: Message, playbackState: ChatViewModel.PlaybackState =
when {
isVoice -> VoiceBubble(message, textColor, isActive, playbackState, onPlayVoice, onSeekVoice)
isMedia -> Text("\uD83D\uDCC4 ${message.text}", fontSize = 15.sp, color = textColor)
isMedia -> FileBubble(message, textColor)
else -> {
Text(message.author, fontSize = 13.sp, fontWeight = FontWeight.Bold, color = textColor.copy(alpha = 0.7f))
Spacer(Modifier.height(2.dp))
@ -167,6 +167,26 @@ 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)
Spacer(Modifier.height(2.dp))
Text(formatFileSize(message.fileSize), fontSize = 12.sp, color = textColor.copy(alpha = 0.5f))
}
}
private fun formatFileSize(bytes: Long): String {
if (bytes <= 0) return ""
if (bytes < 1024) return "${bytes}B"
val kb = bytes / 1024.0
if (kb < 1024.0) return if (kb < 10.0) "%.1fKB".format(kb) else "%.0fKB".format(kb)
val mb = kb / 1024.0
if (mb < 1024.0) return if (mb < 10.0) "%.1fMB".format(mb) else "%.0fMB".format(mb)
val gb = mb / 1024.0
return "%.1fGB".format(gb)
}
private fun formatDurationMs(ms: Int): String {
val secs = ms / 1000
return "%d:%02d".format(secs / 60, secs % 60)

49
tools/chatgui-android/app/src/main/java/com/utun/chat/ui/screens/SettingsScreen.kt

@ -42,7 +42,7 @@ fun SettingsScreen(vm: ChatViewModel, onBack: () -> Unit, onExit: () -> Unit = {
var dirty by remember { mutableStateOf(false) }
var tabIndex by remember { mutableStateOf(0) }
val tabs = listOf("Settings", "Audio")
val tabs = listOf("Settings", "Audio", "Storage")
val handleBack = {
if (dirty) {
@ -73,6 +73,7 @@ fun SettingsScreen(vm: ChatViewModel, onBack: () -> Unit, onExit: () -> Unit = {
when (tabIndex) {
0 -> GeneralSettingsTab(vm, dirty, onDirty = { dirty = it }, onExit)
1 -> AudioSettingsTab(vm)
2 -> StorageSettingsTab()
}
}
}
@ -414,6 +415,52 @@ private fun AudioSettingsTab(vm: ChatViewModel) {
}
}
/* ── Storage Settings Tab ── */
@Composable
private fun StorageSettingsTab() {
val provider = ChatApplication.instance.configProvider
val scope = rememberCoroutineScope()
var autoload by remember { mutableStateOf(true) }
var autoloadMaxSizeMb by remember { mutableStateOf("10") }
var maxStorageGb by remember { mutableStateOf("1") }
LaunchedEffect(Unit) {
autoload = provider.getInt("storage_autoload") != 0
autoloadMaxSizeMb = provider.getInt("storage_autoload_maxsize_mb").toString()
maxStorageGb = provider.getInt("storage_maxsize_gb").toString()
}
Column(
modifier = Modifier.padding(16.dp).verticalScroll(rememberScrollState()),
verticalArrangement = Arrangement.spacedBy(12.dp)
) {
Text("Media Storage", style = MaterialTheme.typography.titleMedium)
Row(verticalAlignment = Alignment.CenterVertically) {
Text("Autoload media", modifier = Modifier.weight(1f))
Switch(checked = autoload, onCheckedChange = {
autoload = it
scope.launch { provider.setValue("storage_autoload", if (it) "1" else "0") }
})
}
HorizontalDivider(Modifier.padding(vertical = 4.dp))
OutlinedTextField(autoloadMaxSizeMb,
{ v -> autoloadMaxSizeMb = v; val nb = v.toIntOrNull(); if (nb != null && nb in 1..500) scope.launch { provider.setValue("storage_autoload_maxsize_mb", v) } },
label = { Text("Autoload max size (MB)") },
keyboardOptions = KeyboardOptions(keyboardType = KeyboardType.Number),
singleLine = true, modifier = Modifier.padding(top = 4.dp))
OutlinedTextField(maxStorageGb,
{ v -> maxStorageGb = v; val nb = v.toIntOrNull(); if (nb != null && nb in 1..100) scope.launch { provider.setValue("storage_maxsize_gb", v) } },
label = { Text("Max storage size (GB)") },
keyboardOptions = KeyboardOptions(keyboardType = KeyboardType.Number),
singleLine = true)
}
}
/* ── Labeled Slider ── */
@Composable

4
tools/chatgui-android/jni_bridge/android_jni_bridge.c

@ -489,8 +489,8 @@ int utun_bridge_voice_get_compressor(void) {
}
void utun_bridge_voice_set_compressor_config(int max_gain_db, int lookback_ms, int lookahead_ms,
float rise_rate_per_500ms, float target_level_db) {
voice_recorder_set_compressor_config(max_gain_db, lookback_ms, lookahead_ms, rise_rate_per_500ms, target_level_db);
float rise_rate_per_sec, float target_level_db) {
voice_recorder_set_compressor_config(max_gain_db, lookback_ms, lookahead_ms, rise_rate_per_sec, target_level_db);
}
float utun_bridge_voice_get_peak_level(void) {

2
tools/chatgui-android/jni_bridge/android_jni_bridge.h

@ -86,7 +86,7 @@ int utun_bridge_voice_get_preset(void);
void utun_bridge_voice_set_compressor(int enabled);
int utun_bridge_voice_get_compressor(void);
void utun_bridge_voice_set_compressor_config(int max_gain_db, int lookback_ms, int lookahead_ms,
float rise_rate_per_500ms, float target_level_db);
float rise_rate_per_sec, float target_level_db);
float utun_bridge_voice_get_peak_level(void);
/* ── Attachment ── */

75
tools/chatgui-android/libutun_lite/attachment_sender.c

@ -12,18 +12,6 @@
#include <unistd.h>
#include <libgen.h>
#define MEDIA_BLOCK_MIN (10 * 1024 * 1024)
#define MEDIA_BLOCK_MAX (25 * 1024 * 1024)
#define MEDIA_BLOCK_TARGET 30
static uint32_t media_calc_block_size(uint64_t file_size) {
if (file_size == 0) return 0;
uint64_t target = (file_size + MEDIA_BLOCK_TARGET - 1) / MEDIA_BLOCK_TARGET;
if (target < MEDIA_BLOCK_MIN) target = MEDIA_BLOCK_MIN;
if (target > MEDIA_BLOCK_MAX) target = MEDIA_BLOCK_MAX;
return (uint32_t)target;
}
int attachment_send(const char* channel_id, const char* src_file_path, const char* db_path) {
if (!channel_id || !src_file_path || !db_path) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "attachment_send: invalid args");
@ -33,27 +21,15 @@ int attachment_send(const char* channel_id, const char* src_file_path, const cha
struct UASYNC* ua = instance_lite_get_uasync();
if (!ua) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "attachment_send: no uasync"); return -1; }
FILE* src = fopen(src_file_path, "rb");
if (!src) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "attachment_send: cannot open %s", src_file_path); return -1; }
fseek(src, 0, SEEK_END);
uint64_t file_size = (uint64_t)ftell(src);
fseek(src, 0, SEEK_SET);
if (file_size == 0) { DEBUG_WARN(DEBUG_CATEGORY_GENERAL, "attachment_send: empty file %s", src_file_path); fclose(src); return -1; }
/* extract basename and ext from path */
/* extract basename and ext */
char path_copy[1024];
snprintf(path_copy, sizeof(path_copy), "%s", src_file_path);
char* fname = basename(path_copy);
char* dot = strrchr(fname, '.');
char ext[32] = "";
if (dot) {
snprintf(ext, sizeof(ext), "%s", dot + 1);
*dot = '\0';
}
if (dot) { snprintf(ext, sizeof(ext), "%s", dot + 1); *dot = '\0'; }
/* generate dt/suffix */
/* generate dt/suffix for unique filename */
time_t t = time(NULL);
struct tm tm_buf; localtime_r(&t, &tm_buf);
char dt[32]; snprintf(dt, sizeof(dt), "%04d%02d%02d-%02d%02d%02d",
@ -62,10 +38,7 @@ int attachment_send(const char* channel_id, const char* src_file_path, const cha
int rnd = rand() & 0xFFFF;
char suffix[8]; snprintf(suffix, sizeof(suffix), "%04x", rnd);
uint32_t block_size = media_calc_block_size(file_size);
int num_blocks = (int)((file_size + block_size - 1) / block_size);
/* ensure media dir */
/* ensure media dir exists */
char media_dir[1024];
snprintf(media_dir, sizeof(media_dir), "%s/media/%s", db_path, channel_id);
{
@ -77,46 +50,26 @@ int attachment_send(const char* channel_id, const char* src_file_path, const cha
mkdir(tmp, 0755);
}
/* write blocks */
for (int n = 0; n < num_blocks; n++) {
char block_path[1280];
snprintf(block_path, sizeof(block_path), "%s/%s_%d_%s_%s.%s",
media_dir, dt, n, fname, suffix, ext[0] ? ext : "bin");
FILE* dst = fopen(block_path, "wb");
if (!dst) {
DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "attachment_send: cannot create block %s", block_path);
continue;
}
uint8_t buf[65536];
uint64_t remaining = block_size;
while (remaining > 0) {
size_t to_read = remaining < sizeof(buf) ? (size_t)remaining : sizeof(buf);
size_t rd = fread(buf, 1, to_read, src);
if (rd == 0) break;
fwrite(buf, 1, rd, dst);
remaining -= rd;
}
fclose(dst);
}
fclose(src);
/* single file path — no splitting */
char media_file[1280];
snprintf(media_file, sizeof(media_file), "%s/%s_%s_%s.%s",
media_dir, dt, fname, suffix, ext[0] ? ext : "bin");
/* post chat_msg_submit */
/* post chat_msg_submit — media_copy=1 so media_index copies src to dest */
struct chat_msg_submit* req = u_calloc(1, sizeof(struct chat_msg_submit) + 1);
if (!req) return -1;
snprintf(req->channel_id, sizeof(req->channel_id), "%s", channel_id);
snprintf(req->content_type, sizeof(req->content_type), "application/octet-stream");
snprintf(req->media_dt, sizeof(req->media_dt), "%s", dt);
snprintf(req->media_basename, sizeof(req->media_basename), "%s", fname);
snprintf(req->media_suffix, sizeof(req->media_suffix), "%s", suffix);
snprintf(req->media_ext, sizeof(req->media_ext), "%s", ext[0] ? ext : "bin");
req->media_num_blocks = (uint32_t)num_blocks;
snprintf(req->media_src, sizeof(req->media_src), "%s", src_file_path);
snprintf(req->media_dest, sizeof(req->media_dest), "%s", media_file);
req->media_copy = 1;
req->data = (uint8_t*)(req + 1);
req->data_len = 0;
req->timestamp = 0;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "attachment_send: ch=%s file=%s blocks=%d size=%llu",
channel_id, src_file_path, num_blocks, (unsigned long long)file_size);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "attachment_send: ch=%s src=%s dst=%s",
channel_id, src_file_path, media_file);
uasync_post(ua, chat_core_submit_trampoline, req);
return 0;

2
tools/chatgui-android/libutun_lite/utun_sources.cmake

@ -58,6 +58,8 @@ function(utun_setup_sources LIB_DIR SRC_DIR CONFIG_DIR TOOLS_DIR)
${SRC_DIR}/proxy
${SRC_DIR}/lwip_tcp
${SRC_DIR}/BBR
${SRC_DIR}/media_delivery
${SRC_DIR}/media_async
${TOOLS_DIR}
${CONFIG_DIR}
PARENT_SCOPE

68
tools/chatgui-android/libutun_lite/voice_recorder.c

@ -17,9 +17,6 @@
#include <unistd.h>
#include <errno.h>
#define MEDIA_BLOCK_MIN (10 * 1024 * 1024)
#define MEDIA_BLOCK_MAX (25 * 1024 * 1024)
#define MEDIA_BLOCK_TARGET 30
#define OPUS_MAGIC 0x5355504F
#define FRAME_MS 20
#define MAX_PACKET 4000
@ -85,9 +82,9 @@ int voice_recorder_init(const char* db_path) {
cfg.block_duration_ms = 20;
cfg.lookback_ms = 200;
cfg.lookahead_ms = 100;
cfg.max_gain_db = 30.0f;
cfg.rise_rate_per_500ms = 2.0f;
cfg.target_level = 0.25f;
cfg.max_gain_db = 25.0f;
cfg.rise_rate_per_sec = 10.0f;
cfg.target_level = 1.0f;
audio_compressor_configure(g_rec->compressor, &cfg);
}
@ -199,14 +196,6 @@ int voice_recorder_feed(const int16_t* samples, int count) {
return 0;
}
static uint32_t media_calc_block_size(uint64_t file_size) {
if (file_size == 0) return 0;
uint64_t target = (file_size + MEDIA_BLOCK_TARGET - 1) / MEDIA_BLOCK_TARGET;
if (target < MEDIA_BLOCK_MIN) target = MEDIA_BLOCK_MIN;
if (target > MEDIA_BLOCK_MAX) target = MEDIA_BLOCK_MAX;
return (uint32_t)target;
}
static void voice_cleanup_locked(struct voice_recorder* rec) {
u_free(rec->pcm_buffer); rec->pcm_buffer = NULL;
rec->pcm_count = 0; rec->pcm_cap = 0;
@ -364,37 +353,8 @@ int voice_recorder_stop(int* out_duration_ms) {
return -1;
}
/* get file size and calculate blocks */
FILE* fsize = fopen(temp_path, "rb"); uint64_t file_size = 0;
if (fsize) { fseek(fsize, 0, SEEK_END); file_size = (uint64_t)ftell(fsize); fclose(fsize); }
uint32_t block_size = media_calc_block_size(file_size);
int num_blocks = (int)((file_size + block_size - 1) / block_size);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "voice_recorder_stop: file_size=%llu block_size=%u num_blocks=%d",
(unsigned long long)file_size, block_size, num_blocks);
/* split file into blocks */
FILE* src = fopen(temp_path, "rb");
if (!src) { unlink(temp_path); voice_cleanup_locked(g_rec); pthread_mutex_unlock(&g_rec->mtx); pthread_mutex_unlock(&g_init_mtx); return -1; }
for (int n = 0; n < num_blocks; n++) {
char block_path[1280];
snprintf(block_path, sizeof(block_path), "%s/%s_%d_%s_%s.opus",
media_dir, dt, n, basename, suffix);
FILE* dst = fopen(block_path, "wb");
if (!dst) { DEBUG_ERROR(DEBUG_CATEGORY_GENERAL, "voice_recorder_stop: cannot create block %s", block_path); continue; }
uint8_t buf[65536];
uint64_t remaining = block_size;
while (remaining > 0) {
size_t rd = fread(buf, 1, remaining < sizeof(buf) ? (size_t)remaining : sizeof(buf), src);
if (rd == 0) break;
fwrite(buf, 1, rd, dst);
remaining -= rd;
}
fclose(dst);
}
fclose(src);
unlink(temp_path);
/* file is already at temp_path in media dir – no splitting needed */
/* media_index_register_async will handle block size calculation, signing, and DB */
/* build compact waveform: 100 bytes base64 (A-Za-z0-9+/), log scale */
static const char b64_table[] = "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789+/";
@ -421,18 +381,16 @@ int voice_recorder_stop(int* out_duration_ms) {
snprintf(req->channel_id, sizeof(req->channel_id), "%s", g_rec->channel_id);
snprintf(req->content_type, sizeof(req->content_type), "audio/opus");
snprintf(req->media_dt, sizeof(req->media_dt), "%s", dt);
snprintf(req->media_basename, sizeof(req->media_basename), "%s", basename);
snprintf(req->media_suffix, sizeof(req->media_suffix), "%s", suffix);
snprintf(req->media_ext, sizeof(req->media_ext), "opus");
req->media_num_blocks = (uint32_t)num_blocks;
snprintf(req->media_src, sizeof(req->media_src), "%s", temp_path);
snprintf(req->media_dest, sizeof(req->media_dest), "%s", temp_path);
req->media_copy = 0;
req->data = (uint8_t*)(req + 1);
req->data_len = (uint32_t)total_data_len;
memcpy(req->data, wf_data, total_data_len);
req->timestamp = 0;
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "voice_recorder_stop: posting media ch=%s dt=%s stem=%s blocks=%d",
req->channel_id, dt, basename, num_blocks);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "voice_recorder_stop: posting media ch=%s path=%s",
req->channel_id, temp_path);
uasync_post(ua, chat_core_submit_trampoline, req);
@ -505,7 +463,7 @@ int voice_recorder_is_compressor_enabled(void) {
}
void voice_recorder_set_compressor_config(int max_gain_db, int lookback_ms, int lookahead_ms,
float rise_rate_per_500ms, float target_level_db) {
float rise_rate_per_sec, float target_level_db) {
pthread_mutex_lock(&g_init_mtx);
if (g_rec && g_rec->compressor) {
pthread_mutex_lock(&g_rec->mtx);
@ -516,11 +474,11 @@ void voice_recorder_set_compressor_config(int max_gain_db, int lookback_ms, int
cfg.lookback_ms = lookback_ms;
cfg.lookahead_ms = lookahead_ms;
cfg.max_gain_db = (float)max_gain_db;
cfg.rise_rate_per_500ms = rise_rate_per_500ms;
cfg.rise_rate_per_sec = rise_rate_per_sec;
cfg.target_level = powf(10.0f, target_level_db / 20.0f);
audio_compressor_configure(g_rec->compressor, &cfg);
DEBUG_INFO(DEBUG_CATEGORY_GENERAL, "voice_recorder compressor config: gain=%ddB lookback=%d lookahead=%d rise=%.1f target=%.0fdBFS",
max_gain_db, lookback_ms, lookahead_ms, rise_rate_per_500ms, target_level_db);
max_gain_db, lookback_ms, lookahead_ms, rise_rate_per_sec, target_level_db);
pthread_mutex_unlock(&g_rec->mtx);
}
pthread_mutex_unlock(&g_init_mtx);

2
tools/chatgui-android/libutun_lite/voice_recorder.h

@ -24,7 +24,7 @@ int voice_recorder_get_preset(void);
void voice_recorder_set_compressor_enabled(int enabled);
int voice_recorder_is_compressor_enabled(void);
void voice_recorder_set_compressor_config(int max_gain_db, int lookback_ms, int lookahead_ms,
float rise_rate_per_500ms, float target_level_db);
float rise_rate_per_sec, float target_level_db);
float voice_recorder_get_peak_level(void);
#ifdef __cplusplus

1
tools/chatgui/CMakeLists.txt

@ -119,6 +119,7 @@ add_executable(chatgui
src/voiceplayback.cpp
src/media_blocks.cpp
src/audiodevicesettingspage.cpp
src/storagesettingspage.cpp
transport/utun_node.cpp
transport/node_config.cpp
transport/config_updater.cpp

2
tools/chatgui/libutun/CMakeLists.txt

@ -67,6 +67,8 @@ target_include_directories(utun PUBLIC
${SRC_DIR}
${SRC_DIR}/transport_layer
${SRC_DIR}/routing_layer
${SRC_DIR}/media_delivery
${SRC_DIR}/media_async
${LIB_DIR}
${SRC_DIR}/uip
${CMAKE_SOURCE_DIR}/db

65
tools/chatgui/src/audiodevicesettingspage.cpp

@ -158,11 +158,8 @@ void AudioDeviceSettingsPage::loadFromDb() {
m_savedInputIdx = m_db->getUiStateInt("audio_input_device", -1);
m_savedCodecPreset = m_db->getUiStateInt("opus_codec_preset", 1);
m_savedCompressorEnabled = m_db->getUiStateInt("compressor_enabled", 0);
m_savedCompressorMaxGainDb = m_db->getUiStateInt("compressor_max_gain_db", 30);
m_savedCompressorLookbackMs = m_db->getUiStateInt("compressor_lookback_ms", 200);
m_savedCompressorLookaheadMs = m_db->getUiStateInt("compressor_lookahead_ms", 100);
m_savedCompressorRiseRateTenths = m_db->getUiStateInt("compressor_rise_rate_tenths", 20);
m_savedCompressorTargetLevelDb = m_db->getUiStateInt("compressor_target_level_db", -12);
m_savedCompressorMaxGainDb = m_db->getUiStateInt("compressor_max_gain_db", 25);
m_savedCompressorRiseRate = m_db->getUiStateInt("compressor_rise_rate", 10);
}
/* Select saved indices */
@ -181,10 +178,7 @@ void AudioDeviceSettingsPage::loadFromDb() {
m_compressorEnabledCb->setChecked(m_savedCompressorEnabled != 0);
m_compressorMaxGainSlider->setValue(m_savedCompressorMaxGainDb);
m_compressorLookbackSlider->setValue(m_savedCompressorLookbackMs);
m_compressorLookaheadSlider->setValue(m_savedCompressorLookaheadMs);
m_compressorRiseRateSlider->setValue(m_savedCompressorRiseRateTenths);
m_compressorTargetLevelSlider->setValue(m_savedCompressorTargetLevelDb);
m_compressorRiseRateSlider->setValue(m_savedCompressorRiseRate);
updateCompressorLabels();
}
@ -211,30 +205,20 @@ void AudioDeviceSettingsPage::applyAndSave() {
int compEnabled = m_compressorEnabledCb->isChecked() ? 1 : 0;
int compMaxGain = m_compressorMaxGainSlider->value();
int compLookback = m_compressorLookbackSlider->value();
int compLookahead = m_compressorLookaheadSlider->value();
int compRiseRateTenths = m_compressorRiseRateSlider->value();
int compTargetLevel = m_compressorTargetLevelSlider->value();
int compRiseRate = m_compressorRiseRateSlider->value();
saveKey("compressor_enabled", compEnabled);
saveKey("compressor_max_gain_db", compMaxGain);
saveKey("compressor_lookback_ms", compLookback);
saveKey("compressor_lookahead_ms", compLookahead);
saveKey("compressor_rise_rate_tenths", compRiseRateTenths);
saveKey("compressor_target_level_db", compTargetLevel);
saveKey("compressor_rise_rate", compRiseRate);
m_savedOutputIdx = outIdx;
m_savedInputIdx = inIdx;
m_savedCodecPreset = codecPreset;
m_savedCompressorEnabled = compEnabled;
m_savedCompressorMaxGainDb = compMaxGain;
m_savedCompressorLookbackMs = compLookback;
m_savedCompressorLookaheadMs = compLookahead;
m_savedCompressorRiseRateTenths = compRiseRateTenths;
m_savedCompressorTargetLevelDb = compTargetLevel;
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "saved output=%d input=%d codec=%d compressor=%d gain=%ddB lookback=%d lookahead=%d rise=%.1fx target=%ddB",
outIdx, inIdx, codecPreset, compEnabled, compMaxGain, compLookback, compLookahead,
(float)compRiseRateTenths / 10.0f, compTargetLevel);
m_savedCompressorRiseRate = compRiseRate;
DEBUG_INFO(DEBUG_CATEGORY_DEBUG, "saved output=%d input=%d codec=%d compressor=%d gain=%ddB rise=%.1fx",
outIdx, inIdx, codecPreset, compEnabled, compMaxGain, (float)compRiseRate);
if (m_recorder) {
m_recorder->setCompressorEnabled(compEnabled != 0);
@ -243,10 +227,10 @@ void AudioDeviceSettingsPage::applyAndSave() {
compCfg.sample_rate = 48000;
compCfg.channels = 1;
compCfg.max_gain_db = (float)compMaxGain;
compCfg.lookback_ms = compLookback;
compCfg.lookahead_ms = compLookahead;
compCfg.rise_rate_per_500ms = (float)compRiseRateTenths / 10.0f;
compCfg.target_level = powf(10.0f, (float)compTargetLevel / 20.0f);
compCfg.lookback_ms = 200;
compCfg.lookahead_ms = 100;
compCfg.rise_rate_per_sec = (float)compRiseRate;
compCfg.target_level = 1.0f;
audio_compressor_configure(m_recorder->compressor(), &compCfg);
}
}
@ -342,37 +326,22 @@ void AudioDeviceSettingsPage::buildCompressorUI(QVBoxLayout* layout) {
};
makeSliderRow("Max gain:", m_compressorMaxGainSlider, m_compressorMaxGainLabel,
0, 30, 1);
makeSliderRow("Lookback time:", m_compressorLookbackSlider, m_compressorLookbackLabel,
50, 500, 10);
makeSliderRow("Lookahead time:", m_compressorLookaheadSlider, m_compressorLookaheadLabel,
20, 500, 10);
makeSliderRow("Rise rate (x/500ms):", m_compressorRiseRateSlider, m_compressorRiseRateLabel,
11, 100, 1);
makeSliderRow("Target level:", m_compressorTargetLevelSlider, m_compressorTargetLevelLabel,
-40, 0, 1);
0, 50, 1);
makeSliderRow("Rise rate (x/sec):", m_compressorRiseRateSlider, m_compressorRiseRateLabel,
2, 10, 1);
auto connectSlider = [this](QSlider* slider) {
connect(slider, &QSlider::valueChanged, this, &AudioDeviceSettingsPage::updateCompressorLabels);
};
connectSlider(m_compressorMaxGainSlider);
connectSlider(m_compressorLookbackSlider);
connectSlider(m_compressorLookaheadSlider);
connectSlider(m_compressorRiseRateSlider);
connectSlider(m_compressorTargetLevelSlider);
}
void AudioDeviceSettingsPage::updateCompressorLabels() {
if (m_compressorMaxGainLabel)
m_compressorMaxGainLabel->setText(QString::number(m_compressorMaxGainSlider->value()) + " dB");
if (m_compressorLookbackLabel)
m_compressorLookbackLabel->setText(QString::number(m_compressorLookbackSlider->value()) + " ms");
if (m_compressorLookaheadLabel)
m_compressorLookaheadLabel->setText(QString::number(m_compressorLookaheadSlider->value()) + " ms");
if (m_compressorRiseRateLabel) {
float v = (float)m_compressorRiseRateSlider->value() / 10.0f;
m_compressorRiseRateLabel->setText(QString::number(v, 'f', 1) + "x");
int v = m_compressorRiseRateSlider->value();
m_compressorRiseRateLabel->setText(QString::number(v) + "x");
}
if (m_compressorTargetLevelLabel)
m_compressorTargetLevelLabel->setText(QString::number(m_compressorTargetLevelSlider->value()) + " dBFS");
}

13
tools/chatgui/src/audiodevicesettingspage.h

@ -45,19 +45,10 @@ private:
QCheckBox* m_compressorEnabledCb;
QSlider* m_compressorMaxGainSlider;
QSlider* m_compressorLookbackSlider;
QSlider* m_compressorLookaheadSlider;
QSlider* m_compressorRiseRateSlider;
QSlider* m_compressorTargetLevelSlider;
QLabel* m_compressorMaxGainLabel;
QLabel* m_compressorLookbackLabel;
QLabel* m_compressorLookaheadLabel;
QLabel* m_compressorRiseRateLabel;
QLabel* m_compressorTargetLevelLabel;
int m_savedCompressorEnabled = 0;
int m_savedCompressorMaxGainDb = 30;
int m_savedCompressorLookbackMs = 200;
int m_savedCompressorLookaheadMs = 100;
int m_savedCompressorRiseRateTenths = 20;
int m_savedCompressorTargetLevelDb = -12;
int m_savedCompressorMaxGainDb = 25;
int m_savedCompressorRiseRate = 10;
};

35
tools/chatgui/src/inputbar.cpp

@ -233,21 +233,15 @@ void InputBar::onPttReleased() {
qPrintable(mediaPath), QDir(mediaPath).exists(), qPrintable(m_channelIdForRecord));
QString fileName = VoiceEncoder::generateFileName();
QString tempPath = mediaPath + "/" + fileName;
QString mediaFile = mediaPath + "/" + fileName;
float durationSec = 0;
int frames = VoiceEncoder::encodeToFile(pcm, pcmSize, 48000, 1, tempPath, durationSec);
int frames = VoiceEncoder::encodeToFile(pcm, pcmSize, 48000, 1, mediaFile, durationSec);
if (frames <= 0) {
QFile::remove(tempPath);
QFile::remove(mediaFile);
return;
}
MediaBlockInfo blk = media_split_to_blocks(tempPath, mediaPath, true);
if (blk.numBlocks <= 0) {
DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "InputBar: block split failed");
return;
}
m_voiceFilePath = QString("%1/%2_%3_%4.%5").arg(mediaPath, blk.dt, blk.basename, blk.suffix, blk.ext);
m_voiceFilePath = mediaFile;
const auto& wf = m_recorder->waveformLevels();
QString wfStr = "wf=";
@ -267,10 +261,9 @@ void InputBar::onPttReleased() {
wfStr += b64[idx];
}
wfStr += QString(";dur=%1;").arg(durationMs);
GUI_INFO("InputBar: voice message ready: dt=%s stem=%s blocks=%d dur=%.1fs",
qPrintable(blk.dt), qPrintable(blk.basename), blk.numBlocks, durationSec);
emit sendVoiceMessage(blk.dt, blk.basename, blk.suffix, blk.ext,
blk.numBlocks, blk.blockSize, durationMs, wfStr);
GUI_INFO("InputBar: voice message ready: file=%s dur=%.1fs",
qPrintable(mediaFile), durationSec);
emit sendVoiceMessage(mediaFile, mediaFile, durationMs, wfStr);
}
void InputBar::onAttachClicked() {
@ -285,13 +278,15 @@ void InputBar::onAttachClicked() {
VoiceEncoder::ensureMediaDir(mediaPath);
QDir().mkpath(mediaPath);
MediaBlockInfo blk = media_split_to_blocks(srcPath, mediaPath, false);
if (blk.numBlocks <= 0) { DEBUG_WARN(DEBUG_CATEGORY_DEBUG, "InputBar::onAttachClicked: split failed"); return; }
QDateTime now = QDateTime::currentDateTime();
QString ts = now.toString("yyyyMMdd-HHmmss");
int rnd = rand() & 0xFFFF;
QString destPath = QString("%1/%2_%3_%4.%5")
.arg(mediaPath, ts).arg(fi.completeBaseName()).arg(rnd, 4, 16, QChar('0')).arg(fi.suffix().isEmpty() ? "bin" : fi.suffix());
GUI_INFO("InputBar::onAttachClicked: file=%s dt=%s blocks=%d",
qPrintable(fi.fileName()), qPrintable(blk.dt), blk.numBlocks);
emit sendFileAttachment(blk.dt, blk.basename, blk.suffix, blk.ext,
blk.numBlocks, blk.blockSize, fi.fileName());
GUI_INFO("InputBar::onAttachClicked: src=%s dest=%s",
qPrintable(srcPath), qPrintable(destPath));
emit sendFileAttachment(srcPath, destPath, fi.fileName());
}
void InputBar::resizeEvent(QResizeEvent* event) {

8
tools/chatgui/src/inputbar.h

@ -19,13 +19,9 @@ public:
signals:
void sendMessage(const QString &text);
void sendVoiceMessage(const QString& mediaDt, const QString& mediaBasename,
const QString& mediaSuffix, const QString& mediaExt,
int numBlocks, qint64 blockSize,
void sendVoiceMessage(const QString& mediaSrc, const QString& mediaDest,
int durationMs, const QString& waveformStr);
void sendFileAttachment(const QString& mediaDt, const QString& mediaBasename,
const QString& mediaSuffix, const QString& mediaExt,
int numBlocks, qint64 blockSize,
void sendFileAttachment(const QString& mediaSrc, const QString& mediaDest,
const QString& displayName);
protected:

12
tools/chatgui/src/mainwindow.cpp

@ -217,11 +217,13 @@ void MainWindow::setupSoundFromConfig() {
m_recorder->setCompressorEnabled(compEnabled != 0);
if (m_recorder->compressor()) {
audio_compressor_config_t compCfg = {0};
compCfg.max_gain_db = (float)m_db->getUiStateInt("compressor_max_gain_db", 30);
compCfg.lookback_ms = m_db->getUiStateInt("compressor_lookback_ms", 200);
compCfg.lookahead_ms = m_db->getUiStateInt("compressor_lookahead_ms", 100);
compCfg.rise_rate_per_500ms = (float)m_db->getUiStateInt("compressor_rise_rate_tenths", 20) / 10.0f;
compCfg.target_level = powf(10.0f, (float)m_db->getUiStateInt("compressor_target_level_db", -12) / 20.0f);
compCfg.sample_rate = 48000;
compCfg.channels = 1;
compCfg.max_gain_db = (float)m_db->getUiStateInt("compressor_max_gain_db", 25);
compCfg.lookback_ms = 200;
compCfg.lookahead_ms = 100;
compCfg.rise_rate_per_sec = (float)m_db->getUiStateInt("compressor_rise_rate", 10);
compCfg.target_level = 1.0f;
audio_compressor_configure(m_recorder->compressor(), &compCfg);
}
}

79
tools/chatgui/src/messagedelegate.cpp

@ -185,6 +185,16 @@ static int measureTextHeight(const QFont &font, const QString &text, int width)
return r.height();
}
static QString formatFileSize(qint64 bytes) {
if (bytes < 1024) return QString("%1B").arg(bytes);
double kb = bytes / 1024.0;
if (kb < 1024.0) return QString("%1KB").arg(kb, 0, 'f', kb < 10.0 ? 1 : 0);
double mb = kb / 1024.0;
if (mb < 1024.0) return QString("%1MB").arg(mb, 0, 'f', mb < 10.0 ? 1 : 0);
double gb = mb / 1024.0;
return QString("%1GB").arg(gb, 0, 'f', 1);
}
static int calcMaxLineWidth(const QFont &font, const QString &text, int availW) {
int maxW = 0;
int pos = 0;
@ -259,6 +269,35 @@ MessageDelegate::Layout MessageDelegate::calcLayout(
return L;
}
/* File attachment detection */
L.isFile = (contentType == "application/octet-stream");
if (L.isFile) {
QFontMetrics fm(option.font);
int th = fm.height();
int ih = 24; /* icon height */
L.bubbleWidth = 230;
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;
L.bubbleRect = QRect(bx, 0, L.bubbleWidth, bubbleH);
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);
if (!L.narrow) {
int avatarY = L.bubbleRect.bottom() - kAvatar;
if (!L.isOutgoing)
L.avatarRect = QRect(kMarginH, avatarY, kAvatar, kAvatar);
else
L.avatarRect = QRect(L.viewWidth - kMarginH - kAvatar, avatarY, kAvatar, kAvatar);
}
return L;
}
L.emojiCount = detectEmojiOnly(text);
L.emojiOnly = (L.emojiCount >= 1);
@ -758,6 +797,46 @@ void MessageDelegate::paint(QPainter *painter, const QStyleOptionViewItem &optio
return;
}
if (L.isFile) {
drawBubble(painter, L, option);
if (!L.narrow) drawAvatar(painter, L, index);
/* File icon */
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"));
QFont fileFont = option.font;
fileFont.setBold(true);
painter->setFont(fileFont);
painter->setPen(option.palette.text().color());
QString displayName = index.data(MsgFileDisplayNameRole).toString();
painter->drawText(L.textRect, Qt::AlignLeft | Qt::AlignVCenter, displayName);
QFont sizeFont = option.font;
sizeFont.setPointSize(sizeFont.pointSize() - 2);
painter->setFont(sizeFont);
painter->setPen(QColor(0x80, 0x80, 0x80));
qint64 fsz = index.data(MsgFileSizeRole).toLongLong();
QString sizeStr = formatFileSize(fsz);
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->restore();
return;
}
auto *cv = qobject_cast<const ChatView *>(option.widget);
bool singleTextSel = false;
bool multiRowSel = false;

2
tools/chatgui/src/messagedelegate.h

@ -26,6 +26,7 @@ enum MessageDataRole {
MsgVoiceNumBlocksRole = Qt::UserRole + 21,
MsgVoiceSigsRole = Qt::UserRole + 22,
MsgFileDisplayNameRole = Qt::UserRole + 23,
MsgFileSizeRole = Qt::UserRole + 24,
};
class MessageDelegate : public QStyledItemDelegate {
@ -61,6 +62,7 @@ private:
QRect emojiArea;
bool isOutgoing = false;
bool isVoice = false;
bool isFile = false;
QRect voicePlayBtnRect;
QRect voiceWaveformRect;
QRect voiceDurationRect;

39
tools/chatgui/src/messagelist.cpp

@ -153,8 +153,17 @@ static void setVoiceMessageRoles(QStandardItem* item, const QByteArray& data,
}
}
if (ct == "application/octet-stream") {
QString fname = QString::fromUtf8(data);
QString raw = QString::fromUtf8(data);
int pipePos = raw.indexOf('|');
QString fname = (pipePos > 0) ? raw.left(pipePos) : raw;
item->setData(fname, MsgFileDisplayNameRole);
if (pipePos > 0) {
QString meta = raw.mid(pipePos + 1);
QStringList parts = meta.split('|');
if (parts.size() >= 1)
item->setData(parts[0].toULongLong(), MsgFileSizeRole);
}
}
}
@ -246,10 +255,8 @@ MessageList::MessageList(DbManager* db, AudioRecorder* recorder, QWidget *parent
gui_bridge_post_uasync_fn(chat_core_submit_trampoline, req);
});
connect(m_inputBar, &InputBar::sendVoiceMessage, this, [this](const QString& dt, const QString& basename,
const QString& suffix, const QString& ext,
int numBlocks, qint64 blockSize,
int durationMs, const QString& wfStr) {
connect(m_inputBar, &InputBar::sendVoiceMessage, this, [this](const QString& mediaSrc, const QString& mediaDest,
int durationMs, const QString& wfStr) {
if (m_currentChannelId.isEmpty()) return;
QByteArray textData = wfStr.toUtf8();
@ -261,11 +268,9 @@ MessageList::MessageList(DbManager* db, AudioRecorder* recorder, QWidget *parent
snprintf(req->channel_id, sizeof(req->channel_id), "%s",
m_currentChannelId.toUtf8().constData());
snprintf(req->content_type, sizeof(req->content_type), "audio/opus");
snprintf(req->media_dt, sizeof(req->media_dt), "%s", dt.toUtf8().constData());
snprintf(req->media_basename, sizeof(req->media_basename), "%s", basename.toUtf8().constData());
snprintf(req->media_suffix, sizeof(req->media_suffix), "%s", suffix.toUtf8().constData());
snprintf(req->media_ext, sizeof(req->media_ext), "%s", ext.toUtf8().constData());
req->media_num_blocks = (uint32_t)numBlocks;
snprintf(req->media_src, sizeof(req->media_src), "%s", mediaSrc.toUtf8().constData());
snprintf(req->media_dest, sizeof(req->media_dest), "%s", mediaDest.toUtf8().constData());
req->media_copy = 0;
req->data = (uint8_t*)(req + 1);
req->data_len = (uint32_t)textData.size();
memcpy(req->data, textData.constData(), textData.size());
@ -273,10 +278,8 @@ MessageList::MessageList(DbManager* db, AudioRecorder* recorder, QWidget *parent
gui_bridge_post_uasync_fn(chat_core_submit_trampoline, req);
});
connect(m_inputBar, &InputBar::sendFileAttachment, this, [this](const QString& dt, const QString& basename,
const QString& suffix, const QString& ext,
int numBlocks, qint64 blockSize,
const QString& displayName) {
connect(m_inputBar, &InputBar::sendFileAttachment, this, [this](const QString& mediaSrc, const QString& mediaDest,
const QString& displayName) {
if (m_currentChannelId.isEmpty()) return;
QByteArray textData = displayName.toUtf8();
@ -288,11 +291,9 @@ MessageList::MessageList(DbManager* db, AudioRecorder* recorder, QWidget *parent
snprintf(req->channel_id, sizeof(req->channel_id), "%s",
m_currentChannelId.toUtf8().constData());
snprintf(req->content_type, sizeof(req->content_type), "application/octet-stream");
snprintf(req->media_dt, sizeof(req->media_dt), "%s", dt.toUtf8().constData());
snprintf(req->media_basename, sizeof(req->media_basename), "%s", basename.toUtf8().constData());
snprintf(req->media_suffix, sizeof(req->media_suffix), "%s", suffix.toUtf8().constData());
snprintf(req->media_ext, sizeof(req->media_ext), "%s", ext.toUtf8().constData());
req->media_num_blocks = (uint32_t)numBlocks;
snprintf(req->media_src, sizeof(req->media_src), "%s", mediaSrc.toUtf8().constData());
snprintf(req->media_dest, sizeof(req->media_dest), "%s", mediaDest.toUtf8().constData());
req->media_copy = (mediaSrc != mediaDest) ? 1 : 0;
req->data = (uint8_t*)(req + 1);
req->data_len = (uint32_t)textData.size();
memcpy(req->data, textData.constData(), textData.size());

18
tools/chatgui/src/settingsdialog.cpp

@ -6,6 +6,7 @@
#include "nodespage.h"
#include "soundsettingspage.h"
#include "audiodevicesettingspage.h"
#include "storagesettingspage.h"
#include "../transport/gui_bridge.h"
#include "../../lib/mem.h"
#include "../db/db_manager.h"
@ -45,6 +46,7 @@ SettingsDialog::SettingsDialog(const QString& configPath, DbManager* db, UtunNod
"QListWidget::item { padding: 10px 14px; }"
"QListWidget::item:selected { background: #3390EC; color: white; }");
m_categoryList->addItem("Profile");
m_categoryList->addItem("Storage");
m_categoryList->addItem("Network");
m_categoryList->addItem("Database");
m_categoryList->addItem("Status");
@ -84,6 +86,10 @@ SettingsDialog::SettingsDialog(const QString& configPath, DbManager* db, UtunNod
profLayout->addRow(sortLabel, m_sortCombo);
m_pages->addWidget(profilePage);
// Storage page
m_storagePage = new StorageSettingsPage(db, m_pages);
m_pages->addWidget(m_storagePage);
// Network page
m_networkPage = new NetworkSettingsPage(m_pages);
m_pages->addWidget(m_networkPage);
@ -166,11 +172,12 @@ void SettingsDialog::loadProfileFromConfig(const QString& path) {
void SettingsDialog::onCategoryChanged(int row) {
m_pages->setCurrentIndex(row);
if (row == 2) m_databasePage->refreshStatus();
if (row == 3) m_statusPage->refreshStatus();
if (row == 4) m_nodesPage->refreshNodes();
if (row == 5) m_soundPage->loadFromDb();
if (row == 6) m_audioDevicePage->loadFromDb();
if (row == 1) m_storagePage->loadFromDb();
if (row == 3) m_databasePage->refreshStatus();
if (row == 4) m_statusPage->refreshStatus();
if (row == 5) m_nodesPage->refreshNodes();
if (row == 6) m_soundPage->loadFromDb();
if (row == 7) m_audioDevicePage->loadFromDb();
}
void SettingsDialog::onSave() {
@ -227,5 +234,6 @@ void SettingsDialog::onSave() {
m_soundPage->applyAndSave();
m_audioDevicePage->applyAndSave();
m_storagePage->applyAndSave();
accept();
}

2
tools/chatgui/src/settingsdialog.h

@ -15,6 +15,7 @@ class StatusPage;
class NodesPage;
class SoundSettingsPage;
class AudioDeviceSettingsPage;
class StorageSettingsPage;
class DbManager;
class UtunNode;
class AudioRecorder;
@ -45,4 +46,5 @@ private:
NodesPage* m_nodesPage;
SoundSettingsPage* m_soundPage;
AudioDeviceSettingsPage* m_audioDevicePage;
StorageSettingsPage* m_storagePage;
};

90
tools/chatgui/src/storagesettingspage.cpp

@ -0,0 +1,90 @@
#include "storagesettingspage.h"
#include "../db/db_manager.h"
#include "../transport/gui_bridge.h"
#include "../../lib/mem.h"
extern "C" {
#include "chat/chat_core.h"
}
#include <QVBoxLayout>
#include <QFormLayout>
#include <QCheckBox>
#include <QSpinBox>
#include <QLabel>
#include <QFrame>
static void saveUiStateInt(const char* key, int val) {
QByteArray v = QByteArray::number(val);
size_t klen = strlen(key), vlen = (size_t)v.size(), total = klen + 1 + vlen + 1;
if (total > 256) return;
void* arg = u_malloc(total);
memcpy(arg, key, klen + 1);
memcpy((char*)arg + klen + 1, v.constData(), vlen + 1);
gui_bridge_post_uasync_fn(chat_core_save_ui_state_trampoline, arg);
}
StorageSettingsPage::StorageSettingsPage(DbManager* db, QWidget* parent)
: QWidget(parent)
, m_db(db)
{
auto* layout = new QVBoxLayout(this);
layout->setContentsMargins(0, 12, 0, 0);
layout->setSpacing(14);
auto* title = new QLabel("Media Storage", this);
title->setStyleSheet("font-weight: bold; font-size: 13px;");
layout->addWidget(title);
auto* form = new QFormLayout();
form->setSpacing(10);
form->setContentsMargins(0, 8, 0, 0);
m_autoloadCheck = new QCheckBox("Autoload media", this);
m_autoloadCheck->setStyleSheet("font-size: 13px;");
form->addRow(m_autoloadCheck);
m_autoloadMaxSize = new QSpinBox(this);
m_autoloadMaxSize->setRange(1, 500);
m_autoloadMaxSize->setSuffix(" MB");
m_autoloadMaxSize->setFixedWidth(110);
m_autoloadMaxSize->setStyleSheet(
"QSpinBox { border: 1px solid palette(mid); border-radius: 4px;"
" padding: 5px 8px; font-size: 13px; background: palette(base); }");
form->addRow("Autoload max size:", m_autoloadMaxSize);
m_maxStorageSize = new QSpinBox(this);
m_maxStorageSize->setRange(1, 100);
m_maxStorageSize->setSuffix(" GB");
m_maxStorageSize->setFixedWidth(110);
m_maxStorageSize->setStyleSheet(
"QSpinBox { border: 1px solid palette(mid); border-radius: 4px;"
" padding: 5px 8px; font-size: 13px; background: palette(base); }");
form->addRow("Max storage size:", m_maxStorageSize);
layout->addLayout(form);
auto* sep = new QFrame(this);
sep->setFrameShape(QFrame::HLine);
sep->setStyleSheet("QFrame { color: palette(mid); }");
layout->addWidget(sep);
layout->addStretch();
}
void StorageSettingsPage::loadFromDb() {
if (!m_db) return;
bool autoload = m_db->getUiStateInt("storage_autoload", 1);
m_autoloadCheck->setChecked(autoload);
int maxSizeMb = m_db->getUiStateInt("storage_autoload_maxsize_mb", 10);
m_autoloadMaxSize->setValue(maxSizeMb);
int maxGb = m_db->getUiStateInt("storage_maxsize_gb", 1);
m_maxStorageSize->setValue(maxGb);
}
void StorageSettingsPage::applyAndSave() {
saveUiStateInt("storage_autoload", m_autoloadCheck->isChecked() ? 1 : 0);
saveUiStateInt("storage_autoload_maxsize_mb", m_autoloadMaxSize->value());
saveUiStateInt("storage_maxsize_gb", m_maxStorageSize->value());
}

22
tools/chatgui/src/storagesettingspage.h

@ -0,0 +1,22 @@
#pragma once
#include <QWidget>
class DbManager;
class QCheckBox;
class QSpinBox;
class StorageSettingsPage : public QWidget {
Q_OBJECT
public:
explicit StorageSettingsPage(DbManager* db, QWidget* parent = nullptr);
void loadFromDb();
void applyAndSave();
private:
DbManager* m_db = nullptr;
QCheckBox* m_autoloadCheck;
QSpinBox* m_autoloadMaxSize;
QSpinBox* m_maxStorageSize;
};
Loading…
Cancel
Save