|
|
|
|
@ -13,9 +13,12 @@
|
|
|
|
|
|
|
|
|
|
#include "radio_audio.h" |
|
|
|
|
#include "radio.h" |
|
|
|
|
#include "radio_proto.h" |
|
|
|
|
#include "../call/voice_jitter.h" |
|
|
|
|
#include "../chat/chat_setting.h" |
|
|
|
|
#include "../utun_instance.h" |
|
|
|
|
#include "../../lib/opus_codec.h" |
|
|
|
|
#include "../../lib/audio_compressor.h" |
|
|
|
|
#include "../../lib/debug_config.h" |
|
|
|
|
#include "../../lib/mem.h" |
|
|
|
|
#include "../../lib/u_async.h" |
|
|
|
|
@ -30,6 +33,21 @@
|
|
|
|
|
#define RADIO_AUDIO_MAX_FRAMES 50 /* ~1000мс @20мс */ |
|
|
|
|
#define RADIO_AUDIO_MIX_SAMPLES 960 /* максимум сэмплов за pull */ |
|
|
|
|
#define RADIO_AUDIO_SOURCE_TIMEOUT_TB 20000 /* 2с без кадров → источник завершён */ |
|
|
|
|
#define RADIO_AUDIO_REORDER_MAX 8 /* кадров в окне переупорядочивания (~160мс) */ |
|
|
|
|
#define RADIO_AUDIO_REORDER_TO_TB 1000 /* 100мс: force-flush pending, если gap не закрылся */ |
|
|
|
|
#define RADIO_AUDIO_STATS_TB 10000 /* 1с: период лога статистики буфера */ |
|
|
|
|
|
|
|
|
|
/* Параметры джиттер-буфера рации (адаптивный, симметричный темп + pre-roll). */ |
|
|
|
|
#define RADIO_AUDIO_TARGET_MS 80.0 /* целевая глубина */ |
|
|
|
|
#define RADIO_AUDIO_MIN_TEMPO 0.75 /* темп при пустом буфере (замедление) */ |
|
|
|
|
#define RADIO_AUDIO_MAX_TEMPO 1.6 /* темп при глубоком буфере (догон) */ |
|
|
|
|
|
|
|
|
|
struct radio_pending { |
|
|
|
|
uint8_t used; |
|
|
|
|
uint16_t seq; |
|
|
|
|
int len; |
|
|
|
|
uint8_t opus[RADIO_MAX_OPUS]; |
|
|
|
|
}; |
|
|
|
|
|
|
|
|
|
struct radio_source { |
|
|
|
|
uint64_t src_node_id; |
|
|
|
|
@ -39,6 +57,12 @@ struct radio_source {
|
|
|
|
|
uint64_t last_frame_tb; |
|
|
|
|
opus_codec_decoder_t* dec; /* decode-аргумент для vj (stateless per call) */ |
|
|
|
|
struct vj* vj; /* jitter-буфер с time-stretch */ |
|
|
|
|
/* reorder по seq (mesh: кадры могут прийти не по порядку через разных соседей) */ |
|
|
|
|
uint16_t next_seq; /* ожидаемый следующий seq */ |
|
|
|
|
uint8_t have_seq; |
|
|
|
|
uint64_t reorder_age_tb; /* момент появления первого gap (для таймаута) */ |
|
|
|
|
int reorder_count; |
|
|
|
|
struct radio_pending pending[RADIO_AUDIO_REORDER_MAX]; |
|
|
|
|
}; |
|
|
|
|
|
|
|
|
|
static pthread_mutex_t g_mtx = PTHREAD_MUTEX_INITIALIZER; |
|
|
|
|
@ -46,8 +70,10 @@ static struct UTUN_INSTANCE* g_inst = NULL;
|
|
|
|
|
static int g_active = 0; |
|
|
|
|
static uint64_t g_group_id = 0; |
|
|
|
|
static opus_codec_encoder_t* g_encoder = NULL; |
|
|
|
|
static struct audio_compressor* g_compressor = NULL; |
|
|
|
|
static struct radio_source g_sources[RADIO_AUDIO_MAX_SOURCES]; |
|
|
|
|
static int32_t g_mix[RADIO_AUDIO_MIX_SAMPLES]; |
|
|
|
|
static uint64_t g_last_stats_tb = 0; |
|
|
|
|
|
|
|
|
|
/* vj decode-коллбэк: дёргается только из аудио-потока (vj_pull). */ |
|
|
|
|
static int radio_audio_decode_cb(void* arg, const uint8_t* enc, int len, int16_t* pcm, int max_samples) { |
|
|
|
|
@ -55,6 +81,106 @@ static int radio_audio_decode_cb(void* arg, const uint8_t* enc, int len, int16_t
|
|
|
|
|
return opus_codec_decode(dec, enc, len, pcm, max_samples); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* ── джиттер-буфер рации (адаптивный, симметричный темп + pre-roll) ── */ |
|
|
|
|
|
|
|
|
|
static struct vj* radio_vj_create(opus_codec_decoder_t* dec) { |
|
|
|
|
struct vj_cfg cfg; |
|
|
|
|
cfg.sample_rate = RADIO_AUDIO_SAMPLE_RATE; |
|
|
|
|
cfg.channels = 1; |
|
|
|
|
cfg.frame_samples = RADIO_AUDIO_FRAME_SAMPLES; |
|
|
|
|
cfg.frame_ms = RADIO_AUDIO_FRAME_MS; |
|
|
|
|
cfg.max_frames = RADIO_AUDIO_MAX_FRAMES; |
|
|
|
|
cfg.target_ms = RADIO_AUDIO_TARGET_MS; |
|
|
|
|
cfg.min_tempo = RADIO_AUDIO_MIN_TEMPO; |
|
|
|
|
cfg.max_tempo = RADIO_AUDIO_MAX_TEMPO; |
|
|
|
|
cfg.pre_roll = 1; |
|
|
|
|
cfg.decode = radio_audio_decode_cb; |
|
|
|
|
cfg.decode_arg = dec; |
|
|
|
|
return vj_create_cfg(&cfg); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* ── reorder по seq (mesh: кадры одного источника могут прийти не по порядку) ── */ |
|
|
|
|
|
|
|
|
|
static void radio_src_reorder_reset(struct radio_source* s) { |
|
|
|
|
s->next_seq = 0; |
|
|
|
|
s->have_seq = 0; |
|
|
|
|
s->reorder_age_tb = 0; |
|
|
|
|
s->reorder_count = 0; |
|
|
|
|
memset(s->pending, 0, sizeof(s->pending)); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* force-flush: выдать все pending в порядке seq и продвинуть next_seq. */ |
|
|
|
|
static void radio_src_flush_pending(struct radio_source* s) { |
|
|
|
|
while (s->reorder_count > 0) { |
|
|
|
|
int mi = -1; |
|
|
|
|
for (int i = 0; i < RADIO_AUDIO_REORDER_MAX; i++) { |
|
|
|
|
if (!s->pending[i].used) continue; |
|
|
|
|
if (mi < 0 || (int16_t)(s->pending[i].seq - s->pending[mi].seq) < 0) mi = i; |
|
|
|
|
} |
|
|
|
|
vj_push(s->vj, s->pending[mi].opus, s->pending[mi].len); |
|
|
|
|
s->next_seq = (uint16_t)(s->pending[mi].seq + 1); |
|
|
|
|
s->pending[mi].used = 0; |
|
|
|
|
s->reorder_count--; |
|
|
|
|
} |
|
|
|
|
s->reorder_age_tb = 0; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* Выдать pending, ставшие contiguous (seq == next_seq). */ |
|
|
|
|
static void radio_src_flush_contiguous(struct radio_source* s) { |
|
|
|
|
int progressed = 1; |
|
|
|
|
while (progressed) { |
|
|
|
|
progressed = 0; |
|
|
|
|
for (int i = 0; i < RADIO_AUDIO_REORDER_MAX; i++) { |
|
|
|
|
if (!s->pending[i].used || s->pending[i].seq != s->next_seq) continue; |
|
|
|
|
vj_push(s->vj, s->pending[i].opus, s->pending[i].len); |
|
|
|
|
s->pending[i].used = 0; |
|
|
|
|
s->reorder_count--; |
|
|
|
|
s->next_seq++; |
|
|
|
|
progressed = 1; |
|
|
|
|
break; |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
if (s->reorder_count == 0) s->reorder_age_tb = 0; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* Положить кадр в источник с учётом seq: поздние дропаем, gap — в pending. */ |
|
|
|
|
static void radio_src_push_frame(struct radio_source* s, uint16_t seq, |
|
|
|
|
const uint8_t* opus, int len, uint64_t now) { |
|
|
|
|
if (!s->have_seq) { s->next_seq = seq; s->have_seq = 1; } |
|
|
|
|
|
|
|
|
|
int16_t rel = (int16_t)(seq - s->next_seq); |
|
|
|
|
if (rel < 0) { |
|
|
|
|
DEBUG_DEBUG(DEBUG_CATEGORY_RADIO, "%s: late frame drop src=%016llx seq=%u next=%u", |
|
|
|
|
RADIO_AUDIO_ID, (unsigned long long)s->src_node_id, seq, s->next_seq); |
|
|
|
|
return; |
|
|
|
|
} |
|
|
|
|
if (rel > 0) { |
|
|
|
|
if (s->reorder_count >= RADIO_AUDIO_REORDER_MAX) { |
|
|
|
|
DEBUG_WARN(DEBUG_CATEGORY_RADIO, "%s: reorder overflow src=%016llx seq=%u next=%u (flush)", |
|
|
|
|
RADIO_AUDIO_ID, (unsigned long long)s->src_node_id, seq, s->next_seq); |
|
|
|
|
radio_src_flush_pending(s); |
|
|
|
|
s->next_seq = seq; |
|
|
|
|
} else { |
|
|
|
|
for (int i = 0; i < RADIO_AUDIO_REORDER_MAX; i++) { |
|
|
|
|
if (s->pending[i].used) continue; |
|
|
|
|
memcpy(s->pending[i].opus, opus, (size_t)len); |
|
|
|
|
s->pending[i].len = len; |
|
|
|
|
s->pending[i].seq = seq; |
|
|
|
|
s->pending[i].used = 1; |
|
|
|
|
s->reorder_count++; |
|
|
|
|
if (!s->reorder_age_tb) s->reorder_age_tb = now; |
|
|
|
|
DEBUG_DEBUG(DEBUG_CATEGORY_RADIO, "%s: reorder hold src=%016llx seq=%u next=%u", |
|
|
|
|
RADIO_AUDIO_ID, (unsigned long long)s->src_node_id, seq, s->next_seq); |
|
|
|
|
break; |
|
|
|
|
} |
|
|
|
|
return; |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
vj_push(s->vj, opus, len); |
|
|
|
|
s->next_seq++; |
|
|
|
|
radio_src_flush_contiguous(s); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* ── таблица говорящих ── */ |
|
|
|
|
|
|
|
|
|
static struct radio_source* radio_src_find(uint64_t src) { |
|
|
|
|
@ -67,6 +193,7 @@ static void radio_src_free(struct radio_source* s) {
|
|
|
|
|
if (s->vj) { vj_destroy(s->vj); s->vj = NULL; } |
|
|
|
|
if (s->dec) { opus_codec_decoder_destroy(s->dec); s->dec = NULL; } |
|
|
|
|
s->used = 0; s->ending = 0; s->last_frame_tb = 0; |
|
|
|
|
radio_src_reorder_reset(s); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* Получить/создать инстанс источника. Один пир — один инстанс: если stream_id
|
|
|
|
|
@ -77,9 +204,9 @@ static struct radio_source* radio_src_acquire(uint64_t src, uint16_t stream) {
|
|
|
|
|
if (s->stream_id != stream) { |
|
|
|
|
s->stream_id = stream; |
|
|
|
|
s->ending = 0; |
|
|
|
|
radio_src_reorder_reset(s); |
|
|
|
|
vj_destroy(s->vj); |
|
|
|
|
s->vj = vj_create(RADIO_AUDIO_SAMPLE_RATE, 1, RADIO_AUDIO_FRAME_SAMPLES, RADIO_AUDIO_FRAME_MS, |
|
|
|
|
RADIO_AUDIO_MAX_FRAMES, radio_audio_decode_cb, s->dec); |
|
|
|
|
s->vj = radio_vj_create(s->dec); |
|
|
|
|
if (!s->vj) { radio_src_free(s); return NULL; } |
|
|
|
|
DEBUG_DEBUG(DEBUG_CATEGORY_RADIO, "%s: source reset src=%016llx stream=%u", RADIO_AUDIO_ID, (unsigned long long)src, stream); |
|
|
|
|
} |
|
|
|
|
@ -90,12 +217,12 @@ static struct radio_source* radio_src_acquire(uint64_t src, uint16_t stream) {
|
|
|
|
|
s = &g_sources[i]; |
|
|
|
|
s->dec = opus_codec_decoder_create(RADIO_AUDIO_SAMPLE_RATE, 1); |
|
|
|
|
if (!s->dec) { DEBUG_ERROR(DEBUG_CATEGORY_RADIO, "%s: decoder create failed", RADIO_AUDIO_ID); return NULL; } |
|
|
|
|
s->vj = vj_create(RADIO_AUDIO_SAMPLE_RATE, 1, RADIO_AUDIO_FRAME_SAMPLES, RADIO_AUDIO_FRAME_MS, |
|
|
|
|
RADIO_AUDIO_MAX_FRAMES, radio_audio_decode_cb, s->dec); |
|
|
|
|
s->vj = radio_vj_create(s->dec); |
|
|
|
|
if (!s->vj) { opus_codec_decoder_destroy(s->dec); s->dec = NULL; return NULL; } |
|
|
|
|
s->src_node_id = src; |
|
|
|
|
s->stream_id = stream; |
|
|
|
|
s->used = 1; s->ending = 0; s->last_frame_tb = 0; |
|
|
|
|
radio_src_reorder_reset(s); |
|
|
|
|
DEBUG_DEBUG(DEBUG_CATEGORY_RADIO, "%s: new source src=%016llx stream=%u", RADIO_AUDIO_ID, (unsigned long long)src, stream); |
|
|
|
|
return s; |
|
|
|
|
} |
|
|
|
|
@ -133,18 +260,38 @@ int radio_audio_start(uint64_t group_id) {
|
|
|
|
|
if (!enc) { pthread_mutex_unlock(&g_mtx); DEBUG_ERROR(DEBUG_CATEGORY_RADIO, "%s: encoder create failed", RADIO_AUDIO_ID); return -1; } |
|
|
|
|
opus_codec_encoder_bitrate_set(enc, RADIO_AUDIO_BITRATE); |
|
|
|
|
|
|
|
|
|
/* компрессор микрофона (AGC): конфиг из [chatserver], lookahead=0 — без задержки и хвоста.
|
|
|
|
|
* gain_smoothed не сбрасывается между PTT — адаптивная громкость запоминается. */ |
|
|
|
|
struct audio_compressor* ac = audio_compressor_create(); |
|
|
|
|
if (ac) { |
|
|
|
|
audio_compressor_config_t accfg = {0}; |
|
|
|
|
accfg.sample_rate = RADIO_AUDIO_SAMPLE_RATE; |
|
|
|
|
accfg.channels = 1; |
|
|
|
|
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; |
|
|
|
|
audio_compressor_configure(ac, &accfg); |
|
|
|
|
audio_compressor_set_enabled(ac, chat_setting_get_int(g_inst, "compressor_enabled", 0)); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
g_encoder = enc; |
|
|
|
|
g_compressor = ac; |
|
|
|
|
g_group_id = group_id; |
|
|
|
|
g_active = 1; |
|
|
|
|
memset(g_sources, 0, sizeof(g_sources)); |
|
|
|
|
pthread_mutex_unlock(&g_mtx); |
|
|
|
|
|
|
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_RADIO, "%s: started grp=%016llx", RADIO_AUDIO_ID, (unsigned long long)group_id); |
|
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_RADIO, "%s: started grp=%016llx compressor=%s", RADIO_AUDIO_ID, |
|
|
|
|
(unsigned long long)group_id, ac ? (audio_compressor_is_enabled(ac) ? "on" : "off") : "unavailable"); |
|
|
|
|
return 0; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
void radio_audio_stop(void) { |
|
|
|
|
opus_codec_encoder_t* enc = NULL; |
|
|
|
|
struct audio_compressor* ac = NULL; |
|
|
|
|
int was_active = 0; |
|
|
|
|
|
|
|
|
|
pthread_mutex_lock(&g_mtx); |
|
|
|
|
@ -153,6 +300,7 @@ void radio_audio_stop(void) {
|
|
|
|
|
g_active = 0; |
|
|
|
|
g_group_id = 0; |
|
|
|
|
enc = g_encoder; g_encoder = NULL; |
|
|
|
|
ac = g_compressor; g_compressor = NULL; |
|
|
|
|
for (int i = 0; i < RADIO_AUDIO_MAX_SOURCES; i++) |
|
|
|
|
if (g_sources[i].used) radio_src_free(&g_sources[i]); |
|
|
|
|
} |
|
|
|
|
@ -160,6 +308,7 @@ void radio_audio_stop(void) {
|
|
|
|
|
|
|
|
|
|
if (was_active) { |
|
|
|
|
if (enc) opus_codec_encoder_destroy(enc); |
|
|
|
|
if (ac) audio_compressor_destroy(ac); |
|
|
|
|
DEBUG_INFO(DEBUG_CATEGORY_RADIO, "%s: stopped", RADIO_AUDIO_ID); |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
@ -197,34 +346,59 @@ void radio_audio_talk_end(uint64_t group_id) {
|
|
|
|
|
int radio_audio_feed_pcm(uint64_t group_id, const int16_t* pcm, int count) { |
|
|
|
|
if (!pcm || count != RADIO_AUDIO_FRAME_SAMPLES) return -1; |
|
|
|
|
|
|
|
|
|
uint8_t opus[256]; |
|
|
|
|
enum { MAX_OUT_FRAMES = 4 }; |
|
|
|
|
uint8_t opus[MAX_OUT_FRAMES][256]; |
|
|
|
|
int opus_len[MAX_OUT_FRAMES]; |
|
|
|
|
int n_frames = 0; |
|
|
|
|
struct UASYNC* ua = NULL; |
|
|
|
|
struct UTUN_INSTANCE* inst = NULL; |
|
|
|
|
int len = -1; |
|
|
|
|
|
|
|
|
|
pthread_mutex_lock(&g_mtx); |
|
|
|
|
if (!g_active || group_id != g_group_id || !g_encoder) { |
|
|
|
|
pthread_mutex_unlock(&g_mtx); |
|
|
|
|
return -1; |
|
|
|
|
} |
|
|
|
|
len = opus_codec_encode(g_encoder, pcm, count, opus, (int)sizeof(opus)); |
|
|
|
|
|
|
|
|
|
/* компрессор микрофона: push → слив выходных блоков → encode каждого.
|
|
|
|
|
* lookahead=0 ⇒ ровно 1 frame-блок на вход; при disabled — passthrough. */ |
|
|
|
|
const int16_t* src = pcm; |
|
|
|
|
size_t src_count = (size_t)count; |
|
|
|
|
if (g_compressor) { |
|
|
|
|
audio_compressor_push(g_compressor, pcm, (size_t)count); |
|
|
|
|
src = audio_compressor_output(g_compressor); |
|
|
|
|
src_count = audio_compressor_output_size(g_compressor); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
size_t off = 0; |
|
|
|
|
while (off + (size_t)RADIO_AUDIO_FRAME_SAMPLES <= src_count && n_frames < MAX_OUT_FRAMES) { |
|
|
|
|
int l = opus_codec_encode(g_encoder, src + off, RADIO_AUDIO_FRAME_SAMPLES, |
|
|
|
|
opus[n_frames], (int)sizeof(opus[n_frames])); |
|
|
|
|
if (l > 0) { opus_len[n_frames] = l; n_frames++; } |
|
|
|
|
off += RADIO_AUDIO_FRAME_SAMPLES; |
|
|
|
|
} |
|
|
|
|
if (off != src_count) |
|
|
|
|
DEBUG_WARN(DEBUG_CATEGORY_RADIO, "%s: feed_pcm dropped %zu tail samples (compressor)", RADIO_AUDIO_ID, src_count - off); |
|
|
|
|
if (g_compressor) audio_compressor_clear_output(g_compressor); |
|
|
|
|
|
|
|
|
|
ua = g_inst ? g_inst->ua : NULL; |
|
|
|
|
inst = g_inst; |
|
|
|
|
pthread_mutex_unlock(&g_mtx); |
|
|
|
|
|
|
|
|
|
if (len <= 0) return -1; |
|
|
|
|
if (n_frames == 0) return -1; |
|
|
|
|
if (!ua || !inst) { |
|
|
|
|
DEBUG_ERROR(DEBUG_CATEGORY_RADIO, "%s: feed_pcm: instance not set", RADIO_AUDIO_ID); |
|
|
|
|
return -1; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
struct radio_media_arg* a = (struct radio_media_arg*)u_malloc(sizeof(*a)); |
|
|
|
|
if (!a) return -1; |
|
|
|
|
a->inst = inst; |
|
|
|
|
a->group_id = group_id; |
|
|
|
|
a->len = (uint16_t)len; |
|
|
|
|
memcpy(a->opus, opus, (size_t)len); |
|
|
|
|
uasync_post(ua, radio_talk_send_trampoline, a); |
|
|
|
|
for (int i = 0; i < n_frames; i++) { |
|
|
|
|
struct radio_media_arg* a = (struct radio_media_arg*)u_malloc(sizeof(*a)); |
|
|
|
|
if (!a) return -1; |
|
|
|
|
a->inst = inst; |
|
|
|
|
a->group_id = group_id; |
|
|
|
|
a->len = (uint16_t)opus_len[i]; |
|
|
|
|
memcpy(a->opus, opus[i], (size_t)opus_len[i]); |
|
|
|
|
uasync_post(ua, radio_talk_send_trampoline, a); |
|
|
|
|
} |
|
|
|
|
return 0; |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
@ -233,16 +407,21 @@ int radio_audio_feed_pcm(uint64_t group_id, const int16_t* pcm, int count) {
|
|
|
|
|
void radio_audio_on_frame(struct UTUN_INSTANCE* inst, uint64_t group_id, |
|
|
|
|
uint64_t src_node_id, uint16_t stream_id, uint16_t seq, |
|
|
|
|
uint8_t fin, const uint8_t* opus, int len, void* arg) { |
|
|
|
|
(void)inst; (void)seq; (void)arg; |
|
|
|
|
if (!opus || len <= 0) return; |
|
|
|
|
(void)inst; (void)arg; |
|
|
|
|
if (!fin && (!opus || len <= 0)) return; |
|
|
|
|
|
|
|
|
|
uint64_t now = get_time_tb(); |
|
|
|
|
pthread_mutex_lock(&g_mtx); |
|
|
|
|
if (g_active && group_id == g_group_id) { |
|
|
|
|
struct radio_source* s = radio_src_acquire(src_node_id, stream_id); |
|
|
|
|
if (s) { |
|
|
|
|
vj_push(s->vj, opus, len); |
|
|
|
|
s->last_frame_tb = get_time_tb(); |
|
|
|
|
if (fin) s->ending = 1; |
|
|
|
|
if (len > 0) radio_src_push_frame(s, seq, opus, len, now); |
|
|
|
|
s->last_frame_tb = now; |
|
|
|
|
if (fin) { |
|
|
|
|
/* FIN: выдать накопленное pending (не ждём закрытия gap) и дренируем. */ |
|
|
|
|
if (s->reorder_count > 0) radio_src_flush_pending(s); |
|
|
|
|
s->ending = 1; |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
} |
|
|
|
|
pthread_mutex_unlock(&g_mtx); |
|
|
|
|
@ -263,6 +442,9 @@ int radio_audio_pull_pcm(uint64_t group_id, int16_t* out, int max_samples) {
|
|
|
|
|
memset(g_mix, 0, (size_t)max_samples * sizeof(int32_t)); |
|
|
|
|
|
|
|
|
|
uint64_t now = get_time_tb(); |
|
|
|
|
int do_stats = (g_last_stats_tb == 0 || now - g_last_stats_tb >= RADIO_AUDIO_STATS_TB); |
|
|
|
|
struct { uint64_t src; uint16_t stream; int depth; int tempo; uint32_t dropped, under; int reorder; } st[RADIO_AUDIO_MAX_SOURCES]; |
|
|
|
|
int nst = 0; |
|
|
|
|
|
|
|
|
|
pthread_mutex_lock(&g_mtx); |
|
|
|
|
for (int i = 0; i < RADIO_AUDIO_MAX_SOURCES; i++) { |
|
|
|
|
@ -275,6 +457,21 @@ int radio_audio_pull_pcm(uint64_t group_id, int16_t* out, int max_samples) {
|
|
|
|
|
radio_src_free(s); |
|
|
|
|
continue; |
|
|
|
|
} |
|
|
|
|
/* reorder: gap висит дольше таймаута — force-flush, чтобы не задерживать звук. */ |
|
|
|
|
if (s->reorder_count > 0 && s->reorder_age_tb && (now - s->reorder_age_tb > RADIO_AUDIO_REORDER_TO_TB)) { |
|
|
|
|
DEBUG_WARN(DEBUG_CATEGORY_RADIO, "%s: reorder timeout src=%016llx pending=%d (flush)", |
|
|
|
|
RADIO_AUDIO_ID, (unsigned long long)s->src_node_id, s->reorder_count); |
|
|
|
|
radio_src_flush_pending(s); |
|
|
|
|
} |
|
|
|
|
/* ending: снять pre-roll (если ещё не снят) и дренировать накопленный остаток. */ |
|
|
|
|
if (s->ending) vj_end(s->vj); |
|
|
|
|
if (do_stats && nst < RADIO_AUDIO_MAX_SOURCES) { |
|
|
|
|
vj_get_stats(s->vj, &st[nst].depth, &st[nst].tempo, &st[nst].dropped, &st[nst].under); |
|
|
|
|
st[nst].src = s->src_node_id; |
|
|
|
|
st[nst].stream = s->stream_id; |
|
|
|
|
st[nst].reorder = s->reorder_count; |
|
|
|
|
nst++; |
|
|
|
|
} |
|
|
|
|
int16_t tmp[RADIO_AUDIO_MIX_SAMPLES]; |
|
|
|
|
int n = vj_pull(s->vj, tmp, max_samples); |
|
|
|
|
for (int k = 0; k < n; k++) g_mix[k] += tmp[k]; |
|
|
|
|
@ -288,6 +485,14 @@ int radio_audio_pull_pcm(uint64_t group_id, int16_t* out, int max_samples) {
|
|
|
|
|
} |
|
|
|
|
pthread_mutex_unlock(&g_mtx); |
|
|
|
|
|
|
|
|
|
if (do_stats) { |
|
|
|
|
g_last_stats_tb = now; |
|
|
|
|
for (int i = 0; i < nst; i++) |
|
|
|
|
DEBUG_DEBUG(DEBUG_CATEGORY_RADIO, "%s: stats src=%016llx stream=%u depth=%dms tempo=%u.%02u dropped=%u underruns=%u reorder=%d", |
|
|
|
|
RADIO_AUDIO_ID, (unsigned long long)st[i].src, st[i].stream, st[i].depth, |
|
|
|
|
st[i].tempo / 100, st[i].tempo % 100, (unsigned)st[i].dropped, (unsigned)st[i].under, st[i].reorder); |
|
|
|
|
} |
|
|
|
|
|
|
|
|
|
/* насыщающее суммирование (soft-clip) в int16 */ |
|
|
|
|
for (int k = 0; k < max_samples; k++) { |
|
|
|
|
int32_t v = g_mix[k]; |
|
|
|
|
|