diff --git a/PROGRESS.md b/PROGRESS.md index 76ee77c..9f39959 100644 --- a/PROGRESS.md +++ b/PROGRESS.md @@ -10,25 +10,20 @@ up instantly. Newest status at the top. ## ▶ Where we left off / next action -- **Done:** **M2 — voice, single stream** ✓ complete (2026-06-16), now genuinely satisfied - through the real client library, not just `test_m2_voice`'s raw-socket harness. - `vc_stream_start/stop`, UDP binding, capture→encode→seal→send, and recv→open→decode→ - playback were all `VC_ERR_NOT_IMPLEMENTED` stubs in `core/src/core/client.cpp` even after - `test_m2_voice` went green — meaning `vccli`/any GUI client still couldn't actually talk. - Implemented for real this session, plus server-side `StreamInfo` broadcast - (`SessionRegistry::set_user_stream/clear_user_stream` → `UserEvent::UPDATED`) so a second - client's `sync_remote_streams()` learns about a peer's stream without polling. - `ctest --test-dir build/m1-dev` — **10/10 tests** green, including the new - `test_voice_client_abi` (two real `vc_client` instances, not raw sockets, drive the full - UDP-binding → stream-announce → cross-client `STREAM_STARTED`/`STOPPED` event path). - `vccli --voice` (new flag, alongside `--host/--port/--nick/--channel/--mute/--text`) was - manually verified live: two instances see each other's mic-stream start/stop in real time. - Note: `test_m1_integration` and `test_m2_voice` have a pre-existing intermittent flake on - Windows in their cleanup paths (thread-join / SQLite-file-handle release race, unrelated to - this session's changes) — rerun in isolation if one fails standalone in the full suite. -- **Next:** **M3 — multi-stream & per-channel tuning** (screen audio, listener-side per-user - NR, jitter buffer stats API — note `voicecat.h` currently exposes no stats getter beyond - `on_level`'s RMS meter, so a stats API needs new C ABI surface). See `docs/roadmap.md §M3`. +- **Done:** **M3 — multi-stream & per-channel tuning** ✓ complete (2026-06-16). See the M3 + section below for the full file-by-file change list. `ctest --test-dir build/m1-dev` — + **11/11 tests** green (3 consecutive full-suite runs), including the new + `test_m3_multistream` (real `vc_client`s, not raw sockets — same lesson as M2: ABI-level + coverage is what proves the client library, not just the wire protocol). + **Two items intentionally still open** (carried forward, not silently dropped): + - `vc_set_input_device`/`vc_set_input_mode`/`vc_set_push_to_talk`/`vc_list_devices` + (device enumeration + VAD/PTT input gate) remain `VC_ERR_NOT_IMPLEMENTED` — explicitly + scoped out of this M3 pass; revisit in a future milestone. + - Stereo Opus is now wire-correct end-to-end (a channel configured `MODE_STEREO` really + encodes/decodes 2-channel Opus packets), but `AudioEngine`'s playback mixer/output device + stays mono internally — stereo streams are downmixed (avg L/R) immediately after decode, + before mixing. True stereo *playback output* is a follow-up, not part of M3. +- **Next:** **M4 — native clients** (Windows C#, macOS/iOS Swift). See `docs/roadmap.md §M4`. --- @@ -37,9 +32,8 @@ up instantly. Newest status at the top. - [x] **M0 — Scaffolding** ✓ complete - [x] **M1 — Control plane** ✓ complete (2026-06-15) - [x] **M2 — Voice, single stream** ✓ complete (2026-06-16) -- [ ] **M3 — Multi-stream & per-channel tuning** ← next -- [ ] **M3 — Multi-stream & per-channel tuning** (screen audio, listener-side per-user NR) -- [ ] **M4 — Native clients** (Windows C#, macOS/iOS Swift) +- [x] **M3 — Multi-stream & per-channel tuning** ✓ complete (2026-06-16) +- [ ] **M4 — Native clients** (Windows C#, macOS/iOS Swift) ← next - [ ] **M5 — Moderation, polish, beyond** (perms, bans, DRED; then file transfer, E2EE, …) --- @@ -140,6 +134,105 @@ wasn't met. Closed the gap: --- +## M3 — Multi-stream & per-channel tuning ✓ (completed 2026-06-16) + +**Exit criterion:** ✓ `test_m3_multistream` — a real `vc_client` (A) runs two concurrent local +streams (MIC + SCREEN_AUDIO) with distinct stream ids; a second client (B) sees both as +separate `STREAM_STARTED` events and a `VC_EVENT_TALK_STATE` talking edge for A's MIC stream; +B independently sets gain/mute/noise-reduction on each of A's streams without one call +affecting the other; A then joins "Music Room" (channel 2: stereo/128kbps/`OPUS_AUDIO`/no +DTX) and announces a fresh MIC stream there, while B stays in "Lobby" (channel 1: mono/24kbps/ +`OPUS_VOIP`/DTX on) — `vc_get_stream_audio_config` shows their effective Opus config differs +exactly as the server enforces per channel. Passes in ~2.4s; verified across 8 consecutive +standalone runs + 3 consecutive full-suite runs with no flakes. + +Exploration before implementing turned up several bugs/gaps where the wire format already +supported this milestone but the client/server logic didn't — these were fixed as part of M3, +not treated as pre-existing-and-out-of-scope: + +- [x] **Server `stream_id` bug** — `handle_stream_announce` always wrote `stream_id=1`, so a + second stream from the same user silently overwrote the first in + `SessionRegistry::set_user_stream`'s replace-by-id logic. Fixed with a per-session counter + (`ConnSession::next_stream_id_`) + `announced_stream_ids_` (also now validated in + `handle_stream_stop`, rejecting stops for ids the session never announced). +- [x] **Per-channel `AudioConfig` was modeled but never populated/enforced.** + `SessionRegistry::init_default_channels()` now seeds Lobby (id=1: mono, 24kbps, `OPUS_VOIP`, + FEC+DTX on) and a new "Music Room" (id=2: stereo, 128kbps, `OPUS_AUDIO`, FEC+DTX off) with + real `AudioConfig`s; new `SessionRegistry::channel_audio_config(channel_id)` accessor (there + was no per-id channel getter before, only `channel_snapshot()`). `handle_stream_announce` + now treats the channel's config as authoritative (mode/frame_ms/application/fec/dtx/ + complexity), clamping (not overriding) `bitrate_bps` to the channel's ceiling. +- [x] **Client silently dropped `mode`/`dtx`/`complexity`/`application` from `effective_audio`** + even for the single M2 stream — `handle_stream_announce_result` and `sync_remote_streams` + only copied `sample_rate`/`bitrate_bps`/`frame_ms`/`fec` into `OpusParams`. New shared + `opus_params_from_audio_config()` helper (`client.cpp`) fixes both the send and receive + paths. +- [x] `core/src/codec/opus_codec.h/.cpp` — new `OpusApplication` enum + `OpusParams::application` + field; `OpusEncoder::init` now honors it instead of hardcoding `OPUS_APPLICATION_VOIP`. +- [x] `core/src/session/session.h/.cpp` — `Stream` struct extended with the full `AudioConfig` + (mode/bitrate_bps/application/fec/expected_packet_loss/dtx/complexity), not just + sample_rate/frame_ms; `copy_streams()` now copies all of it. +- [x] `core/src/core/client.h/.cpp` — local-stream state is now a `std::unordered_map` keyed by `vc_stream_kind` (one active stream per kind — MIC/SCREEN_AUDIO/ + AUX_DEVICE are each singletons for a client), replacing the M2 single-stream fields. + `StreamAnnounce`/`StreamAnnounceResult` round-trips are now correlated by `request_id` + (already round-tripped on the wire; just wasn't read) via `pending_announce_kind_`, so + multiple concurrent announces from one client resolve to the right `LocalStream`. + `on_capture_frame` takes a `kind` parameter and upmixes mono capture to stereo (duplicate + L=R) when a stream's channel config calls for it. `vc_set_self_mute`'s `mic_muted` only + gates the `MIC` kind — a concurrent `SCREEN_AUDIO` share keeps playing while muted. + `set_remote_stream` now actually wires `noise_reduction` through (previously parsed and + discarded). New `run_talk_timer()` (a small dedicated thread, started alongside the UDP + media path, never the miniaudio callback thread) polls both remote talk-state edges + (`AudioEngine::poll_talk_transitions()`) and local capture-activity edges, emitting + `VC_EVENT_TALK_STATE`. +- [x] **Fixed a thread-join race in `teardown_voice()`** — it's called both from `run_io()`'s + own cleanup and from `disconnect()`, on different threads; without serialization both could + see `udp_thread_`/`talk_timer_thread_` as `joinable()` simultaneously and race to `join()` + the same `std::thread` (UB; surfaced as an intermittent `std::system_error: No such process` + under `ctest`). Added a `teardown_mu_` guard around the whole function. This pre-existed for + `udp_thread_` alone (likely the same root cause as the `test_m1_integration`/`test_m2_voice` + cleanup-path flake noted in the M2 section above) — adding `talk_timer_thread_`'s join just + made it surface more often, so it was fixed properly here rather than carried forward again. +- [x] `core/src/audio/audio_engine.h/.cpp` — `CaptureCallback` gained a `kind` parameter + (the real miniaudio capture device is always tagged `kind=0`/MIC; a second concurrent local + stream is fed via its own `inject_capture(kind, ...)` ring buffer — `inject_taps_`, keyed by + kind — since there is only one real hardware capture device in M3). Fixed a buffer-sizing + bug in `on_playback`'s per-stream decode (`opus_decode`'s `frame_size` parameter is + samples-*per-channel*, not total samples — the old code passed `frames * params_.channels`, + which would have overflowed the decode buffer for any stereo stream). Stereo decoder output + is downmixed (avg L/R) into the engine's mono mix accumulator immediately after decode. + `RemoteStream` gained `recv_ns`/`noise_reduction_enabled` (lazy `ApmProcessor` instantiation + — freed on disable, so no separate instance cap is needed per the roadmap's guidance) and + `last_voice_ms`/`talking` (talk-indicator edge state, updated in `push_recv_frame`); new + `set_stream_noise_reduction()` and `poll_talk_transitions()`. Note: until `VOICECAT_HAS_APM` + is wired to a real WebRTC APM build, the NS toggle is plumbed end-to-end but behaviorally a + passthrough no-op (`ApmPassthrough` doesn't touch PCM) — same situation send-side APM has + been in since M2; M3's job was the plumbing, not the DSP backend. +- [x] **New C ABI surface** (`core/include/voicecat.h`, additive only): + `vc_audio_config` struct + `vc_get_stream_audio_config(c, user_id, stream_id, out)` — the + effective Opus config for a stream you own or a peer's, reading from the (now richer) + `LocalStream`/`session::Stream`. `vc_test_inject_capture(c, stream_id, pcm, samples)` — + clearly-marked **test-only**, forwards to `AudioEngine::inject_capture`, so + `test_m3_multistream` can drive two concurrent synthetic-audio streams through the real ABI + without a microphone. +- [x] `tests/test_m3_multistream.cpp` — the M3 exit criterion (ABI-level, mirrors + `test_voice_client_abi.cpp`'s approach per the M2 lesson). Registered in `tests/CMakeLists.txt`. + +**Explicitly out of scope for this pass** (confirmed with the user before implementing): +- `vc_set_input_device`/`vc_set_input_mode`/`vc_set_push_to_talk`/`vc_list_devices` (device + enumeration + VAD/PTT input gate) — still `VC_ERR_NOT_IMPLEMENTED`. These were mentioned as + "scoped to M3" in the M2 follow-up notes above, but docs/roadmap.md's M3 bullets never + actually listed them — deferred again, now tracked explicitly rather than implicitly. +- Real WASAPI desktop-audio loopback capture for `SCREEN_AUDIO` — the engine now supports + feeding a second concurrent local stream via `inject_capture`, but only synthetic PCM is + wired up; a real loopback capture device is a follow-up. +- True stereo *playback output* — `AudioEngine`'s mixer/output device stays mono; stereo + streams are downmixed after decode (see above). The Opus wire format itself is fully + stereo-correct. + +--- + ## Decisions log All architecture/scope decisions are settled and recorded in diff --git a/core/include/voicecat.h b/core/include/voicecat.h index 3af6d48..5fe9903 100644 --- a/core/include/voicecat.h +++ b/core/include/voicecat.h @@ -158,6 +158,22 @@ typedef struct vc_stream_desc { const char* label; /* human label, e.g. "Microphone" */ } vc_stream_desc; +/* The effective Opus configuration in use for a stream — for a stream you own, this is + * StreamAnnounceResult.effective_audio (channel-enforced, docs/voice.md §3); for a remote + * stream, it's the peer's broadcast StreamInfo.audio. See vc_get_stream_audio_config. */ +typedef struct vc_audio_config { + uint32_t codec; /* 0 = OPUS */ + uint32_t mode; /* 0 = mono, 1 = stereo */ + uint32_t sample_rate; + uint32_t bitrate_bps; + uint32_t frame_ms; + uint32_t application; /* 0 = VOIP, 1 = AUDIO, 2 = LOWDELAY */ + int fec; /* bool */ + uint32_t expected_packet_loss; /* % 0..100 */ + int dtx; /* bool */ + uint32_t complexity; /* 0..10 */ +} vc_audio_config; + typedef struct vc_device { const char* id; const char* name; @@ -208,6 +224,18 @@ VC_API vc_result vc_set_self_mute(vc_client* c, int mic_muted, int deafened); VC_API vc_result vc_set_remote_stream(vc_client* c, uint32_t user_id, uint32_t stream_id, float gain, int muted, int noise_reduction); +/* Effective Opus config in use for (user_id, stream_id) — your own stream or a peer's. + * VC_ERR_INVALID_ARG if the user/stream isn't known. */ +VC_API vc_result vc_get_stream_audio_config(vc_client* c, uint32_t user_id, uint32_t stream_id, + vc_audio_config* out); + +/* TEST-ONLY — not for production use. Bypasses the real capture device, injecting raw PCM + * directly into the named local stream's encode pipeline (see AudioEngine::inject_capture). + * Exists so automated tests can drive the real vc_client/ABI path end-to-end without a + * microphone. `stream_id` is the id returned by vc_stream_start. */ +VC_API vc_result vc_test_inject_capture(vc_client* c, uint32_t stream_id, const int16_t* pcm, + size_t samples); + /* ── Text ─────────────────────────────────────────────────────────────────── */ VC_API vc_result vc_send_text(vc_client* c, vc_text_scope scope, uint32_t target_id, const char* utf8); diff --git a/core/src/audio/audio_engine.cpp b/core/src/audio/audio_engine.cpp index 8efdafb..a4701dc 100644 --- a/core/src/audio/audio_engine.cpp +++ b/core/src/audio/audio_engine.cpp @@ -6,10 +6,19 @@ #include "audio/audio_engine.h" #include +#include #include namespace voicecat::audio { +namespace { +int64_t now_ms() { + return std::chrono::duration_cast( + 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(); + 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(frame_samples_)) break; std::vector 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> AudioEngine::poll_talk_transitions() { + std::vector> 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(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(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 pcm(frames * params_.channels); + std::vector pcm(frames * static_cast(dec_channels)); int n; if (maybe_frame) { n = stream.decoder.decode( maybe_frame->payload.data(), static_cast(maybe_frame->payload.size()), - pcm.data(), static_cast(pcm.size())); + pcm.data(), static_cast(frames)); } else { - n = stream.decoder.decode(nullptr, 0, pcm.data(), - static_cast(pcm.size())); + n = stream.decoder.decode(nullptr, 0, pcm.data(), static_cast(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(params_.sample_rate)); + float g = stream.gain; - for (int i = 0; i < n * static_cast(params_.channels); ++i) - mix[i] += static_cast(static_cast(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(pcm[i * 2]) + + static_cast(pcm[i * 2 + 1])) / 2 + : static_cast(pcm[i]); + for (uint32_t c = 0; c < params_.channels; ++c) + mix[i * params_.channels + c] += static_cast(sample * g); + } } stream.playout_ts += frames; } diff --git a/core/src/audio/audio_engine.h b/core/src/audio/audio_engine.h index ae9e8a7..56192b1 100644 --- a/core/src/audio/audio_engine.h +++ b/core/src/audio/audio_engine.h @@ -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; + // 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; 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.0–2.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> 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 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 inject_ring_; // circular, size = frame_samples_ - std::atomic inject_write_{0}; - std::atomic 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 ring; // circular, size = kInjectCapSamples + std::atomic write{0}; + std::atomic read{0}; + }; + std::mutex inject_mu_; + std::unordered_map> 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 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 last_voice_ms{0}; + bool talking = false; }; mutable std::mutex streams_mu_; std::unordered_map streams_; int frame_samples_ = 960; // 20 ms @48 kHz + + static constexpr int64_t kTalkHangoverMs = 300; }; } // namespace voicecat::audio diff --git a/core/src/codec/opus_codec.cpp b/core/src/codec/opus_codec.cpp index eab637f..1b82fd0 100644 --- a/core/src/codec/opus_codec.cpp +++ b/core/src/codec/opus_codec.cpp @@ -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(p.sample_rate), channels_, - OPUS_APPLICATION_VOIP, &err); + opus_application, &err); if (err != OPUS_OK || !enc_) { err_ = opus_strerror(err); return false; diff --git a/core/src/codec/opus_codec.h b/core/src/codec/opus_codec.h index e4381a2..84ff541 100644 --- a/core/src/codec/opus_codec.h +++ b/core/src/codec/opus_codec.h @@ -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. diff --git a/core/src/core/client.cpp b/core/src/core/client.cpp index 871ac68..cd7bb6e 100644 --- a/core/src/core/client.cpp +++ b/core/src/core/client.cpp @@ -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::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(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(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 stereo_pcm(static_cast(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(media_send_crypto_->peek_send_counter()); - hdr.timestamp = local_timestamp_; - local_timestamp_ += static_cast(samples); + hdr.timestamp = ls.timestamp; + ls.timestamp += static_cast(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(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(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(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(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 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 diff --git a/core/src/core/client.h b/core/src/core/client.h index 9515157..f80f1d1 100644 --- a/core/src/core/client.h +++ b/core/src/core/client.h @@ -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 media_recv_crypto_; voicecat::audio::AudioEngine audio_engine_; - voicecat::codec::OpusEncoder local_encoder_; - std::atomic 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 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 last_capture_ms{0}; + bool talking = false; + }; + mutable std::mutex local_streams_mu_; + std::unordered_map local_streams_; // keyed by vc_stream_kind + std::unordered_map 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> 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 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 self_mic_muted_{false}; std::atomic 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_. diff --git a/core/src/session/session.cpp b/core/src/session/session.cpp index 98f2515..fda1ee6 100644 --- a/core/src/session/session.cpp +++ b/core/src/session/session.cpp @@ -38,6 +38,13 @@ std::vector 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(pb.audio().mode()); + s.bitrate_bps = pb.audio().bitrate_bps(); + s.application = static_cast(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; diff --git a/core/src/session/session.h b/core/src/session/session.h index 4838bfd..521bc59 100644 --- a/core/src/session/session.h +++ b/core/src/session/session.h @@ -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 { diff --git a/core/src/voicecat.cpp b/core/src/voicecat.cpp index c38dda4..098c5b0 100644 --- a/core/src/voicecat.cpp +++ b/core/src/voicecat.cpp @@ -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; diff --git a/server/src/conn_session.cpp b/server/src/conn_session.cpp index 21ef286..335b85e 100644 --- a/server/src/conn_session.cpp +++ b/server/src/conn_session.cpp @@ -2,6 +2,7 @@ #ifdef VOICECAT_HAS_NET +#include #include #include #include @@ -329,15 +330,31 @@ void ConnSession::handle_udp_binding(uint64_t req_id, const voicecat::v1::UdpBin void ConnSession::handle_stream_announce(uint64_t req_id, const voicecat::v1::StreamAnnounce& msg) { uint32_t ssrc = registry_->assign_ssrc(session_id_); + uint32_t stream_id = next_stream_id_++; auto env = make_env(req_id); auto* res = env.mutable_stream_announce_result(); res->set_ok(true); - res->set_stream_id(1); + res->set_stream_id(stream_id); res->set_ssrc(ssrc); + // Per-channel AudioConfig is authoritative (docs/voice.md §3): the channel's mode/ + // frame_ms/application/fec/dtx/complexity/expected_packet_loss apply to every stream + // announced into it, regardless of kind. bitrate_bps is clamped (not overridden) to the + // channel's ceiling so a client may still request less. sample_rate stays + // client-requested-or-48000 — everything runs at 48kHz internally per voice.md §3. auto* eff = res->mutable_effective_audio(); - if (msg.has_requested_audio()) { + auto chan_cfg = registry_->channel_audio_config(registry_->user_channel(user_id_.load())); + uint32_t requested_bps = + msg.has_requested_audio() ? msg.requested_audio().bitrate_bps() : 0; + uint32_t requested_rate = + msg.has_requested_audio() ? msg.requested_audio().sample_rate() : 0; + if (chan_cfg) { + *eff = *chan_cfg; + eff->set_bitrate_bps(requested_bps > 0 ? std::min(requested_bps, chan_cfg->bitrate_bps()) + : chan_cfg->bitrate_bps()); + eff->set_sample_rate(requested_rate > 0 ? requested_rate : 48000); + } else if (msg.has_requested_audio()) { *eff = msg.requested_audio(); } else { eff->set_codec(0); // OPUS @@ -350,10 +367,10 @@ void ConnSession::handle_stream_announce(uint64_t req_id, if (eff->bitrate_bps() == 0) eff->set_bitrate_bps(24000); if (eff->frame_ms() == 0) eff->set_frame_ms(20); - announced_stream_id_ = res->stream_id(); + announced_stream_ids_.push_back(stream_id); voicecat::v1::StreamInfo info; - info.set_stream_id(res->stream_id()); + info.set_stream_id(stream_id); info.set_ssrc(ssrc); info.set_kind(msg.kind()); *info.mutable_audio() = *eff; @@ -372,8 +389,12 @@ void ConnSession::handle_stream_announce(uint64_t req_id, } void ConnSession::handle_stream_stop(const voicecat::v1::StreamStop& msg) { + auto it = std::find(announced_stream_ids_.begin(), announced_stream_ids_.end(), + msg.stream_id()); + if (it == announced_stream_ids_.end()) return; // not ours — ignore (no spoofed stops) + announced_stream_ids_.erase(it); + auto updated = registry_->clear_user_stream(user_id_.load(), msg.stream_id()); - if (msg.stream_id() == announced_stream_id_) announced_stream_id_ = 0; if (updated) { auto bcast = make_env(); auto* ue = bcast.mutable_user_event(); diff --git a/server/src/conn_session.h b/server/src/conn_session.h index 12928cc..31b6107 100644 --- a/server/src/conn_session.h +++ b/server/src/conn_session.h @@ -121,8 +121,11 @@ class ConnSession : public std::enable_shared_from_this { std::unique_ptr send_crypto_; std::unique_ptr recv_crypto_; - // M2: locally-announced stream (single MIC stream per session for now) - uint32_t announced_stream_id_{0}; + // M2/M3: locally-announced streams. The server assigns the stream_id (unique per + // session), so a per-session counter + the set of currently-active ids is enough to + // support multiple concurrent streams (MIC + SCREEN_AUDIO + AUX_DEVICE) per user. + uint32_t next_stream_id_{1}; + std::vector announced_stream_ids_; }; } // namespace voicecat::server diff --git a/server/src/session_registry.cpp b/server/src/session_registry.cpp index 17324f9..1beb464 100644 --- a/server/src/session_registry.cpp +++ b/server/src/session_registry.cpp @@ -17,7 +17,44 @@ void SessionRegistry::init_default_channels() { lobby.proto.set_name("Lobby"); lobby.proto.set_type(voicecat::v1::CHANNEL_PERMANENT); lobby.proto.set_order(0); + { + // Speech profile: mono, low bitrate, FEC+DTX on for resilience/silence-suppression. + auto* a = lobby.proto.mutable_audio(); + a->set_codec(0); + a->set_mode(voicecat::v1::MODE_MONO); + a->set_sample_rate(48000); + a->set_bitrate_bps(24000); + a->set_frame_ms(20); + a->set_application(voicecat::v1::OPUS_VOIP); + a->set_fec(true); + a->set_expected_packet_loss(10); + a->set_dtx(true); + a->set_complexity(5); + } channels_[1] = std::move(lobby); + + ChannelEntry music; + music.proto.set_id(2); + music.proto.set_name("Music Room"); + music.proto.set_type(voicecat::v1::CHANNEL_PERMANENT); + music.proto.set_order(1); + { + // Music/screen-audio profile: stereo, high bitrate, FEC/DTX off (continuous signal). + auto* a = music.proto.mutable_audio(); + a->set_codec(0); + a->set_mode(voicecat::v1::MODE_STEREO); + a->set_sample_rate(48000); + a->set_bitrate_bps(128000); + a->set_frame_ms(20); + a->set_application(voicecat::v1::OPUS_AUDIO); + a->set_fec(false); + a->set_expected_packet_loss(0); + a->set_dtx(false); + a->set_complexity(8); + } + channels_[2] = std::move(music); + + next_channel_id_ = 3; // 1 and 2 are now reserved (Lobby, Music Room) } uint64_t SessionRegistry::register_session(std::weak_ptr session) { @@ -203,6 +240,14 @@ uint32_t SessionRegistry::user_channel(uint32_t user_id) const { return (it == users_.end()) ? 0 : it->second.proto.channel_id(); } +std::optional SessionRegistry::channel_audio_config( + uint32_t channel_id) const { + std::shared_lock lk(mu_); + auto it = channels_.find(channel_id); + if (it == channels_.end()) return std::nullopt; + return it->second.proto.audio(); +} + } // namespace voicecat::server #endif // VOICECAT_HAS_NET diff --git a/server/src/session_registry.h b/server/src/session_registry.h index 5d7ae07..4d72a86 100644 --- a/server/src/session_registry.h +++ b/server/src/session_registry.h @@ -112,6 +112,12 @@ class SessionRegistry { // Return the channel_id of a user (0 if not found). uint32_t user_channel(uint32_t user_id) const; + // Return a channel's authoritative AudioConfig (M3 per-channel Opus tuning), or nullopt + // if the channel doesn't exist. There is no per-id Channel getter today otherwise — + // channel_snapshot() copies every channel, which callers needing just one config should + // avoid. + std::optional channel_audio_config(uint32_t channel_id) const; + private: mutable std::shared_mutex mu_; diff --git a/tests/CMakeLists.txt b/tests/CMakeLists.txt index 23feaab..f06c3e5 100644 --- a/tests/CMakeLists.txt +++ b/tests/CMakeLists.txt @@ -76,4 +76,14 @@ if(VOICECAT_USE_VCPKG_DEPS) target_include_directories(test_voice_client_abi PRIVATE ${VC_TEST_INTERNAL_INCLUDES}) add_test(NAME voice_client_abi COMMAND test_voice_client_abi) set_tests_properties(voice_client_abi PROPERTIES TIMEOUT 60) + + # M3 exit criterion: multi-stream (mic + desktop audio), independent per-stream + # gain/mute/NS, per-channel Opus configurability, talk indicators -- all through the + # real C ABI (vc_client), not raw sockets. + add_executable(test_m3_multistream test_m3_multistream.cpp) + target_link_libraries(test_m3_multistream PRIVATE voicecat::server) + target_compile_features(test_m3_multistream PRIVATE cxx_std_20) + target_include_directories(test_m3_multistream PRIVATE ${VC_TEST_INTERNAL_INCLUDES}) + add_test(NAME m3_multistream COMMAND test_m3_multistream) + set_tests_properties(m3_multistream PROPERTIES TIMEOUT 60) endif() diff --git a/tests/test_m3_multistream.cpp b/tests/test_m3_multistream.cpp new file mode 100644 index 0000000..41064e3 --- /dev/null +++ b/tests/test_m3_multistream.cpp @@ -0,0 +1,356 @@ +/* + * test_m3_multistream — M3 exit criterion, exercised through the real C ABI. + * + * Mirrors test_voice_client_abi.cpp's approach (real vc_client instances, not raw sockets — + * the M2 lesson is that ABI-level coverage is what actually proves the client library works). + * Covers the whole M3 milestone in one flow: + * + * 1. A starts two concurrent local streams (MIC + SCREEN_AUDIO) -- distinct stream ids, + * both visible to B as separate STREAM_STARTED events for the same user. + * 2. Synthetic PCM (vc_test_inject_capture) flows into both of A's streams without crashing + * and without disrupting the control/voice plane; B observes a VC_EVENT_TALK_STATE + * talking=true edge for A's MIC stream while both are still in the same channel (voice + * only relays within a channel, so this must happen before step 4 moves A elsewhere). + * 3. B independently gains/mutes/NS-toggles A's two streams (vc_set_remote_stream) -- + * one call doesn't clobber the other's routing; a bogus stream_id is rejected. + * 4. Per-channel Opus configurability: A joins "Music Room" (channel 2, stereo/128kbps/ + * OPUS_AUDIO/no DTX) before announcing there, while B stays in "Lobby" (channel 1, + * mono/24kbps/OPUS_VOIP/DTX) -- vc_get_stream_audio_config shows the two streams' + * effective config differs exactly as the server enforces it. + */ +#include + +#ifdef VOICECAT_HAS_NET + +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include "voicecat.h" +#include "server.h" +#include "db.h" + +// ── Event tracking ──────────────────────────────────────────────────────────── + +struct StreamEvent { + bool started; // true = STARTED, false = STOPPED + uint32_t user_id; + uint32_t stream_id; +}; + +struct TalkEvent { + uint32_t user_id; + uint32_t stream_id; + bool talking; +}; + +struct EventStore { + std::mutex mu; + std::condition_variable cv; + + bool auth_ok{false}; + uint32_t self_user_id{0}; + bool channel_list_received{false}; + std::vector stream_events; + std::vector talk_events; + bool disconnected{false}; + + const char* label{nullptr}; +}; + +static void on_event(void* user, const vc_event* ev) { + auto* s = static_cast(user); + std::lock_guard lk(s->mu); + switch (ev->type) { + case VC_EVENT_AUTH_RESULT: + s->auth_ok = (ev->result == VC_OK); + s->self_user_id = ev->user_id; + if (!s->auth_ok) std::fprintf(stderr, "[%s] AUTH FAILED: %s\n", + s->label ? s->label : "?", ev->text ? ev->text : "(no msg)"); + break; + case VC_EVENT_CHANNEL_LIST: + s->channel_list_received = true; + break; + case VC_EVENT_STREAM_STARTED: + s->stream_events.push_back({true, ev->user_id, ev->stream_id}); + break; + case VC_EVENT_STREAM_STOPPED: + s->stream_events.push_back({false, ev->user_id, ev->stream_id}); + break; + case VC_EVENT_TALK_STATE: + s->talk_events.push_back({ev->user_id, ev->stream_id, ev->u32a != 0}); + break; + case VC_EVENT_ERROR: + std::fprintf(stderr, "[%s] ERROR rc=%d: %s\n", + s->label ? s->label : "?", ev->result, ev->text ? ev->text : ""); + break; + case VC_EVENT_DISCONNECTED: + s->disconnected = true; + break; + default: + break; + } + s->cv.notify_all(); +} + +template +static bool wait_for(EventStore& s, Pred pred, int timeout_ms) { + auto deadline = std::chrono::steady_clock::now() + std::chrono::milliseconds(timeout_ms); + std::unique_lock lk(s.mu); + return s.cv.wait_until(lk, deadline, [&] { return pred(s); }); +} + +static std::vector make_sine_frame(int frame_idx, float freq_hz, + int frame_samples = 960) { + std::vector pcm(frame_samples); + for (int i = 0; i < frame_samples; ++i) { + float t = static_cast(frame_idx * frame_samples + i) / 48000.0f; + pcm[i] = static_cast(std::sin(2.0f * 3.14159265f * freq_hz * t) * 16000.0f); + } + return pcm; +} + +// ── Test harness ────────────────────────────────────────────────────────────── + +static int g_failures = 0; +#define CHECK(cond) \ + do { \ + if (!(cond)) { \ + std::printf("FAIL: %s (%s:%d)\n", #cond, __FILE__, __LINE__); \ + ++g_failures; \ + } \ + } while (0) + +int main() { + auto tmp = std::filesystem::temp_directory_path() / + ("vctest_m3_" + std::to_string( + std::chrono::steady_clock::now().time_since_epoch().count())); + std::filesystem::create_directories(tmp); + std::string data_dir = tmp.string(); + + std::atomic bound_port{0}; + std::mutex ready_mu; + std::condition_variable ready_cv; + bool ready{false}; + + voicecat::server::Config cfg; + cfg.data_dir = data_dir; + cfg.bind_port = 0; + cfg.media_port = 0; + cfg.server_name = "VoiceCat-M3Test"; + cfg.allow_guests = true; + cfg.on_ready = [&](uint16_t p) { + bound_port.store(p); + { std::lock_guard lk(ready_mu); ready = true; } + ready_cv.notify_all(); + }; + + voicecat::server::Server server(cfg); + std::thread server_thread([&] { server.run(); }); + + { + std::unique_lock lk(ready_mu); + bool ok = ready_cv.wait_for(lk, std::chrono::seconds(10), [&] { return ready; }); + if (!ok) { + std::printf("FAIL: server did not become ready within 10s\n"); + server.stop(); + server_thread.join(); + std::filesystem::remove_all(tmp); + return 1; + } + } + + uint16_t port = bound_port.load(); + std::printf("m3_multistream: server ready on :%u\n", port); + + // ── Client A: guest "M3-A" ──────────────────────────────────────────────── + EventStore evA; + evA.label = "clientA"; + vc_callbacks cbA{on_event, nullptr, &evA}; + vc_config cfgA{"test-clientA", "0.1", VC_LOG_OFF}; + vc_client* clientA = vc_client_create(&cfgA, cbA); + CHECK(clientA != nullptr); + + CHECK(vc_connect(clientA, "127.0.0.1", port) == VC_OK); + CHECK(vc_authenticate_guest(clientA, "M3-A") == VC_OK); + CHECK(wait_for(evA, [](EventStore& s) { return s.auth_ok; }, 8000)); + CHECK(wait_for(evA, [](EventStore& s) { return s.channel_list_received; }, 3000)); + + // ── Client B: guest "M3-B" ──────────────────────────────────────────────── + EventStore evB; + evB.label = "clientB"; + vc_callbacks cbB{on_event, nullptr, &evB}; + vc_config cfgB{"test-clientB", "0.1", VC_LOG_OFF}; + vc_client* clientB = vc_client_create(&cfgB, cbB); + CHECK(clientB != nullptr); + + CHECK(vc_connect(clientB, "127.0.0.1", port) == VC_OK); + CHECK(vc_authenticate_guest(clientB, "M3-B") == VC_OK); + CHECK(wait_for(evB, [](EventStore& s) { return s.auth_ok; }, 8000)); + CHECK(wait_for(evB, [](EventStore& s) { return s.channel_list_received; }, 3000)); + + uint32_t a_uid = 0; + { std::lock_guard lk(evA.mu); a_uid = evA.self_user_id; } + + // Both guests land in channel 1 (Lobby) automatically; give the async UDP binding + // handshake a moment to complete on both clients before announcing streams. + std::this_thread::sleep_for(std::chrono::milliseconds(500)); + + // ── 1. A starts MIC + SCREEN_AUDIO concurrently ────────────────────────── + vc_stream_desc mic_desc{}; + mic_desc.kind = VC_STREAM_MIC; + mic_desc.label = "mic"; + uint32_t mic_sid = 0; + CHECK(vc_stream_start(clientA, &mic_desc, &mic_sid) == VC_OK); + + vc_stream_desc screen_desc{}; + screen_desc.kind = VC_STREAM_SCREEN_AUDIO; + screen_desc.label = "desktop audio"; + uint32_t screen_sid = 0; + CHECK(vc_stream_start(clientA, &screen_desc, &screen_sid) == VC_OK); + + CHECK(mic_sid != 0 && screen_sid != 0 && mic_sid != screen_sid); + + // B observes two distinct STREAM_STARTED events for user A. + bool b_saw_both = wait_for(evB, [&](EventStore& s) { + bool saw_mic = false, saw_screen = false; + for (auto& e : s.stream_events) { + if (!e.started || e.user_id != a_uid) continue; + if (e.stream_id == mic_sid) saw_mic = true; + if (e.stream_id == screen_sid) saw_screen = true; + } + return saw_mic && saw_screen; + }, 5000); + CHECK(b_saw_both); + + // Also wait for A's own view of both streams (vc_test_inject_capture requires the + // LocalStream to be active, which flips on A's io_thread_ independently of -- and not + // necessarily before -- the broadcast B observes above). + bool a_self_saw_both = wait_for(evA, [&](EventStore& s) { + bool saw_mic = false, saw_screen = false; + for (auto& e : s.stream_events) { + if (!e.started || e.user_id != a_uid) continue; + if (e.stream_id == mic_sid) saw_mic = true; + if (e.stream_id == screen_sid) saw_screen = true; + } + return saw_mic && saw_screen; + }, 5000); + CHECK(a_self_saw_both); + + // ── 2. Inject synthetic PCM into both of A's local streams ────────────── + for (int i = 0; i < 25; ++i) { + auto mic_pcm = make_sine_frame(i, 440.0f); + auto screen_pcm = make_sine_frame(i, 880.0f); + CHECK(vc_test_inject_capture(clientA, mic_sid, mic_pcm.data(), mic_pcm.size()) == VC_OK); + CHECK(vc_test_inject_capture(clientA, screen_sid, screen_pcm.data(), screen_pcm.size()) == VC_OK); + std::this_thread::sleep_for(std::chrono::milliseconds(20)); + } + + // No disconnects/errors should have resulted from the dual-stream PCM flow. + { std::lock_guard lk(evA.mu); CHECK(!evA.disconnected); } + { std::lock_guard lk(evB.mu); CHECK(!evB.disconnected); } + + // ── 5. Talk indicators ──────────────────────────────────────────────────── + // While A and B are still both in Lobby (voice actually relays between them here -- + // the SFU forwards within a channel, so this must happen before A moves to Music Room + // in step 4 below), confirm B observed a talking=true edge for A's MIC stream. + bool b_saw_talking = wait_for(evB, [&](EventStore& s) { + for (auto& e : s.talk_events) + if (e.user_id == a_uid && e.stream_id == mic_sid && e.talking) return true; + return false; + }, 3000); + CHECK(b_saw_talking); + + // ── 3. B independently controls gain/mute/NS on each of A's streams ───── + CHECK(vc_set_remote_stream(clientB, a_uid, mic_sid, 1.0f, 0, 0) == VC_OK); + CHECK(vc_set_remote_stream(clientB, a_uid, screen_sid, 0.3f, 1, 1) == VC_OK); + CHECK(vc_set_remote_stream(clientB, a_uid, 0xDEADBEEF, 1.0f, 0, 0) == VC_ERR_INVALID_ARG); + + // Toggle NS on/off a few times -- plumbing should never fault or disrupt the stream. + for (int i = 0; i < 3; ++i) { + CHECK(vc_set_remote_stream(clientB, a_uid, mic_sid, 1.0f, 0, 1) == VC_OK); + CHECK(vc_set_remote_stream(clientB, a_uid, mic_sid, 1.0f, 0, 0) == VC_OK); + } + { std::lock_guard lk(evB.mu); CHECK(!evB.disconnected); } + + // ── 4. Per-channel Opus configurability ────────────────────────────────── + // A moves to "Music Room" (channel 2: stereo/128kbps/OPUS_AUDIO/no DTX) and announces a + // fresh MIC stream there; B stays in "Lobby" (channel 1: mono/24kbps/OPUS_VOIP/DTX) with + // its own MIC stream. Their effective_audio should differ exactly as configured server-side. + CHECK(vc_stream_stop(clientA, mic_sid) == VC_OK); + CHECK(vc_join_channel(clientA, 2, nullptr) == VC_OK); + std::this_thread::sleep_for(std::chrono::milliseconds(300)); + + vc_stream_desc music_mic_desc{}; + music_mic_desc.kind = VC_STREAM_MIC; + music_mic_desc.label = "music-mic"; + uint32_t a_music_mic_sid = 0; + CHECK(vc_stream_start(clientA, &music_mic_desc, &a_music_mic_sid) == VC_OK); + CHECK(wait_for(evA, [&](EventStore& s) { + for (auto& e : s.stream_events) + if (e.started && e.user_id == a_uid && e.stream_id == a_music_mic_sid) return true; + return false; + }, 5000)); + + vc_stream_desc b_mic_desc{}; + b_mic_desc.kind = VC_STREAM_MIC; + b_mic_desc.label = "lobby-mic"; + uint32_t b_mic_sid = 0; + CHECK(vc_stream_start(clientB, &b_mic_desc, &b_mic_sid) == VC_OK); + CHECK(wait_for(evB, [&](EventStore& s) { + uint32_t self = s.self_user_id; + for (auto& e : s.stream_events) + if (e.started && e.user_id == self && e.stream_id == b_mic_sid) return true; + return false; + }, 5000)); + + vc_audio_config a_cfg{}; + vc_audio_config b_cfg{}; + CHECK(vc_get_stream_audio_config(clientA, a_uid, a_music_mic_sid, &a_cfg) == VC_OK); + uint32_t b_uid = 0; + { std::lock_guard lk(evB.mu); b_uid = evB.self_user_id; } + CHECK(vc_get_stream_audio_config(clientB, b_uid, b_mic_sid, &b_cfg) == VC_OK); + + // Music Room: stereo, 128kbps, OPUS_AUDIO, DTX off. Lobby: mono, 24kbps, OPUS_VOIP, DTX on. + CHECK(a_cfg.mode == 1 /* stereo */); + CHECK(b_cfg.mode == 0 /* mono */); + CHECK(a_cfg.bitrate_bps == 128000); + CHECK(b_cfg.bitrate_bps == 24000); + CHECK(a_cfg.application == 1 /* OPUS_AUDIO */); + CHECK(b_cfg.application == 0 /* OPUS_VOIP */); + CHECK(a_cfg.dtx == 0); + CHECK(b_cfg.dtx != 0); + + // ── Cleanup ─────────────────────────────────────────────────────────────── + vc_disconnect(clientA); + vc_disconnect(clientB); + vc_client_destroy(clientA); + vc_client_destroy(clientB); + + server.stop(); + server_thread.join(); + + std::filesystem::remove_all(tmp); + + if (g_failures == 0) { + std::printf("m3_multistream: all checks passed\n"); + return 0; + } + std::printf("m3_multistream: %d failure(s)\n", g_failures); + return 1; +} + +#else // !VOICECAT_HAS_NET + +int main() { + std::printf("m3_multistream: SKIP (VOICECAT_HAS_NET not defined)\n"); + return 0; +} + +#endif // VOICECAT_HAS_NET