You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
331 lines
17 KiB
331 lines
17 KiB
#include "voice_audio_io.h" |
|
extern "C" { |
|
#include "call/call_audio.h" |
|
#include "radio/radio_audio.h" |
|
#include "speex_aec.h" |
|
#include "debug_config.h" |
|
} |
|
#include <algorithm> |
|
#include <array> |
|
#include <chrono> |
|
#include <cmath> |
|
#include <cstring> |
|
#include <system_error> |
|
|
|
static int64_t voice_time_us() { |
|
return std::chrono::duration_cast<std::chrono::microseconds>(std::chrono::steady_clock::now().time_since_epoch()).count(); |
|
} |
|
|
|
VoiceAudioIo::~VoiceAudioIo() { stop(); } |
|
|
|
bool VoiceAudioIo::start(Kind kind, uint64_t session) { |
|
stop(); |
|
m_kind = kind; m_session = session; |
|
m_captureEvents.reset(); m_references.reset(); m_playback.reset(); |
|
m_playPosition = 960; m_playEpoch = 1; |
|
m_captureDropped = m_referenceDropped = m_playbackUnderruns = 0; |
|
m_running = true; |
|
try { m_playbackThread = std::thread(&VoiceAudioIo::playbackLoop, this); } |
|
catch (const std::system_error& error) { |
|
m_running = false; |
|
DEBUG_ERROR(DEBUG_CATEGORY_CALL, "voice-io: RX worker start failed reason=%s", error.what()); return false; |
|
} |
|
DEBUG_INFO(DEBUG_CATEGORY_CALL, "voice-io: started session=%llu kind=%s PCM reserve<=40ms", |
|
(unsigned long long)session, kind == Kind::Call ? "call" : "radio"); |
|
return true; |
|
} |
|
|
|
void VoiceAudioIo::stop() { |
|
stopCapture(); |
|
m_running = false; m_wake.notify_all(); |
|
if (m_playbackThread.joinable()) m_playbackThread.join(); |
|
} |
|
|
|
bool VoiceAudioIo::startCapture(int channels, bool aecEnabled) { |
|
stopCapture(); |
|
if (!m_running || channels < 1 || channels > 2 || (m_kind == Kind::Call && channels != 1)) { |
|
DEBUG_ERROR(DEBUG_CATEGORY_CALL, "voice-io: invalid capture channels=%d", channels); return false; |
|
} |
|
m_captureEvents.reset(); m_channels = channels; m_captureSequence = 0; |
|
m_aecRequested = aecEnabled; m_capturing = true; |
|
try { m_captureThread = std::thread(&VoiceAudioIo::captureLoop, this); } |
|
catch (const std::system_error& error) { |
|
m_capturing = false; |
|
DEBUG_ERROR(DEBUG_CATEGORY_CALL, "voice-io: capture worker start failed reason=%s", error.what()); return false; |
|
} |
|
return true; |
|
} |
|
|
|
void VoiceAudioIo::stopCapture() { |
|
if (m_captureThread.joinable()) { |
|
if (m_capturing) command(EventKind::Stop); |
|
m_captureThread.join(); |
|
} |
|
m_capturing = false; |
|
} |
|
|
|
/* GUI-команда не теряется при заполнении PCM: зарезервированы слоты; отказ завершает capture. */ |
|
void VoiceAudioIo::command(EventKind kind, bool value) { |
|
if (!m_capturing) return; |
|
Event event{}; event.kind = kind; event.value = value; |
|
for (int attempt = 0; attempt < 100; ++attempt) { |
|
if (m_captureEvents.push(event, true)) { m_wake.notify_all(); return; } |
|
m_wake.notify_all(); std::this_thread::sleep_for(std::chrono::milliseconds(1)); |
|
} |
|
DEBUG_ERROR(DEBUG_CATEGORY_CALL, "voice-io: control queue full session=%llu command=%d; stopping capture", |
|
(unsigned long long)m_session, (int)kind); |
|
m_capturing = false; m_wake.notify_all(); |
|
} |
|
void VoiceAudioIo::setPtt(bool pressed) { command(EventKind::Ptt, pressed); } |
|
void VoiceAudioIo::setMuted(bool muted) { command(EventKind::Mute, muted); } |
|
void VoiceAudioIo::setAecEnabled(bool enabled) { |
|
if (enabled == m_aecRequested) return; |
|
m_aecRequested = enabled; command(EventKind::Aec, enabled); |
|
} |
|
void VoiceAudioIo::resetPlayback() { ++m_playEpoch; m_wake.notify_all(); } |
|
|
|
void VoiceAudioIo::capture(const int16_t* pcm, unsigned frames, int64_t firstTimeUs, uint64_t route, bool missing) { |
|
if (!pcm || !m_capturing) return; |
|
if (missing) { m_captureSequence += frames; m_captureDropped.fetch_add(frames); return; } |
|
int64_t enqueued = voice_time_us(); |
|
if (!firstTimeUs) firstTimeUs = enqueued - (int64_t)frames * 1000000 / 48000; |
|
unsigned offset = 0; |
|
while (offset < frames) { |
|
Event event{}; event.frames = std::min(frames - offset, 960u); |
|
event.time = firstTimeUs + (int64_t)offset * 1000000 / 48000; |
|
event.sequence = m_captureSequence; m_captureSequence += event.frames; |
|
event.route = route; |
|
event.enqueued = enqueued; |
|
memcpy(event.pcm, pcm + offset * m_channels, event.frames * m_channels * sizeof(int16_t)); |
|
if (!m_captureEvents.push(event)) m_captureDropped.fetch_add(event.frames, std::memory_order_relaxed); |
|
offset += event.frames; |
|
} |
|
m_wake.notify_all(); |
|
} |
|
|
|
void VoiceAudioIo::reference(const float* pcm, unsigned frames, int64_t firstTimeUs) { |
|
if (!m_capturing) return; |
|
if (!firstTimeUs) firstTimeUs = voice_time_us(); |
|
for (unsigned offset = 0; offset < frames;) { |
|
Event event{}; event.frames = std::min(frames - offset, 960u); |
|
event.time = firstTimeUs + (int64_t)offset * 1000000 / 48000; |
|
event.sequence = m_playEpoch.load(); |
|
for (unsigned i = 0; i < event.frames * 2; ++i) { |
|
float sample = pcm[offset * 2 + i]; |
|
event.pcm[i] = (int16_t)std::clamp(sample * 32768.0f, -32768.0f, 32767.0f); |
|
} |
|
if (!m_references.push(event)) m_referenceDropped.fetch_add(event.frames, std::memory_order_relaxed); |
|
offset += event.frames; |
|
} |
|
m_wake.notify_all(); |
|
} |
|
|
|
/* Выход всегда обслуживается независимо от capture/VAD. Уже потреблённый PCM становится референсом у владельца микшера. */ |
|
void VoiceAudioIo::mix(float* stereo, unsigned frames) { |
|
uint64_t epoch = m_playEpoch.load(); |
|
for (unsigned i = 0; i < frames; ++i) { |
|
if (m_playCurrent.epoch != epoch) m_playPosition = 960; |
|
if (m_playPosition == 960) { |
|
bool available = false; |
|
for (int attempt = 0; attempt < 12 && m_playback.pop(m_playCurrent); ++attempt) { |
|
if (m_playCurrent.epoch == epoch) { available = true; m_playPosition = 0; break; } |
|
} |
|
if (!available) { m_playbackUnderruns.fetch_add(frames - i, std::memory_order_relaxed); break; } |
|
} |
|
stereo[2*i] += m_playCurrent.pcm[2*m_playPosition] / 32768.0f; |
|
stereo[2*i+1] += m_playCurrent.pcm[2*m_playPosition+1] / 32768.0f; |
|
++m_playPosition; |
|
} |
|
m_wake.notify_all(); |
|
} |
|
|
|
void VoiceAudioIo::playbackLoop() { |
|
uint64_t epoch = m_playEpoch.load(); |
|
while (m_running) { |
|
uint64_t requested = m_playEpoch.load(); |
|
if (requested != epoch) { |
|
if (m_kind == Kind::Call) call_audio_reset_io(m_session, 1); |
|
else radio_audio_reset_playback(m_session); |
|
epoch = requested; |
|
} |
|
if (m_playback.size() < 2) { |
|
Playback block{}; block.epoch = epoch; |
|
if (m_kind == Kind::Call) { |
|
int16_t mono[960]{}; |
|
int count = std::clamp(call_audio_pull_pcm(m_session, mono, 960), 0, 960); |
|
for (int i = 0; i < count; ++i) block.pcm[2*i] = block.pcm[2*i+1] = mono[i]; |
|
} else radio_audio_pull_pcm(m_session, block.pcm, 1920); |
|
if (!m_playback.push(block)) DEBUG_WARN(DEBUG_CATEGORY_CALL, "voice-io: playback queue contention"); |
|
} else { |
|
std::unique_lock lock(m_waitMutex); |
|
m_wake.wait_for(lock, std::chrono::milliseconds(5)); |
|
} |
|
} |
|
} |
|
|
|
/* Референс хранится как PCM + время, а не как фиксированный двухкадровый FIFO независимых часов. */ |
|
void VoiceAudioIo::captureLoop() { |
|
std::array<Event, 96> history{}; |
|
size_t historyHead = 0, historyCount = 0; |
|
uint64_t referenceEpoch = m_playEpoch.load(), referenceDrops = m_referenceDropped.load(); |
|
speex_aec_t* aec = nullptr; |
|
auto configureAec = [&](bool enabled) { |
|
speex_aec_destroy(aec); aec = nullptr; |
|
if (enabled) aec = speex_aec_create_mc(48000, 960, 14400, 1, m_channels, 2); |
|
if (enabled && !aec) DEBUG_ERROR(DEBUG_CATEGORY_AEC, "voice-io: AEC unavailable session=%llu", (unsigned long long)m_session); |
|
DEBUG_INFO(DEBUG_CATEGORY_AEC, "voice-io: capture AEC=%d mic=%d render=2 session=%llu", |
|
aec != nullptr, m_channels, (unsigned long long)m_session); |
|
}; |
|
configureAec(m_aecRequested); |
|
struct Marker { unsigned position; EventKind kind; bool value; }; |
|
std::array<Marker, 64> markers{}; |
|
size_t markerCount = 0; |
|
int16_t raw[1920]{}, cleaned[1920]{}, referencePcm[1920]{}; |
|
uint8_t muteMask[960]{}; |
|
unsigned filled = 0; |
|
int64_t frameTime = 0, lastReport = voice_time_us(); |
|
uint64_t expectedSequence = 0, missingReference = 0, reportedReferenceDrops = referenceDrops, captureRoute = 0; |
|
int64_t expectedCaptureTime = 0; |
|
bool referenceAvailable = false; |
|
bool muted = false; |
|
auto applyMarker = [&](EventKind kind, bool value) { |
|
if (kind == EventKind::Mute) muted = value; |
|
else if (kind == EventKind::Ptt && m_kind == Kind::Radio) { |
|
if (value) radio_audio_talk_begin(m_session); else radio_audio_talk_end(m_session); |
|
} |
|
}; |
|
auto consumeReference = [&]() { |
|
uint64_t epoch = m_playEpoch.load(), drops = m_referenceDropped.load(); |
|
if (epoch != referenceEpoch || drops != referenceDrops) { |
|
DEBUG_DEBUG(DEBUG_CATEGORY_AEC, "voice-io: reference reset epoch=%llu->%llu dropped=%llu", |
|
(unsigned long long)referenceEpoch, (unsigned long long)epoch, (unsigned long long)(drops - referenceDrops)); |
|
historyHead = historyCount = 0; |
|
if (aec) speex_aec_reset(aec); |
|
referenceEpoch = epoch; referenceDrops = drops; |
|
} |
|
Event event{}; |
|
for (unsigned drained = 0; drained < 32 && m_references.pop(event); ++drained) { |
|
if (event.sequence != referenceEpoch) continue; |
|
if (historyCount && event.time <= history[(historyHead + historyCount - 1) % history.size()].time) { |
|
DEBUG_WARN(DEBUG_CATEGORY_AEC, "voice-io: non-monotonic render clock, resetting AEC"); |
|
historyHead = historyCount = 0; |
|
if (aec) speex_aec_reset(aec); |
|
} |
|
if (historyCount == history.size()) { historyHead = (historyHead + 1) % history.size(); --historyCount; } |
|
history[(historyHead + historyCount++) % history.size()] = event; |
|
} |
|
}; |
|
auto alignReference = [&]() { |
|
if (!historyCount) return false; |
|
size_t index = 0; |
|
for (unsigned i = 0; i < 960; ++i) { |
|
double time = frameTime + (double)i * 1000000.0 / 48000.0; |
|
while (index + 1 < historyCount && history[(historyHead + index + 1) % history.size()].time <= time) ++index; |
|
auto& block = history[(historyHead + index) % history.size()]; |
|
double interval = 1000000.0 / 48000.0; |
|
if (index + 1 < historyCount) { |
|
auto& next = history[(historyHead + index + 1) % history.size()]; |
|
double measured = (double)(next.time - block.time) / block.frames; |
|
if (measured > interval * 0.95 && measured < interval * 1.05) interval = measured; |
|
} |
|
double position = (time - block.time) / interval; |
|
if (position < 0 || position >= block.frames) return false; |
|
unsigned sample = (unsigned)position; |
|
double fraction = position - sample; |
|
for (unsigned channel = 0; channel < 2; ++channel) { |
|
int first = block.pcm[2*sample+channel]; |
|
int second = sample + 1 < block.frames ? block.pcm[2*(sample+1)+channel] : first; |
|
referencePcm[2*i+channel] = (int16_t)(first + (second - first) * fraction); |
|
} |
|
} |
|
return true; |
|
}; |
|
auto deliver = [&](unsigned frames) { |
|
unsigned begin = 0; |
|
for (size_t mark = 0; mark <= markerCount; ++mark) { |
|
unsigned end = mark < markerCount ? markers[mark].position : frames; |
|
if (m_kind == Kind::Call) memset(muteMask + begin, muted, end - begin); |
|
else if (end > begin) radio_audio_feed_pcm(m_session, cleaned + begin * m_channels, (end - begin) * m_channels); |
|
if (mark < markerCount) applyMarker(markers[mark].kind, markers[mark].value); |
|
begin = end; |
|
} |
|
if (m_kind == Kind::Call && frames == 960) call_audio_feed_prepared_pcm(m_session, cleaned, muteMask); |
|
filled = 0; markerCount = 0; |
|
}; |
|
bool stopping = false; |
|
while (m_capturing && !stopping) { |
|
consumeReference(); |
|
Event event{}; |
|
if (!m_captureEvents.pop(event)) { |
|
std::unique_lock lock(m_waitMutex); m_wake.wait_for(lock, std::chrono::milliseconds(5)); continue; |
|
} |
|
if (event.kind == EventKind::Stop) { stopping = true; break; } |
|
if (event.kind == EventKind::Aec) { configureAec(event.value); continue; } |
|
if (event.kind != EventKind::Pcm) { |
|
if (!filled) applyMarker(event.kind, event.value); |
|
else if (markerCount < markers.size()) markers[markerCount++] = {filled, event.kind, event.value}; |
|
else { DEBUG_ERROR(DEBUG_CATEGORY_CALL, "voice-io: too many controls within one capture frame"); stopping = true; } |
|
continue; |
|
} |
|
if (voice_time_us() - event.enqueued > 100000) { |
|
m_captureDropped.fetch_add(event.frames); continue; |
|
} |
|
bool clockJump = expectedCaptureTime && std::abs(event.time - expectedCaptureTime) > 20000; |
|
if (event.sequence != expectedSequence || (captureRoute && event.route != captureRoute) || clockJump) { |
|
DEBUG_WARN(DEBUG_CATEGORY_CALL, "voice-io: capture discontinuity expected=%llu actual=%llu pending=%u route=%016llx->%016llx", |
|
(unsigned long long)expectedSequence, (unsigned long long)event.sequence, filled, |
|
(unsigned long long)captureRoute, (unsigned long long)event.route); |
|
if (clockJump) DEBUG_WARN(DEBUG_CATEGORY_AEC, "voice-io: capture clock jump=%lldus", (long long)(event.time - expectedCaptureTime)); |
|
for (size_t i = 0; i < markerCount; ++i) applyMarker(markers[i].kind, markers[i].value); |
|
filled = 0; markerCount = 0; |
|
if (aec) speex_aec_reset(aec); |
|
if (m_kind == Kind::Radio) radio_audio_capture_discontinuity(m_session); |
|
} |
|
expectedSequence = event.sequence + event.frames; |
|
captureRoute = event.route; |
|
expectedCaptureTime = event.time + (int64_t)event.frames * 1000000 / 48000; |
|
for (unsigned offset = 0; offset < event.frames;) { |
|
if (!filled) frameTime = event.time + (int64_t)offset * 1000000 / 48000; |
|
unsigned take = std::min(960 - filled, event.frames - offset); |
|
memcpy(raw + filled * m_channels, event.pcm + offset * m_channels, take * m_channels * sizeof(int16_t)); |
|
filled += take; offset += take; |
|
if (filled != 960) continue; |
|
memcpy(cleaned, raw, 960 * m_channels * sizeof(int16_t)); |
|
consumeReference(); |
|
bool available = aec && alignReference(); |
|
if (available) { |
|
speex_aec_feed_playback(aec, referencePcm, 1920); |
|
speex_aec_process_capture(aec, raw, 960 * m_channels, cleaned); |
|
} else if (aec) { |
|
++missingReference; |
|
if (referenceAvailable) speex_aec_reset(aec); |
|
} |
|
referenceAvailable = available; |
|
deliver(960); |
|
} |
|
int64_t now = voice_time_us(); |
|
if (now - lastReport >= 1000000) { |
|
uint64_t captureDrops = m_captureDropped.exchange(0), refDrops = m_referenceDropped.load(); |
|
uint64_t underruns = m_playbackUnderruns.exchange(0); |
|
if (captureDrops || refDrops != reportedReferenceDrops || missingReference || underruns) |
|
DEBUG_WARN(DEBUG_CATEGORY_AEC, "voice-io: session=%llu capture_drop=%llu reference_drop=%llu missing_ref=%llu output_missing=%llu", |
|
(unsigned long long)m_session, (unsigned long long)captureDrops, (unsigned long long)refDrops, |
|
(unsigned long long)missingReference, (unsigned long long)underruns); |
|
else DEBUG_DEBUG(DEBUG_CATEGORY_AEC, "voice-io: session=%llu aligned history=%zu capture_queue=%zu", |
|
(unsigned long long)m_session, historyCount, m_captureEvents.size()); |
|
missingReference = 0; reportedReferenceDrops = refDrops; lastReport = now; |
|
} |
|
} |
|
if (filled) { |
|
DEBUG_WARN(DEBUG_CATEGORY_CALL, "voice-io: capture tail=%u frames bypasses AEC; call tail discarded, radio tail preserved", filled); |
|
memcpy(cleaned, raw, filled * m_channels * sizeof(int16_t)); |
|
deliver(filled); |
|
} |
|
uint64_t lostCapture = m_captureDropped.exchange(0), lostReference = m_referenceDropped.load() - reportedReferenceDrops; |
|
if (lostCapture || lostReference || missingReference) |
|
DEBUG_WARN(DEBUG_CATEGORY_AEC, "voice-io: capture stopped lost_capture=%llu lost_reference=%llu missing_reference=%llu", |
|
(unsigned long long)lostCapture, (unsigned long long)lostReference, (unsigned long long)missingReference); |
|
|
|
speex_aec_destroy(aec); |
|
m_capturing = false; |
|
}
|
|
|