feat(M3): multi-stream & per-channel tuning

Implements docs/roadmap.md M3: multiple concurrent streams per user (MIC +
SCREEN_AUDIO + AUX_DEVICE), independent per-stream receiver gain/mute/noise-
reduction, talk indicators, and enforced per-channel Opus configurability
(mono/stereo, bitrate, frame size, FEC/DTX, application).

Bugs fixed along the way (found while implementing, not pre-existing scope):
- Server hard-coded stream_id=1 for every announce, so a second stream from
  the same user silently overwrote the first in SessionRegistry::set_user_stream.
  Now a per-session counter (ConnSession::next_stream_id_); handle_stream_stop
  validates against announced_stream_ids_ before clearing.
- Client dropped mode/dtx/complexity/application from effective_audio even for
  the single M2 stream -- only sample_rate/bitrate_bps/frame_ms/fec were ever
  applied to OpusParams. Fixed on both the send (handle_stream_announce_result)
  and receive (sync_remote_streams) paths via a shared
  opus_params_from_audio_config() helper.
- OpusEncoder always used OPUS_APPLICATION_VOIP; added OpusParams::application
  and wired it through.
- on_playback's per-stream decode passed the wrong frame_size to opus_decode
  (total samples instead of samples-per-channel), which would have overflowed
  the decode buffer for any stereo stream.
- teardown_voice() raced when called concurrently from run_io()'s own cleanup
  and from disconnect() on a different thread -- both could see
  udp_thread_/talk_timer_thread_ as joinable() at once and race to join() the
  same std::thread (intermittent std::system_error under ctest). Fixed with a
  teardown_mu_ guard instead of carrying the flake forward.

New:
- Per-channel AudioConfig: SessionRegistry now seeds Lobby (mono/24kbps/VOIP/
  FEC+DTX) and a new "Music Room" channel (stereo/128kbps/AUDIO/no DTX);
  handle_stream_announce enforces the channel's config, clamping (not
  overriding) bitrate_bps to its ceiling.
