Browse Source

reality_relay: async DNS вместо getaddrinfo + AGC-компрессор по умолчанию (-12 dBFS)

proxy
evgeny 2 weeks ago
parent
commit
ba1c2d85d3
  1. 6
      src/chat/chat_setting.c
  2. 8
      src/radio/radio_audio.c
  3. 2
      src/routing_layer/topo_node_sqlite.c
  4. 143
      src/transport_layer/reality_relay.c
  5. 12
      tools/chatgui-android/app/src/main/java/com/utun/chat/data/ConfigProvider.kt
  6. 15
      tools/chatgui-android/app/src/main/java/com/utun/chat/ui/screens/SettingsScreen.kt
  7. 6
      tools/chatgui-android/libutun_lite/voice_recorder.c

6
src/chat/chat_setting.c

@ -34,9 +34,9 @@ struct chat_setting_def {
static const struct chat_setting_def g_setting_defs[] = {
{"storage_autoload", CHAT_SETTING_BOOL, 1, 0, 1},
{"opus_codec_preset", CHAT_SETTING_INT, 1, 0, 4},
{"compressor_enabled", CHAT_SETTING_BOOL, 0, 0, 1},
{"compressor_max_gain_db", CHAT_SETTING_INT, 25, 1, 100},
{"compressor_rise_rate", CHAT_SETTING_INT, 10, 1, 100},
{"compressor_enabled", CHAT_SETTING_BOOL, 1, 0, 1},
{"compressor_max_gain_db", CHAT_SETTING_INT, 15, 1, 100},
{"compressor_rise_rate", CHAT_SETTING_INT, 2, 1, 100},
{"media_download_max_peers", CHAT_SETTING_INT, 3, 1, 16},
{"storage_backfill_days", CHAT_SETTING_INT, 7, 1, 365},
{"standby_active_sec", CHAT_SETTING_INT, 2, 1, 3600},

8
src/radio/radio_audio.c

@ -270,11 +270,11 @@ int radio_audio_start(uint64_t group_id) {
accfg.block_duration_ms = RADIO_AUDIO_FRAME_MS;
accfg.lookback_ms = 200;
accfg.lookahead_ms = 0;
accfg.max_gain_db = (float)chat_setting_get_int(g_inst, "compressor_max_gain_db", 25);
accfg.rise_rate_per_sec = (float)chat_setting_get_int(g_inst, "compressor_rise_rate", 10);
accfg.target_level = 1.0f;
accfg.max_gain_db = (float)chat_setting_get_int(g_inst, "compressor_max_gain_db", 15);
accfg.rise_rate_per_sec = (float)chat_setting_get_int(g_inst, "compressor_rise_rate", 2);
accfg.target_level = 0.25f; /* -12 dBFS — как у голосовых сообщений */
audio_compressor_configure(ac, &accfg);
audio_compressor_set_enabled(ac, chat_setting_get_int(g_inst, "compressor_enabled", 0));
audio_compressor_set_enabled(ac, chat_setting_get_int(g_inst, "compressor_enabled", 1));
}
g_encoder = enc;

2
src/routing_layer/topo_node_sqlite.c

@ -124,7 +124,7 @@ int topo_node_sqlite_addrs_put(sqlite3* db, uint64_t node_id, struct TOPO_NODE*
sqlite3_stmt* is = NULL;
if (sqlite3_prepare_v2(db,
"INSERT INTO node_addresses(node_id,family,protocol,address,port,addr_type,config_type,nat_type,socket_id,options)"
"INSERT OR IGNORE INTO node_addresses(node_id,family,protocol,address,port,addr_type,config_type,nat_type,socket_id,options)"
" VALUES(?,?,?,?,?,?,?,?,?,?)", -1, &is, NULL) != SQLITE_OK) {
sqlite3_exec(db, "COMMIT", NULL, NULL, NULL);
return -1;

143
src/transport_layer/reality_relay.c

@ -9,6 +9,7 @@
#include "../lib/mem.h"
#include "../lib/debug_config.h"
#include "../lib/platform_compat.h"
#include "../lib/async_dns.h"
#include <string.h>
#include <stdlib.h>
#include <errno.h>
@ -32,6 +33,10 @@ struct reality_relay {
void *dest_sid;
int connecting; // 1 = ждём завершения connect() к dest
char host[256]; // dest hostname (для async DNS и логирования)
uint16_t port; // dest port
struct adns_query *dns_q; // pending async DNS query (NULL если нет)
uint8_t *c2d; size_t c2d_len, c2d_off; // pending client→dest
uint8_t *d2c; size_t d2c_len, d2c_off; // pending dest→client
@ -57,6 +62,7 @@ static void relay_free_cb(void *arg) {
static void relay_free(struct reality_relay *r) {
if (!r || r->closed) return;
r->closed = 1;
if (r->dns_q) { adns_cancel(r->dns_q); r->dns_q = NULL; }
if (r->client_sid) { uasync_remove_socket_t(r->ua, r->client_sock); r->client_sid = NULL; }
if (r->dest_sid) { uasync_remove_socket_t(r->ua, r->dest_sock); r->dest_sid = NULL; }
if (r->client_sock != SOCKET_INVALID) { socket_close_wrapper(r->client_sock); r->client_sock = SOCKET_INVALID; }
@ -182,6 +188,56 @@ static void relay_dest_write_cb(socket_t sock, void *arg) {
relay_flush_c2d(r);
}
// Создать неблокирующий dest-сокет и начать connect() к уже разрезолвенному адресу.
// Вызывается сразу (literal IP) или из коллбэка async DNS.
static void relay_connect_dest(struct reality_relay *r, const struct sockaddr_in *addr) {
socket_t dest_sock = socket(AF_INET, SOCK_STREAM, 0);
if (dest_sock == SOCKET_INVALID) {
DEBUG_WARN(DEBUG_CATEGORY_REALITY, "reality_relay: socket failed for %s:%u", r->host, r->port);
relay_free(r);
return;
}
socket_set_nonblocking(dest_sock);
r->dest_sock = dest_sock;
int cr = connect(dest_sock, (const struct sockaddr *)addr, sizeof(*addr));
if (cr < 0) {
int err = socket_get_error();
if (err != EINPROGRESS && err != ERR_WOULDBLOCK) {
DEBUG_WARN(DEBUG_CATEGORY_REALITY, "reality_relay connect %s:%u failed err=%d", r->host, r->port, err);
relay_free(r);
return;
}
} else {
r->connecting = 0;
int opt = 1; setsockopt(dest_sock, IPPROTO_TCP, TCP_NODELAY, (const char *)&opt, sizeof(opt));
}
r->dest_sid = uasync_add_socket_t(r->ua, dest_sock, relay_dest_read_cb, relay_dest_write_cb, NULL, "reality_dest", r);
if (!r->dest_sid) { relay_free(r); return; }
if (!r->connecting) {
if (r->c2d) relay_flush_c2d(r);
else uasync_set_socket_read(r->ua, r->client_sid, 1);
}
}
static void relay_dns_done_cb(const struct adns_result *res, void *arg) {
struct reality_relay *r = (struct reality_relay *)arg;
r->dns_q = NULL;
if (res->status == ADNS_OK && res->count > 0) {
struct sockaddr_in sin = res->addrs[0];
sin.sin_port = htons(r->port);
const uint8_t *b = (const uint8_t *)&sin.sin_addr;
DEBUG_INFO(DEBUG_CATEGORY_REALITY, "reality_relay: resolved %s -> %u.%u.%u.%u:%u",
r->host, b[0], b[1], b[2], b[3], r->port);
relay_connect_dest(r, &sin);
} else {
DEBUG_WARN(DEBUG_CATEGORY_REALITY, "reality_relay: DNS failed for %s status=%d", r->host, res->status);
relay_free(r);
}
}
int reality_relay_start(struct UASYNC *ua, socket_t client_sock,
const char *dest,
const uint8_t *initial_data, size_t initial_len) {
@ -190,74 +246,59 @@ int reality_relay_start(struct UASYNC *ua, socket_t client_sock,
return -1;
}
// парсим "host:port"
char host[256];
uint16_t port = 443;
{
size_t dl = strlen(dest);
if (dl >= sizeof(host)) { DEBUG_ERROR(DEBUG_CATEGORY_REALITY, "reality_relay_start: dest too long"); goto fail; }
memcpy(host, dest, dl + 1);
char *colon = strrchr(host, ':');
if (colon) { *colon = '\0'; port = (uint16_t)atoi(colon + 1); }
if (!host[0] || port == 0) { DEBUG_ERROR(DEBUG_CATEGORY_REALITY, "reality_relay_start: bad dest '%s'", dest); goto fail; }
}
struct addrinfo hints = {0};
hints.ai_family = AF_UNSPEC;
hints.ai_socktype = SOCK_STREAM;
char port_str[16];
snprintf(port_str, sizeof(port_str), "%u", port);
struct addrinfo *res;
if (getaddrinfo(host, port_str, &hints, &res) != 0) {
DEBUG_WARN(DEBUG_CATEGORY_REALITY, "reality_relay: getaddrinfo %s failed", host);
goto fail;
}
socket_t dest_sock = socket(res->ai_family, res->ai_socktype, res->ai_protocol);
if (dest_sock == SOCKET_INVALID) { freeaddrinfo(res); DEBUG_ERROR(DEBUG_CATEGORY_REALITY, "reality_relay: socket failed"); goto fail; }
socket_set_nonblocking(dest_sock);
struct reality_relay *r = u_calloc(1, sizeof(struct reality_relay));
if (!r) { socket_close_wrapper(dest_sock); freeaddrinfo(res); DEBUG_ERROR(DEBUG_CATEGORY_REALITY, "reality_relay: calloc failed"); goto fail; }
if (!r) { socket_close_wrapper(client_sock); DEBUG_ERROR(DEBUG_CATEGORY_REALITY, "reality_relay_start: calloc failed"); return -1; }
r->ua = ua;
r->client_sock = client_sock;
r->dest_sock = dest_sock;
r->dest_sock = SOCKET_INVALID;
r->connecting = 1;
// берём владение client_sock: снимаем прежнюю регистрацию и ставим свою
// берём владение client_sock сразу (снимаем старую регистрацию) — до парсинга,
// чтобы любая ошибка ниже не оставила сокет зарегистрированным в uasync.
uasync_remove_socket_t(ua, client_sock);
r->client_sid = uasync_add_socket_t(ua, client_sock, relay_client_read_cb, relay_client_write_cb, NULL, "reality_client", r);
if (!r->client_sid) { relay_free(r); freeaddrinfo(res); return -1; }
if (!r->client_sid) { relay_free(r); return -1; }
// не читаем от клиента, пока dest не подключён (иначе буферили бы без границ во время DNS)
uasync_set_socket_read(ua, r->client_sid, 0);
// парсим "host:port"
{
size_t dl = strlen(dest);
if (dl >= sizeof(r->host)) { DEBUG_ERROR(DEBUG_CATEGORY_REALITY, "reality_relay_start: dest too long"); relay_free(r); return -1; }
memcpy(r->host, dest, dl + 1);
char *colon = strrchr(r->host, ':');
if (colon) { *colon = '\0'; r->port = (uint16_t)atoi(colon + 1); }
else r->port = 443;
if (!r->host[0] || r->port == 0) { DEBUG_ERROR(DEBUG_CATEGORY_REALITY, "reality_relay_start: bad dest '%s'", dest); relay_free(r); return -1; }
}
if (initial_len && initial_data) {
r->c2d = u_malloc(initial_len);
if (!r->c2d) { relay_free(r); freeaddrinfo(res); return -1; }
if (!r->c2d) { relay_free(r); return -1; }
memcpy(r->c2d, initial_data, initial_len);
r->c2d_len = initial_len;
r->c2d_off = 0;
}
int cr = connect(dest_sock, res->ai_addr, res->ai_addrlen);
freeaddrinfo(res);
if (cr < 0) {
int err = socket_get_error();
if (err != EINPROGRESS && err != ERR_WOULDBLOCK) {
DEBUG_WARN(DEBUG_CATEGORY_REALITY, "reality_relay connect %s:%u failed err=%d", host, port, err);
relay_free(r); return -1;
}
} else {
r->connecting = 0;
int opt = 1; setsockopt(dest_sock, IPPROTO_TCP, TCP_NODELAY, (const char *)&opt, sizeof(opt));
// literal IPv4 — без DNS
struct sockaddr_in sin;
memset(&sin, 0, sizeof(sin));
sin.sin_family = AF_INET;
sin.sin_port = htons(r->port);
if (inet_pton(AF_INET, r->host, &sin.sin_addr) == 1) {
relay_connect_dest(r, &sin);
return 0;
}
r->dest_sid = uasync_add_socket_t(ua, dest_sock, relay_dest_read_cb, relay_dest_write_cb, NULL, "reality_dest", r);
if (!r->dest_sid) { relay_free(r); return -1; }
if (!r->connecting && r->c2d) relay_flush_c2d(r);
// асинхронный A-запрос — не блокирует event loop
r->dns_q = adns_resolve(ua, r->host, NULL, relay_dns_done_cb, r);
if (!r->dns_q) {
DEBUG_WARN(DEBUG_CATEGORY_REALITY, "reality_relay: async DNS start failed for %s", r->host);
relay_free(r);
return -1;
}
DEBUG_INFO(DEBUG_CATEGORY_REALITY, "reality_relay started: client→%s:%u initial=%zu", host, port, initial_len);
DEBUG_INFO(DEBUG_CATEGORY_REALITY, "reality_relay started: client→%s:%u initial=%zu (resolving)",
r->host, r->port, initial_len);
return 0;
fail:
socket_close_wrapper(client_sock);
return -1;
}

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

@ -124,9 +124,9 @@ class ConfigProvider(private val context: Context) {
key == "chat.storage_autoload_maxsize_mb" -> context.dataStore.data.first()[chatStorageAutoloadMaxsizeMb] ?: 10
key == "chat.storage_maxsize_gb" -> context.dataStore.data.first()[chatStorageMaxsizeGb] ?: 1
key == "chat.opus_codec_preset" -> context.dataStore.data.first()[chatOpusCodecPreset] ?: 1
key == "chat.compressor_enabled" -> context.dataStore.data.first()[chatCompressorEnabled] ?: 0
key == "chat.compressor_max_gain_db" -> context.dataStore.data.first()[chatCompressorMaxGainDb] ?: 25
key == "chat.compressor_rise_rate" -> context.dataStore.data.first()[chatCompressorRiseRate] ?: 10
key == "chat.compressor_enabled" -> context.dataStore.data.first()[chatCompressorEnabled] ?: 1
key == "chat.compressor_max_gain_db" -> context.dataStore.data.first()[chatCompressorMaxGainDb] ?: 15
key == "chat.compressor_rise_rate" -> context.dataStore.data.first()[chatCompressorRiseRate] ?: 2
key == "chat.media_download_max_peers" -> context.dataStore.data.first()[chatMediaDownloadMaxPeers] ?: 3
key == "chat.standby_active_sec" -> context.dataStore.data.first()[chatStandbyActiveSec] ?: 2
key == "chat.standby_sleep_sec" -> context.dataStore.data.first()[chatStandbySleepSec] ?: 60
@ -155,9 +155,9 @@ class ConfigProvider(private val context: Context) {
"chat.storage_autoload_maxsize_mb" -> prefs[chatStorageAutoloadMaxsizeMb] = value.toIntOrNull() ?: 10
"chat.storage_maxsize_gb" -> prefs[chatStorageMaxsizeGb] = value.toIntOrNull() ?: 1
"chat.opus_codec_preset" -> prefs[chatOpusCodecPreset] = value.toIntOrNull() ?: 1
"chat.compressor_enabled" -> prefs[chatCompressorEnabled] = value.toIntOrNull() ?: 0
"chat.compressor_max_gain_db" -> prefs[chatCompressorMaxGainDb] = value.toIntOrNull() ?: 25
"chat.compressor_rise_rate" -> prefs[chatCompressorRiseRate] = value.toIntOrNull() ?: 10
"chat.compressor_enabled" -> prefs[chatCompressorEnabled] = value.toIntOrNull() ?: 1
"chat.compressor_max_gain_db" -> prefs[chatCompressorMaxGainDb] = value.toIntOrNull() ?: 15
"chat.compressor_rise_rate" -> prefs[chatCompressorRiseRate] = value.toIntOrNull() ?: 2
"chat.media_download_max_peers" -> prefs[chatMediaDownloadMaxPeers] = value.toIntOrNull() ?: 3
"chat.standby_active_sec" -> prefs[chatStandbyActiveSec] = value.toIntOrNull() ?: 2
"chat.standby_sleep_sec" -> prefs[chatStandbySleepSec] = value.toIntOrNull() ?: 60

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

@ -552,26 +552,26 @@ private fun AudioSettingsTab(vm: ChatViewModel) {
/* Compressor params with defaults from desktop */
var compressorOn by remember { mutableStateOf(vm.getVoiceCompressor()) }
var maxGainDb by remember { mutableStateOf(30) }
var maxGainDb by remember { mutableStateOf(15) }
var lookbackMs by remember { mutableStateOf(200) }
var lookaheadMs by remember { mutableStateOf(100) }
var riseRateTenths by remember { mutableStateOf(20) }
var riseRate by remember { mutableStateOf(2) }
var targetLevelDb by remember { mutableStateOf(-12) }
val apply = {
vm.setVoicePreset(preset)
vm.setVoiceCompressor(compressorOn)
vm.setCompressorConfig(maxGainDb, lookbackMs, lookaheadMs, riseRateTenths / 10f, targetLevelDb.toFloat())
vm.setCompressorConfig(maxGainDb, lookbackMs, lookaheadMs, riseRate.toFloat(), targetLevelDb.toFloat())
scope.launch {
provider.setValue("chat.opus_codec_preset", preset.toString())
provider.setValue("chat.compressor_enabled", if (compressorOn) "1" else "0")
provider.setValue("chat.compressor_max_gain_db", maxGainDb.toString())
provider.setValue("chat.compressor_rise_rate", riseRateTenths.toString())
provider.setValue("chat.compressor_rise_rate", riseRate.toString())
}
NativeLib.setChatSetting("opus_codec_preset", preset.toString())
NativeLib.setChatSetting("compressor_enabled", if (compressorOn) "1" else "0")
NativeLib.setChatSetting("compressor_max_gain_db", maxGainDb.toString())
NativeLib.setChatSetting("compressor_rise_rate", riseRateTenths.toString())
NativeLib.setChatSetting("compressor_rise_rate", riseRate.toString())
}
Column(
@ -679,9 +679,8 @@ private fun AudioSettingsTab(vm: ChatViewModel) {
{ lookbackMs = it.toInt(); apply() })
LabeledSlider("Lookahead time:", lookaheadMs.toFloat(), 20f..500f, 48, "${lookaheadMs} ms",
{ lookaheadMs = it.toInt(); apply() })
LabeledSlider("Rise rate (x/500ms):", riseRateTenths.toFloat(), 11f..100f, 89,
String.format("%.1fx", riseRateTenths / 10f),
{ riseRateTenths = it.toInt(); apply() })
LabeledSlider("Rise rate (x/sec):", riseRate.toFloat(), 2f..10f, 8, "${riseRate}x",
{ riseRate = it.toInt(); apply() })
LabeledSlider("Target level:", targetLevelDb.toFloat(), -40f..0f, 40, "${targetLevelDb} dBFS",
{ targetLevelDb = it.toInt(); apply() })
}

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

@ -82,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 = 25.0f;
cfg.rise_rate_per_sec = 10.0f;
cfg.target_level = 1.0f;
cfg.max_gain_db = 15.0f;
cfg.rise_rate_per_sec = 2.0f;
cfg.target_level = 0.25f; /* -12 dBFS — едино с рацией и UI */
audio_compressor_configure(g_rec->compressor, &cfg);
}

Loading…
Cancel
Save