- core/src/core/client.h/.cpp: local-stream state is now a
  std::unordered_map<int, LocalStream> keyed by vc_stream_kind, with
  request_id-correlated announce/result handling (request_id already
  round-tripped on the wire; just wasn't read before). on_capture_frame is
  kind-aware and upmixes mono capture to stereo when a stream's config calls
  for it. set_self_mute's mic_muted now only gates the MIC kind. NS is wired
  through set_remote_stream. New run_talk_timer() thread emits
  VC_EVENT_TALK_STATE from both remote and local edge detection.
- core/src/audio/audio_engine.h/.cpp: kind-keyed injection taps
  (inject_capture), stereo-to-mono downmix at the decode/mix boundary,
  RemoteStream gains recv_ns (lazy ApmProcessor) + noise_reduction_enabled
  and last_voice_ms/talking; new set_stream_noise_reduction() and
  poll_talk_transitions().
- core/src/session/session.h/.cpp: Stream now carries the full AudioConfig,
  not just sample_rate/frame_ms.
- New additive C ABI (core/include/voicecat.h): vc_audio_config +
  vc_get_stream_audio_config (effective Opus config for any stream you own or
  a peer's); vc_test_inject_capture (test-only synthetic PCM injection,
  clearly marked, mirrors AudioEngine::inject_capture).
- tests/test_m3_multistream.cpp: the M3 exit criterion through the real ABI
  (mirrors test_voice_client_abi.cpp's approach, not raw sockets) -- two
  concurrent local streams, independent gain/mute/NS control, per-channel
  config divergence via vc_get_stream_audio_config, talk indicators.

Explicitly out of scope for this pass (tracked in PROGRESS.md, not silently
dropped): VAD/PTT input gate + device enumeration; real WASAPI loopback
capture for SCREEN_AUDIO (synthetic injection only); true stereo playback
output (AudioEngine's mixer/output device stays mono -- Opus itself is fully
stereo-correct on the wire).

ctest --test-dir build/m1-dev: 11/11 green, verified across 3 consecutive
full-suite runs plus 8 standalone runs of the new test.

Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
2026-06-16 14:12:37 +02:00
parent c693cab35c
commit 867557eda1
17 changed files with 1102 additions and 133 deletions

View File

@@ -6,10 +6,19 @@
#include "audio/audio_engine.h"
#include <algorithm>
#include <chrono>
#include <cstring>
namespace voicecat::audio {
namespace {
int64_t now_ms() {
return std::chrono::duration_cast<std::chrono::milliseconds>(
std::chrono::steady_clock::now().time_since_epoch())
.count();
}
} // namespace
// ── JitterBuffer ─────────────────────────────────────────────────────────────
void JitterBuffer::push(Frame f) {
@@ -66,9 +75,7 @@ void JitterBuffer::reset() {
// ── AudioEngine ──────────────────────────────────────────────────────────────
AudioEngine::AudioEngine() {
inject_ring_.resize(kInjectCapSamples, 0);
}
AudioEngine::AudioEngine() = default;
AudioEngine::~AudioEngine() { stop(); }
@@ -135,31 +142,43 @@ void AudioEngine::stop() {
#endif
}
void AudioEngine::inject_capture(const int16_t* pcm, size_t n) {
std::lock_guard lk(inject_mu_);
size_t w = inject_write_.load(std::memory_order_relaxed);
void AudioEngine::inject_capture(int kind, const int16_t* pcm, size_t n) {
InjectTap* tap;
{
std::lock_guard lk(inject_mu_);
auto& slot = inject_taps_[kind];
if (!slot) {
slot = std::make_unique<InjectTap>();
slot->ring.resize(kInjectCapSamples, 0);
}
tap = slot.get();
}
size_t w = tap->write.load(std::memory_order_relaxed);
for (size_t i = 0; i < n; ++i)
inject_ring_[(w + i) % kInjectCapSamples] = pcm[i];
inject_write_.store(w + n, std::memory_order_release);
tap->ring[(w + i) % kInjectCapSamples] = pcm[i];
tap->write.store(w + n, std::memory_order_release);
// Fire capture_cb_ for each complete frame now available.
while (true) {
size_t r = inject_read_.load(std::memory_order_relaxed);
size_t avail = inject_write_.load(std::memory_order_acquire) - r;
size_t r = tap->read.load(std::memory_order_relaxed);
size_t avail = tap->write.load(std::memory_order_acquire) - r;
if (avail < static_cast<size_t>(frame_samples_)) break;
std::vector<int16_t> frame(frame_samples_);
for (int i = 0; i < frame_samples_; ++i)
frame[i] = inject_ring_[(r + i) % kInjectCapSamples];
inject_read_.store(r + frame_samples_, std::memory_order_release);
frame[i] = tap->ring[(r + i) % kInjectCapSamples];
tap->read.store(r + frame_samples_, std::memory_order_release);
if (capture_cb_) capture_cb_(frame.data(), frame_samples_);
if (capture_cb_) capture_cb_(kind, frame.data(), frame_samples_);
}
}
void AudioEngine::push_recv_frame(uint32_t ssrc, JitterBuffer::Frame f) {
std::lock_guard lk(streams_mu_);
streams_[ssrc].jitter.push(std::move(f));
auto& s = streams_[ssrc];
s.last_voice_ms.store(now_ms(), std::memory_order_relaxed);
s.jitter.push(std::move(f));
}
void AudioEngine::set_stream_gain(uint32_t ssrc, float gain) {
@@ -172,11 +191,37 @@ void AudioEngine::set_stream_mute(uint32_t ssrc, bool mute) {
streams_[ssrc].mute = mute;
}
void AudioEngine::set_stream_noise_reduction(uint32_t ssrc, bool enable) {
std::lock_guard lk(streams_mu_);
auto& s = streams_[ssrc];
s.noise_reduction_enabled = enable;
if (enable) {
if (!s.recv_ns) s.recv_ns = ApmProcessor::create();
} else {
s.recv_ns.reset();
}
}
void AudioEngine::remove_stream(uint32_t ssrc) {
std::lock_guard lk(streams_mu_);
streams_.erase(ssrc);
}
std::vector<std::pair<uint32_t, bool>> AudioEngine::poll_talk_transitions() {
std::vector<std::pair<uint32_t, bool>> edges;
std::lock_guard lk(streams_mu_);
int64_t now = now_ms();
for (auto& [ssrc, stream] : streams_) {
bool now_talking = (now - stream.last_voice_ms.load(std::memory_order_relaxed)) <
kTalkHangoverMs;
if (now_talking != stream.talking) {
stream.talking = now_talking;
edges.emplace_back(ssrc, now_talking);
}
}
return edges;
}
uint32_t AudioEngine::stream_packets_lost(uint32_t ssrc) const {
std::lock_guard lk(streams_mu_);
auto it = streams_.find(ssrc);
@@ -205,7 +250,10 @@ void AudioEngine::capture_data_cb(ma_device* dev, void* /*out*/,
}
void AudioEngine::on_capture(const int16_t* pcm, ma_uint32 frames) {
if (capture_cb_) capture_cb_(pcm, static_cast<int>(frames));
// The real hardware capture device is always the "primary" tap (kind 0 / MIC). A second
// concurrent local stream (e.g. SCREEN_AUDIO) is fed via inject_capture() in M3 — there is
// only one real capture device.
if (capture_cb_) capture_cb_(0, pcm, static_cast<int>(frames));
}
void AudioEngine::playback_data_cb(ma_device* dev, void* out,
@@ -226,24 +274,43 @@ void AudioEngine::on_playback(int16_t* out, ma_uint32 frames) {
for (auto& [ssrc, stream] : streams_) {
if (stream.mute || !stream.decoder.valid()) continue;
// M3: a stream's Opus channel count (mono/stereo, per-channel AudioConfig) may differ
// from the engine-wide playback channel count (always mono in M3 — see PROGRESS.md).
// Decode into a buffer sized for the *decoder's* channel count (opus_decode's
// frame_size parameter is samples-per-channel, not total samples — pass `frames`,
// not `frames * channels`), then convert at this boundary.
int dec_channels = std::max(1, stream.decoder.channels());
auto maybe_frame = stream.jitter.pop(stream.playout_ts);
std::vector<int16_t> pcm(frames * params_.channels);
std::vector<int16_t> pcm(frames * static_cast<ma_uint32>(dec_channels));
int n;
if (maybe_frame) {
n = stream.decoder.decode(
maybe_frame->payload.data(),
static_cast<int>(maybe_frame->payload.size()),
pcm.data(), static_cast<int>(pcm.size()));
pcm.data(), static_cast<int>(frames));
} else {
n = stream.decoder.decode(nullptr, 0, pcm.data(),
static_cast<int>(pcm.size()));
n = stream.decoder.decode(nullptr, 0, pcm.data(), static_cast<int>(frames));
}
if (n > 0) {
// `n` is samples-per-channel (matches the frame_samples convention used by
// OpusEncoder::encode elsewhere in the codebase).
if (stream.recv_ns)
stream.recv_ns->process_capture(pcm.data(), n,
static_cast<int>(params_.sample_rate));
float g = stream.gain;
for (int i = 0; i < n * static_cast<int>(params_.channels); ++i)
mix[i] += static_cast<int32_t>(static_cast<float>(pcm[i]) * g);
for (int i = 0; i < n; ++i) {
// Downmix decoder output to the engine's mono accumulator if needed
// (average L/R); upmix is unnecessary since the mix buffer is per-channel.
int32_t sample = (dec_channels == 2)
? (static_cast<int32_t>(pcm[i * 2]) +
static_cast<int32_t>(pcm[i * 2 + 1])) / 2
: static_cast<int32_t>(pcm[i]);
for (uint32_t c = 0; c < params_.channels; ++c)
mix[i * params_.channels + c] += static_cast<int32_t>(sample * g);
}
}
stream.playout_ts += frames;
}

View File

@@ -31,6 +31,8 @@
#include "codec/opus_codec.h"
#endif
#include "audio/apm_processor.h"
namespace voicecat::audio {
// ── JitterBuffer ─────────────────────────────────────────────────────────────
@@ -83,8 +85,12 @@ struct AudioParams {
// Owns miniaudio capture/playback, per-ssrc jitter buffers + Opus decoders, and the mixer.
class AudioEngine {
public:
// Callback type for encoded capture frames ready to be sent.
using CaptureCallback = std::function<void(const int16_t* pcm, int samples)>;
// Callback type for encoded capture frames ready to be sent. `kind` identifies which
// local stream this PCM belongs to (a vc_stream_kind value; 0 = MIC for the real capture
// device, which is always the "primary" tap). M3: multiple concurrent local streams are
// possible (e.g. MIC + SCREEN_AUDIO), each fed via its own injection tap (see
// inject_capture) since there is only one real hardware capture device.
using CaptureCallback = std::function<void(int kind, const int16_t* pcm, int samples)>;
AudioEngine();
~AudioEngine();
@@ -100,8 +106,11 @@ class AudioEngine {
bool running() const { return running_.load(std::memory_order_acquire); }
// Inject synthetic PCM directly into the capture pipeline (bypasses real device).
// Thread-safe; can be called from any thread including tests.
void inject_capture(const int16_t* pcm, size_t n);
// Thread-safe; can be called from any thread including tests. `kind` selects which local
// stream's injection tap to feed (each gets its own ring buffer); the 2-arg overload
// targets kind 0 (MIC) for source compatibility with existing callers.
void inject_capture(int kind, const int16_t* pcm, size_t n);
void inject_capture(const int16_t* pcm, size_t n) { inject_capture(0, pcm, n); }
// Called by the net thread when a decoded voice frame arrives for a remote stream.
void push_recv_frame(uint32_t ssrc, JitterBuffer::Frame f);
@@ -109,8 +118,17 @@ class AudioEngine {
// Per-stream receive-side controls (safe from any thread).
void set_stream_gain(uint32_t ssrc, float gain); // 0.02.0, default 1.0
void set_stream_mute(uint32_t ssrc, bool mute);
// Listener-chosen, local-only noise reduction on a specific remote stream (docs/voice.md
// §10) — lazily instantiates an ApmProcessor on first enable, frees it on disable.
void set_stream_noise_reduction(uint32_t ssrc, bool enable);
void remove_stream(uint32_t ssrc);
// Edge-triggered talk-state transitions since the last call (docs/voice.md §7: talk state
// is derived from recent frame arrival, no protocol message). Call from a lightweight
// poller, not the audio callback thread. Returns {ssrc, now_talking} for each stream whose
// state flipped.
std::vector<std::pair<uint32_t, bool>> poll_talk_transitions();
// Get stats for a remote stream's jitter buffer.
uint32_t stream_packets_lost(uint32_t ssrc) const;
uint32_t stream_target_depth_ms(uint32_t ssrc) const;
@@ -138,12 +156,16 @@ class AudioEngine {
CaptureCallback capture_cb_;
std::atomic<bool> running_{false};
// Inject ring: stores raw int16 PCM written by inject_capture().
// The encode thread reads from this (no real capture device needed in tests).
std::mutex inject_mu_;
std::vector<int16_t> inject_ring_; // circular, size = frame_samples_
std::atomic<size_t> inject_write_{0};
std::atomic<size_t> inject_read_{0};
// Inject ring(s): stores raw int16 PCM written by inject_capture(), one ring per local
// stream kind so e.g. MIC and SCREEN_AUDIO can each be fed independently in tests.
// The encode thread reads from these (no real capture device needed in tests).
struct InjectTap {
std::vector<int16_t> ring; // circular, size = kInjectCapSamples
std::atomic<size_t> write{0};
std::atomic<size_t> read{0};
};
std::mutex inject_mu_;
std::unordered_map<int, std::unique_ptr<InjectTap>> inject_taps_;
static constexpr size_t kInjectCapSamples = 48000 * 2; // 2 s @48 kHz mono
// Per remote stream (protected by streams_mu_).
@@ -155,11 +177,24 @@ class AudioEngine {
float gain = 1.0f;
bool mute = false;
uint32_t playout_ts = 0;
// M3: listener-chosen, local-only noise reduction (docs/voice.md §10). Lazily
// created only when enabled — bounded by how many remote streams this listener
// subscribes to, so no separate instance cap is needed.
bool noise_reduction_enabled = false;
std::unique_ptr<ApmProcessor> recv_ns;
// M3: talk-indicator edge detection (docs/voice.md §7) — updated by push_recv_frame
// (already off the real-time audio thread), polled by poll_talk_transitions().
std::atomic<int64_t> last_voice_ms{0};
bool talking = false;
};
mutable std::mutex streams_mu_;
std::unordered_map<uint32_t, RemoteStream> streams_;
int frame_samples_ = 960; // 20 ms @48 kHz
static constexpr int64_t kTalkHangoverMs = 300;
};
} // namespace voicecat::audio

View File

@@ -11,9 +11,16 @@ bool OpusEncoder::init(const OpusParams& p) {
channels_ = p.stereo ? 2 : 1;
frame_samples_ = opus_frame_samples(p);
int opus_application = OPUS_APPLICATION_VOIP;
switch (p.application) {
case OpusApplication::Audio: opus_application = OPUS_APPLICATION_AUDIO; break;
case OpusApplication::LowDelay: opus_application = OPUS_APPLICATION_RESTRICTED_LOWDELAY; break;
case OpusApplication::Voip: default: opus_application = OPUS_APPLICATION_VOIP; break;
}
int err = 0;
enc_ = opus_encoder_create(static_cast<opus_int32>(p.sample_rate), channels_,
OPUS_APPLICATION_VOIP, &err);
opus_application, &err);
if (err != OPUS_OK || !enc_) {
err_ = opus_strerror(err);
return false;

View File

@@ -16,15 +16,24 @@
namespace voicecat::codec {
// Mirrors voicecat::v1::OpusApplication (core/proto/voicecat.proto) without depending on
// generated protobuf headers from this low-level codec module.
enum class OpusApplication {
Voip = 0, // speech, optimized for low-rate intelligibility
Audio = 1, // music/screen-audio, optimized for fidelity
LowDelay = 2, // monitoring, minimal algorithmic delay
};
struct OpusParams {
uint32_t sample_rate = 48000;
uint32_t bitrate_bps = 24000;
uint32_t frame_ms = 20;
bool stereo = false;
bool fec = true;
bool dtx = false;
uint32_t complexity = 10;
uint32_t expected_packet_loss = 0; // % 0..100
uint32_t sample_rate = 48000;
uint32_t bitrate_bps = 24000;
uint32_t frame_ms = 20;
bool stereo = false;
bool fec = true;
bool dtx = false;
uint32_t complexity = 10;
uint32_t expected_packet_loss = 0; // % 0..100
OpusApplication application = OpusApplication::Voip;
};
// Returns frame_samples for a given sample_rate + frame_ms.

View File

@@ -241,6 +241,8 @@ cleanup:
}
void vc_client::teardown_voice() {
std::lock_guard teardown_lk(teardown_mu_);
udp_stop_.store(true, std::memory_order_release);
int ufd = udp_fd_.load(std::memory_order_acquire);
if (ufd != -1) {
@@ -250,10 +252,21 @@ void vc_client::teardown_voice() {
if (udp_thread_.joinable()) udp_thread_.join();
udp_ready_.store(false, std::memory_order_release);
talk_timer_stop_.store(true, std::memory_order_release);
if (talk_timer_thread_.joinable()) talk_timer_thread_.join();
audio_engine_.stop();
local_stream_active_.store(false, std::memory_order_release);
local_stream_pending_ = false;
local_encoder_.destroy();
{
std::lock_guard lk(local_streams_mu_);
for (auto& [kind, ls] : local_streams_) {
(void)kind;
ls.active.store(false, std::memory_order_release);
ls.pending = false;
ls.encoder.destroy();
}
local_streams_.clear();
pending_announce_kind_.clear();
}
media_send_crypto_.reset();
media_recv_crypto_.reset();
@@ -314,7 +327,7 @@ void vc_client::handle_envelope(const voicecat::v1::Envelope& env) {
handle_udp_binding_ack(env.udp_binding());
break;
case voicecat::v1::Envelope::kStreamAnnounceResult:
handle_stream_announce_result(env.stream_announce_result());
handle_stream_announce_result(env.request_id(), env.stream_announce_result());
break;
case voicecat::v1::Envelope::kPong:
break; // ignore keepalive responses
@@ -565,6 +578,9 @@ void vc_client::finish_udp_binding() {
udp_stop_.store(false, std::memory_order_release);
udp_ready_.store(true, std::memory_order_release);
udp_thread_ = std::thread([this] { run_udp_recv(); });
talk_timer_stop_.store(false, std::memory_order_release);
talk_timer_thread_ = std::thread([this] { run_talk_timer(); });
}
void vc_client::run_udp_recv() {
@@ -608,22 +624,67 @@ void vc_client::run_udp_recv() {
}
}
void vc_client::on_capture_frame(const int16_t* pcm, int samples) {
if (!local_stream_active_.load(std::memory_order_acquire)) return;
if (self_mic_muted_.load(std::memory_order_acquire)) return;
namespace {
int64_t client_now_ms() {
return std::chrono::duration_cast<std::chrono::milliseconds>(
std::chrono::steady_clock::now().time_since_epoch())
.count();
}
// Builds an OpusParams from a wire AudioConfig, applying the same field-by-field mapping on
// both the send (local-stream encoder) and receive (remote-stream decoder) paths — fixes the
// M2 gap where mode/dtx/complexity/application were silently dropped.
voicecat::codec::OpusParams opus_params_from_audio_config(const voicecat::v1::AudioConfig& a) {
voicecat::codec::OpusParams p;
p.sample_rate = a.sample_rate() ? a.sample_rate() : 48000;
p.bitrate_bps = a.bitrate_bps() ? a.bitrate_bps() : 24000;
p.frame_ms = a.frame_ms() ? a.frame_ms() : 20;
p.stereo = (a.mode() == voicecat::v1::MODE_STEREO);
p.fec = a.fec();
p.dtx = a.dtx();
p.complexity = a.complexity() ? a.complexity() : 10;
p.expected_packet_loss = a.expected_packet_loss();
p.application = static_cast<voicecat::codec::OpusApplication>(a.application());
return p;
}
} // namespace
void vc_client::on_capture_frame(int kind, const int16_t* pcm, int samples) {
std::lock_guard lk(local_streams_mu_);
auto it = local_streams_.find(kind);
if (it == local_streams_.end() || !it->second.active.load(std::memory_order_acquire)) return;
// "Mic muted" only gates the MIC stream — a concurrently-running SCREEN_AUDIO share keeps
// playing while the user's mic is muted (docs §M3 scope decision).
if (kind == static_cast<int>(VC_STREAM_MIC) &&
self_mic_muted_.load(std::memory_order_acquire)) return;
if (!media_send_crypto_) return;
int fd = udp_fd_.load(std::memory_order_acquire);
if (fd == -1) return;
auto& ls = it->second;
ls.last_capture_ms.store(client_now_ms(), std::memory_order_relaxed);
uint8_t opus_buf[1500];
int opus_len = local_encoder_.encode(pcm, samples, opus_buf, sizeof(opus_buf));
int opus_len;
if (ls.effective_params.stereo) {
// Capture is always mono in M3 (no stereo capture device); upmix L=R so a
// channel configured for stereo still gets a real, spec-correct stereo Opus stream.
std::vector<int16_t> stereo_pcm(static_cast<size_t>(samples) * 2);
for (int i = 0; i < samples; ++i) {
stereo_pcm[i * 2] = pcm[i];
stereo_pcm[i * 2 + 1] = pcm[i];
}
opus_len = ls.encoder.encode(stereo_pcm.data(), samples, opus_buf, sizeof(opus_buf));
} else {
opus_len = ls.encoder.encode(pcm, samples, opus_buf, sizeof(opus_buf));
}
if (opus_len <= 0) return;
voicecat::net::VoiceFrame hdr;
hdr.ssrc = local_ssrc_;
hdr.ssrc = ls.ssrc;
hdr.seq = static_cast<uint16_t>(media_send_crypto_->peek_send_counter());
hdr.timestamp = local_timestamp_;
local_timestamp_ += static_cast<uint32_t>(samples);
hdr.timestamp = ls.timestamp;
ls.timestamp += static_cast<uint32_t>(samples);
uint8_t header_bytes[voicecat::net::kVoiceHeaderSize];
voicecat::net::serialize_header(hdr, header_bytes);
@@ -653,7 +714,9 @@ void vc_client::ensure_audio_running() {
p.sample_rate = 48000;
p.channels = 1;
p.frame_ms = 20;
audio_engine_.start(p, [this](const int16_t* pcm, int samples) { on_capture_frame(pcm, samples); });
audio_engine_.start(p, [this](int kind, const int16_t* pcm, int samples) {
on_capture_frame(kind, pcm, samples);
});
}
void vc_client::sync_remote_streams(const voicecat::v1::User& user) {
@@ -671,9 +734,7 @@ void vc_client::sync_remote_streams(const voicecat::v1::User& user) {
remote_streams_[ssrc] = {user.id(), si.stream_id()};
voicecat::codec::OpusParams p;
p.sample_rate = si.audio().sample_rate() ? si.audio().sample_rate() : 48000;
p.frame_ms = si.audio().frame_ms() ? si.audio().frame_ms() : 20;
voicecat::codec::OpusParams p = opus_params_from_audio_config(si.audio());
audio_engine_.init_recv_stream(ssrc, p);
audio_engine_.set_stream_mute(ssrc, self_deafened_.load(std::memory_order_acquire));
@@ -714,16 +775,26 @@ void vc_client::sync_remote_streams(const voicecat::v1::User& user) {
vc_result vc_client::stream_start(const vc_stream_desc& desc, uint32_t* out_stream_id) {
if (state_net_.load(std::memory_order_acquire) != VC_STATE_CONNECTED) return VC_ERR_NOT_CONNECTED;
if (local_stream_active_.load(std::memory_order_acquire) || local_stream_pending_)
return VC_ERR_ALREADY;
uint32_t sid = next_local_stream_id_++;
local_stream_id_ = sid;
local_stream_pending_ = true;
if (out_stream_id) *out_stream_id = sid;
int kind = static_cast<int>(desc.kind);
uint64_t req_id;
{
std::lock_guard lk(local_streams_mu_);
auto& ls = local_streams_[kind]; // default-constructs if this kind is new
if (ls.active.load(std::memory_order_acquire) || ls.pending) return VC_ERR_ALREADY;
uint32_t sid = next_local_stream_id_++;
ls.stream_id = sid;
ls.pending = true;
if (out_stream_id) *out_stream_id = sid;
req_id = next_req_id_++;
ls.pending_request_id = req_id;
pending_announce_kind_[req_id] = kind;
}
voicecat::v1::Envelope req;
req.set_request_id(next_req_id_++);
req.set_request_id(req_id);
auto* ann = req.mutable_stream_announce();
ann->set_kind(desc.kind == VC_STREAM_SCREEN_AUDIO ? voicecat::v1::STREAM_SCREEN_AUDIO
: desc.kind == VC_STREAM_AUX_DEVICE ? voicecat::v1::STREAM_AUX_DEVICE
@@ -731,55 +802,76 @@ vc_result vc_client::stream_start(const vc_stream_desc& desc, uint32_t* out_stre
if (desc.label) ann->set_label(desc.label);
auto* audio = ann->mutable_requested_audio();
audio->set_sample_rate(48000);
audio->set_bitrate_bps(24000);
// bitrate_bps intentionally left unset (0 = "no preference"): the channel's AudioConfig
// is authoritative when the channel has one (docs/voice.md §3) — a client-side default
// here would otherwise always clamp the channel's bitrate down to this value. frame_ms/
// fec are likewise channel-controlled when a channel config exists; these are only used
// as a fallback when it doesn't.
audio->set_frame_ms(20);
audio->set_fec(true);
queue_envelope(req);
return VC_OK;
}
void vc_client::handle_stream_announce_result(const voicecat::v1::StreamAnnounceResult& msg) {
if (!local_stream_pending_) return;
local_stream_pending_ = false;
void vc_client::handle_stream_announce_result(uint64_t req_id,
const voicecat::v1::StreamAnnounceResult& msg) {
uint32_t self_uid;
uint32_t emit_stream_id;
bool ok_to_emit = false;
{
std::lock_guard lk(local_streams_mu_);
auto pit = pending_announce_kind_.find(req_id);
if (pit == pending_announce_kind_.end()) return; // stray/duplicate — ignore
int kind = pit->second;
pending_announce_kind_.erase(pit);
if (!msg.ok()) {
emit_error(VC_ERR_PROTOCOL, msg.error().c_str());
return;
}
auto sit = local_streams_.find(kind);
if (sit == local_streams_.end() || !sit->second.pending) return;
auto& ls = sit->second;
ls.pending = false;
local_ssrc_ = msg.ssrc();
local_timestamp_ = 0;
if (!msg.ok()) {
emit_error(VC_ERR_PROTOCOL, msg.error().c_str());
return;
}
voicecat::codec::OpusParams p;
const auto& eff = msg.effective_audio();
p.sample_rate = eff.sample_rate() ? eff.sample_rate() : 48000;
p.bitrate_bps = eff.bitrate_bps() ? eff.bitrate_bps() : 24000;
p.frame_ms = eff.frame_ms() ? eff.frame_ms() : 20;
p.fec = eff.fec();
local_frame_samples_ = static_cast<uint16_t>(voicecat::codec::opus_frame_samples(p));
ls.ssrc = msg.ssrc();
ls.timestamp = 0;
ls.effective_params = opus_params_from_audio_config(msg.effective_audio());
ls.frame_samples = static_cast<uint16_t>(voicecat::codec::opus_frame_samples(ls.effective_params));
if (!local_encoder_.init(p)) {
emit_error(VC_ERR_AUDIO, "opus encoder init failed");
return;
if (!ls.encoder.init(ls.effective_params)) {
emit_error(VC_ERR_AUDIO, "opus encoder init failed");
return;
}
ls.active.store(true, std::memory_order_release);
self_uid = self_user_id_;
emit_stream_id = ls.stream_id;
ok_to_emit = true;
}
ensure_audio_running();
local_stream_active_.store(true, std::memory_order_release);
vc_event ev{};
ev.type = VC_EVENT_STREAM_STARTED;
ev.user_id = self_user_id_;
ev.stream_id = local_stream_id_;
emit(ev);
if (ok_to_emit) {
vc_event ev{};
ev.type = VC_EVENT_STREAM_STARTED;
ev.user_id = self_uid;
ev.stream_id = emit_stream_id;
emit(ev);
}
}
vc_result vc_client::stream_stop(uint32_t stream_id) {
if (state_net_.load(std::memory_order_acquire) != VC_STATE_CONNECTED) return VC_ERR_NOT_CONNECTED;
if (!local_stream_active_.load(std::memory_order_acquire) || stream_id != local_stream_id_)
return VC_ERR_INVALID_ARG;
local_stream_active_.store(false, std::memory_order_release);
local_encoder_.destroy();
{
std::lock_guard lk(local_streams_mu_);
LocalStream* ls = find_local_stream_by_id(stream_id);
if (!ls || !ls->active.load(std::memory_order_acquire)) return VC_ERR_INVALID_ARG;
ls->active.store(false, std::memory_order_release);
ls->encoder.destroy();
}
voicecat::v1::Envelope req;
req.set_request_id(next_req_id_++);
@@ -794,6 +886,14 @@ vc_result vc_client::stream_stop(uint32_t stream_id) {
return VC_OK;
}
vc_client::LocalStream* vc_client::find_local_stream_by_id(uint32_t stream_id) {
for (auto& [kind, ls] : local_streams_) {
(void)kind;
if (ls.stream_id == stream_id) return &ls;
}
return nullptr;
}
vc_result vc_client::set_input_device(uint32_t, const char*) { return VC_ERR_NOT_IMPLEMENTED; }
vc_result vc_client::set_input_mode(vc_input_mode) { return VC_ERR_NOT_IMPLEMENTED; }
vc_result vc_client::set_push_to_talk(bool) { return VC_ERR_NOT_IMPLEMENTED; }
@@ -812,7 +912,7 @@ vc_result vc_client::set_self_mute(bool mic_muted, bool deafened) {
}
vc_result vc_client::set_remote_stream(uint32_t user_id, uint32_t stream_id, float gain,
bool muted, bool /*noise_reduction*/) {
bool muted, bool noise_reduction) {
if (state_net_.load(std::memory_order_acquire) != VC_STATE_CONNECTED) return VC_ERR_NOT_CONNECTED;
const auto* user = session_model_.find_user(user_id);
if (!user) return VC_ERR_INVALID_ARG;
@@ -820,18 +920,124 @@ vc_result vc_client::set_remote_stream(uint32_t user_id, uint32_t stream_id, flo
if (s.stream_id == stream_id) {
audio_engine_.set_stream_gain(s.ssrc, gain);
audio_engine_.set_stream_mute(s.ssrc, muted);
audio_engine_.set_stream_noise_reduction(s.ssrc, noise_reduction);
return VC_OK;
}
}
return VC_ERR_INVALID_ARG;
}
vc_result vc_client::get_stream_audio_config(uint32_t user_id, uint32_t stream_id,
vc_audio_config* out) {
if (state_net_.load(std::memory_order_acquire) != VC_STATE_CONNECTED) return VC_ERR_NOT_CONNECTED;
if (user_id == self_user_id_) {
std::lock_guard lk(local_streams_mu_);
LocalStream* ls = find_local_stream_by_id(stream_id);
if (!ls || !ls->active.load(std::memory_order_acquire)) return VC_ERR_INVALID_ARG;
const auto& p = ls->effective_params;
out->codec = 0;
out->mode = p.stereo ? 1u : 0u;
out->sample_rate = p.sample_rate;
out->bitrate_bps = p.bitrate_bps;
out->frame_ms = p.frame_ms;
out->application = static_cast<uint32_t>(p.application);
out->fec = p.fec ? 1 : 0;
out->expected_packet_loss = p.expected_packet_loss;
out->dtx = p.dtx ? 1 : 0;
out->complexity = p.complexity;
return VC_OK;
}
const auto* user = session_model_.find_user(user_id);
if (!user) return VC_ERR_INVALID_ARG;
for (const auto& s : user->streams) {
if (s.stream_id != stream_id) continue;
out->codec = 0;
out->mode = s.mode;
out->sample_rate = s.sample_rate;
out->bitrate_bps = s.bitrate_bps;
out->frame_ms = s.frame_ms;
out->application = s.application;
out->fec = s.fec ? 1 : 0;
out->expected_packet_loss = s.expected_packet_loss;
out->dtx = s.dtx ? 1 : 0;
out->complexity = s.complexity;
return VC_OK;
}
return VC_ERR_INVALID_ARG;
}
vc_result vc_client::test_inject_capture(uint32_t stream_id, const int16_t* pcm, size_t samples) {
int kind = -1;
{
std::lock_guard lk(local_streams_mu_);
for (auto& [k, ls] : local_streams_) {
if (ls.stream_id == stream_id && ls.active.load(std::memory_order_acquire)) {
kind = k;
break;
}
}
}
if (kind < 0) return VC_ERR_INVALID_ARG;
audio_engine_.inject_capture(kind, pcm, samples);
return VC_OK;
}
vc_result vc_client::list_devices(vc_device_kind, vc_device_list* out) {
out->items = nullptr;
out->count = 0;
return VC_ERR_NOT_IMPLEMENTED;
}
void vc_client::run_talk_timer() {
while (!talk_timer_stop_.load(std::memory_order_acquire)) {
// Remote streams: ask the engine for edge-triggered transitions, map ssrc -> (user,
// stream) via remote_streams_, emit.
for (auto& [ssrc, talking] : audio_engine_.poll_talk_transitions()) {
uint32_t uid = 0, sid = 0;
{
std::lock_guard lk(remote_streams_mu_);
auto it = remote_streams_.find(ssrc);
if (it == remote_streams_.end()) continue;
uid = it->second.first;
sid = it->second.second;
}
vc_event ev{};
ev.type = VC_EVENT_TALK_STATE;
ev.user_id = uid;
ev.stream_id = sid;
ev.u32a = talking ? 1u : 0u;
emit(ev);
}
// Local streams: same hangover logic, driven by on_capture_frame's last_capture_ms.
std::vector<vc_event> local_events;
{
std::lock_guard lk(local_streams_mu_);
int64_t now = client_now_ms();
for (auto& [kind, ls] : local_streams_) {
(void)kind;
if (!ls.active.load(std::memory_order_acquire)) continue;
bool now_talking =
(now - ls.last_capture_ms.load(std::memory_order_relaxed)) < kTalkHangoverMs;
if (now_talking != ls.talking) {
ls.talking = now_talking;
vc_event ev{};
ev.type = VC_EVENT_TALK_STATE;
ev.user_id = self_user_id_;
ev.stream_id = ls.stream_id;
ev.u32a = now_talking ? 1u : 0u;
local_events.push_back(ev);
}
}
}
for (auto& ev : local_events) emit(ev);
std::this_thread::sleep_for(std::chrono::milliseconds(kTalkPollMs));
}
}
#else // !VOICECAT_HAS_NET
// ── M0 stub implementations ───────────────────────────────────────────────────
@@ -864,5 +1070,11 @@ vc_result vc_client::list_devices(vc_device_kind, vc_device_list* out) {
out->count = 0;
return VC_ERR_NOT_IMPLEMENTED;
}
vc_result vc_client::get_stream_audio_config(uint32_t, uint32_t, vc_audio_config*) {
return VC_ERR_NOT_IMPLEMENTED;
}
vc_result vc_client::test_inject_capture(uint32_t, const int16_t*, size_t) {
return VC_ERR_NOT_IMPLEMENTED;
}
#endif // VOICECAT_HAS_NET

View File

@@ -58,6 +58,14 @@ struct vc_client {
vc_result list_devices(vc_device_kind kind, vc_device_list* out);
// M3: effective Opus config for a (user_id, stream_id) — our own pending/active local
// streams, or any peer's broadcast StreamInfo.audio.
vc_result get_stream_audio_config(uint32_t user_id, uint32_t stream_id,
vc_audio_config* out);
// TEST-ONLY (see voicecat.h) — inject synthetic PCM into a local stream's encode pipeline.
vc_result test_inject_capture(uint32_t stream_id, const int16_t* pcm, size_t samples);
vc_connection_state state() const {
#ifdef VOICECAT_HAS_NET
return state_net_.load(std::memory_order_acquire);
@@ -122,20 +130,43 @@ struct vc_client {
std::unique_ptr<voicecat::crypto::SodiumMediaCrypto> media_recv_crypto_;
voicecat::audio::AudioEngine audio_engine_;
voicecat::codec::OpusEncoder local_encoder_;
std::atomic<bool> local_stream_active_{false};
bool local_stream_pending_{false}; // announced, waiting for result
uint32_t local_stream_id_{0};
uint32_t local_ssrc_{0};
uint32_t next_local_stream_id_{1};
uint32_t local_timestamp_{0};
uint16_t local_frame_samples_{960};
// M3: one LocalStream per concurrently-active stream kind (MIC/SCREEN_AUDIO/AUX_DEVICE
// are each singletons for a given client — see docs/roadmap.md §M3). Replaces the M2
// single-stream fields (local_encoder_/local_stream_active_/etc).
struct LocalStream {
voicecat::codec::OpusEncoder encoder;
std::atomic<bool> active{false};
bool pending{false}; // announced, awaiting StreamAnnounceResult
uint32_t stream_id{0};
uint32_t ssrc{0};
uint32_t timestamp{0};
uint16_t frame_samples{960};
uint64_t pending_request_id{0};
// Full effective OpusParams from the last StreamAnnounceResult — retained so
// vc_get_stream_audio_config() has something to read back for our own streams.
voicecat::codec::OpusParams effective_params;
// Talk-indicator edge detection (docs/voice.md §7) — updated in on_capture_frame.
std::atomic<int64_t> last_capture_ms{0};
bool talking = false;
};
mutable std::mutex local_streams_mu_;
std::unordered_map<int, LocalStream> local_streams_; // keyed by vc_stream_kind
std::unordered_map<uint64_t, int> pending_announce_kind_; // request_id -> kind
uint32_t next_local_stream_id_{1};
// ssrc → (user_id, stream_id) for remote streams already wired into audio_engine_.
mutable std::mutex remote_streams_mu_;
std::unordered_map<uint32_t, std::pair<uint32_t, uint32_t>> remote_streams_;
// M3: talk-indicator polling thread (separate from udp_thread_ / the miniaudio callback
// thread — see architecture.md §3 real-time rule).
std::thread talk_timer_thread_;
std::atomic<bool> talk_timer_stop_{false};
static constexpr int64_t kTalkPollMs = 100;
static constexpr int64_t kTalkHangoverMs = 300;
// Local UDP destination (server media endpoint), resolved once during binding.
uint32_t udp_dest_addr_{0}; // network byte order
uint16_t udp_dest_port_{0}; // network byte order
@@ -143,6 +174,14 @@ struct vc_client {
std::atomic<bool> self_mic_muted_{false};
std::atomic<bool> self_deafened_{false};
// teardown_voice() is called both from run_io()'s own cleanup (on the io_thread_, when
// the read loop exits) and from disconnect() (on the caller's thread) -- without
// serializing those two call sites, both can see udp_thread_/talk_timer_thread_ as
// joinable() at the same time and race to join() the same std::thread object (UB; an
// intermittent "No such process" std::system_error on Windows is the typical symptom).
// This mutex makes teardown_voice() idempotent under concurrent calls.
std::mutex teardown_mu_;
// ── io_thread_ entry point ──────────────────────────────────────────────────
void run_io(std::string host, uint16_t port);
@@ -155,7 +194,8 @@ struct vc_client {
void handle_text_message(const voicecat::v1::TextMessage& msg);
void handle_disconnect(const voicecat::v1::Disconnect& msg);
void handle_udp_binding_ack(const voicecat::v1::UdpBinding& msg);
void handle_stream_announce_result(const voicecat::v1::StreamAnnounceResult& msg);
void handle_stream_announce_result(uint64_t req_id,
const voicecat::v1::StreamAnnounceResult& msg);
// ── M2: UDP / media helpers ──────────────────────────────────────────────────
// Kicks off TCP UdpBinding request; called once after a successful AuthResult.
@@ -164,8 +204,9 @@ struct vc_client {
void finish_udp_binding();
// udp_thread_ entry point: recv loop, AEAD-open, decode, push to audio_engine_.
void run_udp_recv();
// capture_cb passed to audio_engine_.start(): encode + seal + send one frame.
void on_capture_frame(const int16_t* pcm, int samples);
// capture_cb passed to audio_engine_.start(): encode + seal + send one frame for the
// given local stream `kind` (M3: multiple concurrent local streams are possible).
void on_capture_frame(int kind, const int16_t* pcm, int samples);
// Inspect a User proto's streams and wire up any new remote ssrc into audio_engine_,
// emitting VC_EVENT_STREAM_STARTED/STOPPED as streams appear/disappear.
void sync_remote_streams(const voicecat::v1::User& user);
@@ -174,6 +215,12 @@ struct vc_client {
// Joins udp_thread_, stops audio_engine_, clears media crypto/remote-stream state.
// Safe to call multiple times. Called both from run_io()'s cleanup and disconnect().
void teardown_voice();
// talk_timer_thread_ entry point (M3): polls audio_engine_ for remote talk-state edges
// and local capture activity, emitting VC_EVENT_TALK_STATE. Never the audio RT thread.
void run_talk_timer();
// Find a LocalStream by its client-assigned stream_id (held under local_streams_mu_ by
// the caller, or taken internally). Returns nullptr if not found/not active.
LocalStream* find_local_stream_by_id(uint32_t stream_id);
// ── Helpers (io_thread_ and caller threads) ─────────────────────────────────
// Queue an encoded envelope to be sent on io_thread_.

View File

@@ -38,6 +38,13 @@ std::vector<Stream> copy_streams(
s.label = pb.label();
s.sample_rate = pb.audio().sample_rate() ? pb.audio().sample_rate() : 48000;
s.frame_ms = pb.audio().frame_ms() ? pb.audio().frame_ms() : 20;
s.mode = static_cast<uint32_t>(pb.audio().mode());
s.bitrate_bps = pb.audio().bitrate_bps();
s.application = static_cast<uint32_t>(pb.audio().application());
s.fec = pb.audio().fec();
s.expected_packet_loss = pb.audio().expected_packet_loss();
s.dtx = pb.audio().dtx();
s.complexity = pb.audio().complexity();
out.push_back(std::move(s));
}
return out;

View File

@@ -31,8 +31,19 @@ struct Stream {
uint32_t ssrc{0};
int kind{0};
std::string label;
// Full effective AudioConfig (docs/protocol.md §5), as broadcast by the server in
// StreamInfo.audio — mirrors voicecat::v1::AudioConfig field-for-field so M3 per-channel
// tuning (mono/stereo, bitrate, FEC/DTX, application) is observable client-side, not just
// sample_rate/frame_ms.
uint32_t sample_rate{48000};
uint32_t frame_ms{20};
uint32_t mode{0}; // 0 = mono, 1 = stereo (ChannelMode)
uint32_t bitrate_bps{0};
uint32_t application{0}; // 0=VOIP, 1=AUDIO, 2=LOWDELAY (OpusApplication)
bool fec{false};
uint32_t expected_packet_loss{0};
bool dtx{false};
uint32_t complexity{0};
};
struct User {

View File

@@ -116,6 +116,18 @@ vc_result vc_set_remote_stream(vc_client* c, uint32_t user_id, uint32_t stream_i
return c->set_remote_stream(user_id, stream_id, gain, muted != 0, noise_reduction != 0);
}
vc_result vc_get_stream_audio_config(vc_client* c, uint32_t user_id, uint32_t stream_id,
vc_audio_config* out) {
if (c == nullptr || out == nullptr) return VC_ERR_INVALID_ARG;
return c->get_stream_audio_config(user_id, stream_id, out);
}
vc_result vc_test_inject_capture(vc_client* c, uint32_t stream_id, const int16_t* pcm,
size_t samples) {
if (c == nullptr || pcm == nullptr) return VC_ERR_INVALID_ARG;
return c->test_inject_capture(stream_id, pcm, samples);
}
vc_result vc_send_text(vc_client* c, vc_text_scope scope, uint32_t target_id,
const char* utf8) {
if (c == nullptr || utf8 == nullptr) return VC_ERR_INVALID_ARG;