diff --git a/cmake/compile_definitions/windows.cmake b/cmake/compile_definitions/windows.cmake index dc2b42bcb..3e36d11ad 100644 --- a/cmake/compile_definitions/windows.cmake +++ b/cmake/compile_definitions/windows.cmake @@ -218,6 +218,10 @@ set(PLATFORM_TARGET_FILES "${CMAKE_SOURCE_DIR}/src/platform/windows/display_vram.cpp" "${CMAKE_SOURCE_DIR}/src/platform/windows/display_wgc.cpp" "${CMAKE_SOURCE_DIR}/src/platform/windows/audio.cpp" + "${CMAKE_SOURCE_DIR}/src/platform/windows/mic_write.h" + "${CMAKE_SOURCE_DIR}/src/platform/windows/mic_write.cpp" + "${CMAKE_SOURCE_DIR}/src/platform/windows/vibepollo_vmic.h" + "${CMAKE_SOURCE_DIR}/src/platform/windows/vibepollo_vmic.cpp" "${CMAKE_SOURCE_DIR}/src/platform/windows/virtual_display.h" ${SUNSHINE_WINDOWS_VDISPLAY_SOURCES} "${CMAKE_SOURCE_DIR}/src/platform/windows/utils.h" diff --git a/docs/mic/uplink-rtp.md b/docs/mic/uplink-rtp.md new file mode 100644 index 000000000..e78e196fa --- /dev/null +++ b/docs/mic/uplink-rtp.md @@ -0,0 +1,95 @@ +# Client microphone uplink (UDP/RTP) + +## Goals + +- Real media-plane uplink (not control-stream `0x3003`). +- Compatible wire with Foundation Sunshine clients (`base+12`, RTSP SETUP `mic`). +- Host sink: **Steam Streaming Microphone only** (no VB-Cable auto-install). +- Advertise **Capable / Ready / Port** in `/serverinfo` (not RS-FEC). +- **Opportunistic RS-FEC**: no capability flag; if a client sends parity shards, use them per-session. + +## Ports + +| Offset | Role | +|--------|------| +| +9 | video | +| +10 | control (ENet) | +| +11 | audio downlink | +| +12 | **mic uplink** (`MIC_STREAM_PORT`) | + +## Discovery (`/serverinfo`, HTTPS) + +| Field | Meaning | +|-------|---------| +| `MicrophoneCapable` | Feature enabled in config (Windows build with mic uplink). | +| `MicrophoneReady` | Steam Streaming Microphone render endpoint is present **now**. | +| `MicrophonePort` | `net::map_port(12)` (informational; RTSP SETUP is authoritative). | + +**Not advertised:** RS-FEC. New clients may send parity without probing a flag. + +Legacy Foundation clients ignore unknown nodes. + +## RTSP + +### DESCRIBE (when `stream_mic` and policy allows) + +``` +m=audio RTP/AVP 96 +a=rtpmap:96 opus/48000/2 +a=fmtp:96 minptime=10;useinbandfec=1 +``` + +### SETUP + +`SETUP .../mic/...` → `Transport: server_port=`, sets session `mic.enabled`. + +### Encryption + +- Support `SS_ENC_MIC` (0x08) when the client enables it (Foundation-compatible IV scheme). +- Do not require encryption for baseline plaintext Foundation clients. + +## Wire + +### Baseline (Foundation clients) + +1. UDP to mic port. +2. RTP (PT 96/97) or 16-bit extended header type `0x5504`. +3. Sequence number: **little-endian** (Foundation client quirk). +4. Payload: Opus mono @ 48 kHz. + +### Opportunistic RS-FEC (new clients) + +- Data: same as baseline. +- FEC: RTP `packetType=127` + `AUDIO_FEC_HEADER` + parity (mirror downlink audio layout). +- Geometry: `RTPA_DATA_SHARDS=4`, `RTPA_FEC_SHARDS=2`. +- **Per-session**: first FEC shard seen enables RS assembly for that session only. +- Incomplete blocks fall back to PLC / skip; never break baseline sessions. + +## Host path + +``` +UDP recv → classify data vs FEC + → (optional SS_ENC_MIC decrypt) + → (optional RS recover if session saw FEC) + → Opus decode → float PCM + → Steam Streaming Microphone (WASAPI render) + → optional default capture switch to Microphone (Steam Streaming Microphone) +``` + +## Config (`sunshine.conf`) + +| Key | Default | Meaning | +|-----|---------|---------| +| `stream_mic` | `true` | Master enable | +| `mic_require_steam` | `true` | If Ready=0, skip DESCRIBE mic / reject useful uplink | +| `mic_sink` | `Speakers (Steam Streaming Microphone)` | Render endpoint patterns | +| `mic_capture_device` | `Microphone (Steam Streaming Microphone)` | Default capture switch (empty = don't switch) | +| `mic_buffer_ms` | `50` | WASAPI render buffer hint | +| `mic_buffer_packets` | `2` | Jitter prebuffer depth | + +## Non-goals + +- No VB-Cable download/install. +- No control-stream mic as the primary path. +- No RS-FEC capability advertisement. +- No claiming Opus in-band FEC is transport RS-FEC. diff --git a/src/config.cpp b/src/config.cpp index 70d5815ba..0a0ba6003 100644 --- a/src/config.cpp +++ b/src/config.cpp @@ -954,6 +954,12 @@ namespace config { true, // install_steam_drivers true, // keep_sink_default true, // auto_capture + true, // stream_mic (Windows uplink; no-op on other platforms without backend) + true, // mic_require_steam + "Speakers (Steam Streaming Microphone)", // mic_sink + "Microphone (Steam Streaming Microphone)", // mic_capture_device + 50, // mic_buffer_ms + 2, // mic_buffer_packets }; stream_t stream { @@ -2013,6 +2019,12 @@ namespace config { bool_f(vars, "install_steam_audio_drivers", audio.install_steam_drivers); bool_f(vars, "keep_sink_default", audio.keep_default); bool_f(vars, "auto_capture_sink", audio.auto_capture); + bool_f(vars, "stream_mic", audio.stream_mic); + bool_f(vars, "mic_require_steam", audio.mic_require_steam); + string_f(vars, "mic_sink", audio.mic_sink); + string_f(vars, "mic_capture_device", audio.mic_capture_device); + int_between_f(vars, "mic_buffer_ms", audio.mic_buffer_ms, {10, 200}); + int_between_f(vars, "mic_buffer_packets", audio.mic_buffer_packets, {1, 16}); string_restricted_f(vars, "origin_web_ui_allowed", nvhttp.origin_web_ui_allowed, {"pc"sv, "lan"sv, "wan"sv}); // reflect origin ACL update immediately in HTTP layer diff --git a/src/config.h b/src/config.h index 79ff97803..2a0531d4f 100644 --- a/src/config.h +++ b/src/config.h @@ -252,6 +252,14 @@ namespace config { bool install_steam_drivers; bool keep_default; bool auto_capture; + + // Client microphone uplink (UDP/RTP → Steam Streaming Microphone) + bool stream_mic; ///< Advertise/accept mic uplink when capable + bool mic_require_steam; ///< If true, DESCRIBE/SETUP mic only when Steam mic is present + std::string mic_sink; ///< WASAPI render endpoint (Steam Streaming Microphone speakers side) + std::string mic_capture_device; ///< Optional default capture device to select for host apps + int mic_buffer_ms; ///< WASAPI render buffer hint (10-200) + int mic_buffer_packets; ///< Jitter prebuffer depth in Opus packets (1-16) }; constexpr int ENCRYPTION_MODE_NEVER = 0; // Never use video encryption, even if the client supports it diff --git a/src/nvhttp.cpp b/src/nvhttp.cpp index 0e945b082..118a0457b 100644 --- a/src/nvhttp.cpp +++ b/src/nvhttp.cpp @@ -2391,6 +2391,17 @@ namespace nvhttp { tree.put("root.VirtualDisplayDriverReady", true); } #endif + // Client microphone uplink advertisement and formal protocol discovery. + { + auto mic = stream::get_mic_status(); + tree.put("root.MicrophoneCapable", mic.capable); + tree.put("root.MicrophoneReady", mic.ready); + if (mic.port) { + tree.put("root.MicrophonePort", mic.port); + } + tree.put("root.MicrophoneProtocol", "moonlight-mic"); + tree.put("root.MicrophoneProtocolVersions", "1"); + } } else { tree.put("root.mac", "00:00:00:00:00:00"); tree.put("root.Permission", "0"); diff --git a/src/platform/common.h b/src/platform/common.h index 3ac9219b9..07d47f1a0 100644 --- a/src/platform/common.h +++ b/src/platform/common.h @@ -613,6 +613,15 @@ namespace platf { virtual ~mic_t() = default; }; + /** + * @brief Saved default capture endpoints so mic uplink can restore them later. + */ + struct capture_snapshot_t { + std::string console_id; + std::string comms_id; + std::string multimedia_id; + }; + class audio_control_t { public: virtual int set_sink(const std::string &sink) = 0; @@ -639,6 +648,45 @@ namespace platf { */ virtual void reset_default_device(const std::string &preferred_device = {}) {} + /** + * @brief Initialize Steam Streaming Microphone render backend for client mic uplink. + * @return 0 on success, -1 if unavailable (Steam not running / endpoint missing). + */ + virtual int init_mic_redirect_device() { + return -1; + } + + /** + * @brief Release the mic redirect backend. + */ + virtual void release_mic_redirect_device() {} + + /** + * @brief Write mono float32 PCM into the mic redirect backend. + */ + virtual int write_mic_pcm(const float * /*samples*/, std::uint32_t /*count*/) { + return -1; + } + + /** + * @brief True if the Steam mic render endpoint can be opened right now. + */ + virtual bool mic_redirect_available() { + return false; + } + + virtual capture_snapshot_t snapshot_capture_defaults() { + return {}; + } + + virtual void switch_capture_to(const std::string & /*device_name*/) {} + + virtual void restore_capture_from(const capture_snapshot_t & /*snapshot*/) {} + + virtual std::string get_current_default_capture_name() { + return {}; + } + virtual ~audio_control_t() = default; }; @@ -657,6 +705,11 @@ namespace platf { std::unique_ptr audio_control(); + /** + * @brief Check for an active microphone redirect endpoint without initializing it. + */ + bool mic_redirect_available(); + /** * @brief Get the display_t instance for the given hwdevice_type. * If display_name is empty, use the first monitor that's compatible you can find diff --git a/src/platform/windows/audio.cpp b/src/platform/windows/audio.cpp index 9d589e713..9f45fc3f7 100644 --- a/src/platform/windows/audio.cpp +++ b/src/platform/windows/audio.cpp @@ -24,6 +24,7 @@ #include "src/logging.h" #include "src/platform/common.h" #include "utf_utils.h" +#include "vibepollo_vmic.h" // Must be the last included file // clang-format off @@ -1621,11 +1622,156 @@ namespace platf::audio { return 0; } + + int init_mic_redirect_device() override { + if (mic_redirect_device) { + return 0; + } + auto device = std::make_unique(); + if (device->init() != 0) { + BOOST_LOG(warning) << "[mic] Steam Streaming Microphone unavailable"sv; + return -1; + } + BOOST_LOG(info) << "[mic] Steam Streaming Microphone backend initialized"sv; + mic_redirect_device = std::move(device); + return 0; + } + + void release_mic_redirect_device() override { + mic_redirect_device.reset(); + BOOST_LOG(info) << "[mic] Steam Streaming Microphone backend released"sv; + } + + int write_mic_pcm(const float *data, std::uint32_t count) override { + if (!mic_redirect_device) { + return -1; + } + return mic_redirect_device->write_pcm(data, count); + } + + bool mic_redirect_available() override { + if (mic_redirect_device) { + return true; + } + + if (!config::audio.mic_sink.empty() && is_sink_available(config::audio.mic_sink)) { + return true; + } + return is_sink_available("Steam Streaming Microphone"); + } + + platf::capture_snapshot_t snapshot_capture_defaults() override { + platf::capture_snapshot_t snap; + auto get_id = [&](ERole role) -> std::string { + device_t dev; + if (FAILED(device_enum->GetDefaultAudioEndpoint(eCapture, role, &dev))) { + return {}; + } + wstring_t id; + if (FAILED(dev->GetId(&id))) { + return {}; + } + return utf_utils::to_utf8(id.get()); + }; + snap.console_id = get_id(eConsole); + snap.comms_id = get_id(eCommunications); + snap.multimedia_id = get_id(eMultimedia); + return snap; + } + + void switch_capture_to(const std::string &device_name) override { + auto target_id = find_capture_device_id(utf_utils::from_utf8(device_name)); + if (target_id.empty()) { + BOOST_LOG(warning) << "[mic] switch_capture_to: device not found: " << device_name; + return; + } + bool any_failed = false; + for (int x = 0; x < (int) ERole_enum_count; ++x) { + if (FAILED(policy->SetDefaultEndpoint(target_id.c_str(), (ERole) x))) { + BOOST_LOG(warning) << "[mic] SetDefaultEndpoint failed for role " << x << " on: " << device_name; + any_failed = true; + } + } + if (!any_failed) { + BOOST_LOG(info) << "[mic] default capture switched to: " << device_name; + } + } + + void restore_capture_from(const platf::capture_snapshot_t &snap) override { + auto restore_role = [&](const std::string &id_utf8, ERole role) { + if (id_utf8.empty()) { + return; + } + auto id = utf_utils::from_utf8(id_utf8); + policy->SetDefaultEndpoint(id.c_str(), role); + }; + restore_role(snap.console_id, eConsole); + restore_role(snap.comms_id, eCommunications); + restore_role(snap.multimedia_id, eMultimedia); + BOOST_LOG(info) << "[mic] default capture roles restored"sv; + } + + std::string get_current_default_capture_name() override { + device_t dev; + if (FAILED(device_enum->GetDefaultAudioEndpoint(eCapture, eConsole, &dev))) { + return {}; + } + prop_t prop; + if (FAILED(dev->OpenPropertyStore(STGM_READ, &prop))) { + return {}; + } + prop_var_t pv; + if (SUCCEEDED(prop->GetValue(PKEY_Device_FriendlyName, &pv.prop)) && pv.prop.vt == VT_LPWSTR) { + return utf_utils::to_utf8(pv.prop.pwszVal); + } + return {}; + } + + std::wstring find_capture_device_id(const std::wstring &name) { + collection_t collection; + if (FAILED(device_enum->EnumAudioEndpoints(eCapture, DEVICE_STATE_ACTIVE, &collection))) { + return {}; + } + UINT count = 0; + collection->GetCount(&count); + for (UINT i = 0; i < count; ++i) { + device_t dev; + if (FAILED(collection->Item(i, &dev))) { + continue; + } + prop_t prop; + if (FAILED(dev->OpenPropertyStore(STGM_READ, &prop))) { + continue; + } + prop_var_t pv; + if (SUCCEEDED(prop->GetValue(PKEY_Device_FriendlyName, &pv.prop)) && pv.prop.vt == VT_LPWSTR) { + if (std::wcscmp(pv.prop.pwszVal, name.c_str()) == 0) { + wstring_t id; + if (SUCCEEDED(dev->GetId(&id))) { + return std::wstring(id.get()); + } + } + } + // Also match adapter friendly name / partial contains for Steam mic aliases + prop_var_t adapter; + if (SUCCEEDED(prop->GetValue(PKEY_DeviceInterface_FriendlyName, &adapter.prop)) && adapter.prop.vt == VT_LPWSTR) { + if (std::wcsstr(adapter.prop.pwszVal, name.c_str()) != nullptr || std::wcsstr(name.c_str(), adapter.prop.pwszVal) != nullptr) { + wstring_t id; + if (SUCCEEDED(dev->GetId(&id))) { + return std::wstring(id.get()); + } + } + } + } + return {}; + } + ~audio_control_t() override = default; policy_t policy; audio::device_enum_t device_enum; std::string assigned_sink; + std::unique_ptr mic_redirect_device; }; } // namespace platf::audio @@ -1653,6 +1799,11 @@ namespace platf { return control; } + bool mic_redirect_available() { + auto control = std::make_unique(); + return control->init() == 0 && control->mic_redirect_available(); + } + std::unique_ptr init() { if (dxgi::init()) { return nullptr; diff --git a/src/platform/windows/mic_write.cpp b/src/platform/windows/mic_write.cpp new file mode 100644 index 000000000..bd38c216f --- /dev/null +++ b/src/platform/windows/mic_write.cpp @@ -0,0 +1,512 @@ +/** + * @file src/platform/windows/mic_write.cpp + * @brief WASAPI render backend implementation for Steam Streaming Microphone. + */ + +#include "mic_write.h" + +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include +#include +#include +#include +#include + +// Must come after mmdeviceapi.h +#include "PolicyConfig.h" +#include "misc.h" +#include "src/config.h" +#include "src/logging.h" +#include "src/platform/common.h" + +using namespace std::literals; + +namespace { + + // Local PKEY definitions — avoid relying on audio.cpp's anonymous-namespace definitions. + const PROPERTYKEY local_PKEY_Device_FriendlyName = { + {0xa45c254e, 0xdf1c, 0x4efd, {0x80, 0x20, 0x67, 0xd1, 0x46, 0xa8, 0x50, 0xe0}}, 14 + }; + + /** + * @brief Waveformat used with IPolicyConfig::SetDeviceFormat. + * + * THE CRITICAL FIX: SubFormat must be KSDATAFORMAT_SUBTYPE_PCM here. + * Using IEEE_FLOAT for SetDeviceFormat causes WASAPI Initialize to fail + * with 0x88890008 (AUDCLNT_E_UNSUPPORTED_FORMAT) on Steam Streaming Microphone. + */ + WAVEFORMATEXTENSIBLE make_recommended_steam_mic_device_waveformat() { + WAVEFORMATEXTENSIBLE wfx {}; + wfx.Format.wFormatTag = WAVE_FORMAT_EXTENSIBLE; + wfx.Format.nChannels = 2; + wfx.Format.nSamplesPerSec = 48000; + wfx.Format.wBitsPerSample = 32; + wfx.Format.nBlockAlign = 8; + wfx.Format.nAvgBytesPerSec = 384000; + wfx.Format.cbSize = sizeof(WAVEFORMATEXTENSIBLE) - sizeof(WAVEFORMATEX); + wfx.Samples.wValidBitsPerSample = 32; + wfx.dwChannelMask = SPEAKER_FRONT_LEFT | SPEAKER_FRONT_RIGHT; + wfx.SubFormat = KSDATAFORMAT_SUBTYPE_PCM; // <-- PCM for SetDeviceFormat + return wfx; + } + + /** + * @brief Waveformat used with IAudioClient::Initialize. + * + * After setting the device format to PCM, WASAPI Initialize succeeds with + * IEEE_FLOAT — this is the render stream format. + */ + WAVEFORMATEXTENSIBLE make_required_steam_mic_render_waveformat() { + WAVEFORMATEXTENSIBLE wfx {}; + wfx.Format.wFormatTag = WAVE_FORMAT_EXTENSIBLE; + wfx.Format.nChannels = 2; + wfx.Format.nSamplesPerSec = 48000; + wfx.Format.wBitsPerSample = 32; + wfx.Format.nBlockAlign = 8; + wfx.Format.nAvgBytesPerSec = 384000; + wfx.Format.cbSize = sizeof(WAVEFORMATEXTENSIBLE) - sizeof(WAVEFORMATEX); + wfx.Samples.wValidBitsPerSample = 32; + wfx.dwChannelMask = SPEAKER_FRONT_LEFT | SPEAKER_FRONT_RIGHT; + wfx.SubFormat = KSDATAFORMAT_SUBTYPE_IEEE_FLOAT; // <-- float for Initialize + return wfx; + } + +} // namespace + +namespace platf::audio { + + int mic_write_wasapi_t::init() { + // Create device enumerator + IMMDeviceEnumerator *raw_enum = nullptr; + auto hr = CoCreateInstance(CLSID_MMDeviceEnumerator, nullptr, CLSCTX_ALL, + IID_IMMDeviceEnumerator, (void **)&raw_enum); + if (FAILED(hr) || !raw_enum) { + BOOST_LOG(warning) << "[mic] mic_write_wasapi: CoCreateInstance(DeviceEnumerator) failed 0x"sv + << util::hex(hr).to_string_view(); + return -1; + } + device_enum.reset(raw_enum); + + // Find the Steam Streaming Microphone render endpoint + std::wstring device_id; + if (!find_target_device(device_id)) { + BOOST_LOG(warning) << "[mic] mic_write_wasapi: Steam Streaming Microphone render endpoint not found"sv; + return -1; + } + + // Normalize device format to PCM (critical — must run before Initialize) + ensure_recommended_steam_mic_format(device_id); + + // Initialize WASAPI with IEEE_FLOAT render stream + if (!initialize_device(device_id)) { + return -1; + } + + return 0; + } + + bool mic_write_wasapi_t::find_target_device(std::wstring &out_device_id) { + // Match a Steam mic render endpoint by friendly-name substring within a given + // device-state mask. Returns the device id of the first match and records its + // friendly name in target_device_name. + auto search_mask = [&](DWORD state_mask, std::wstring &id_out) -> bool { + IMMDeviceCollection *raw_collection = nullptr; + auto hr = device_enum->EnumAudioEndpoints(eRender, state_mask, &raw_collection); + if (FAILED(hr) || !raw_collection) { + return false; + } + util::safe_ptr> collection(raw_collection); + + UINT count = 0; + collection->GetCount(&count); + + for (UINT i = 0; i < count; ++i) { + IMMDevice *raw_device = nullptr; + if (FAILED(collection->Item(i, &raw_device)) || !raw_device) continue; + util::safe_ptr> device(raw_device); + + IPropertyStore *raw_prop = nullptr; + if (FAILED(device->OpenPropertyStore(STGM_READ, &raw_prop)) || !raw_prop) continue; + util::safe_ptr> prop(raw_prop); + + PROPVARIANT pv; + PropVariantInit(&pv); + if (FAILED(prop->GetValue(local_PKEY_Device_FriendlyName, &pv)) || + pv.vt != VT_LPWSTR || !pv.pwszVal) { + PropVariantClear(&pv); + continue; + } + std::wstring name(pv.pwszVal); + PropVariantClear(&pv); + + // Case-insensitive substring match against autodetect patterns + std::wstring name_lower = name; + std::transform(name_lower.begin(), name_lower.end(), name_lower.begin(), ::towlower); + + bool matched = false; + for (const auto &pattern : autodetect_patterns) { + std::wstring pat_lower = pattern; + std::transform(pat_lower.begin(), pat_lower.end(), pat_lower.begin(), ::towlower); + if (name_lower.find(pat_lower) != std::wstring::npos) { + matched = true; + break; + } + } + if (!matched) continue; + + WCHAR *raw_id = nullptr; + if (SUCCEEDED(device->GetId(&raw_id)) && raw_id) { + id_out = std::wstring(raw_id); + CoTaskMemFree(raw_id); + target_device_name = platf::to_utf8(name); + return true; + } + } + return false; + }; + + // 1. Fast path: an ACTIVE endpoint is ready to use. + if (search_mask(DEVICE_STATE_ACTIVE, out_device_id)) { + BOOST_LOG(info) << "[mic] found Steam mic render endpoint: "sv << target_device_name; + return true; + } + + // 2. Self-heal: the Steam Streaming Microphone render endpoint (the loopback + // sink we write decoded mic audio into) can be left disabled — e.g. a user + // disabling it in Sound settings, or a Steam Remote Play toggle. WASAPI then + // won't enumerate it as ACTIVE and passthrough silently dies. Look in the + // recoverable inactive states and re-enable via IPolicyConfig. NOTPRESENT is + // excluded: a driver-absent endpoint can't be re-enabled by visibility, and + // matching a stale not-present entry first would falsely abort the heal. + std::wstring inactive_id; + if (search_mask(DEVICE_STATE_DISABLED | DEVICE_STATE_UNPLUGGED, inactive_id)) { + BOOST_LOG(warning) << "[mic] Steam mic render endpoint ["sv << target_device_name + << "] is present but inactive — attempting to re-enable it."sv; + if (try_enable_endpoint(inactive_id)) { + // Endpoint visibility/state propagation through MMDevice is asynchronous; + // poll for the endpoint to surface as ACTIVE rather than assuming a delay. + for (int attempt = 0; attempt < 20; ++attempt) { + if (search_mask(DEVICE_STATE_ACTIVE, out_device_id)) { + BOOST_LOG(info) << "[mic] re-enabled and found Steam mic render endpoint: "sv << target_device_name; + return true; + } + std::this_thread::sleep_for(std::chrono::milliseconds(100)); + } + } + BOOST_LOG(warning) << "[mic] could not re-enable the Steam mic render endpoint — enable " + "'Speakers (Steam Streaming Microphone)' in Windows Sound settings, or " + "toggle Steam Remote Play microphone."sv; + return false; + } + + return false; + } + + bool mic_write_wasapi_t::try_enable_endpoint(const std::wstring &device_id) { + IPolicyConfig *raw_policy = nullptr; + auto hr = CoCreateInstance(CLSID_CPolicyConfigClient, nullptr, CLSCTX_ALL, + IID_IPolicyConfig, (void **)&raw_policy); + if (FAILED(hr) || !raw_policy) { + BOOST_LOG(warning) << "[mic] try_enable_endpoint: CoCreateInstance(IPolicyConfig) failed 0x"sv + << util::hex(hr).to_string_view(); + return false; + } + util::safe_ptr> policy(raw_policy); + + hr = policy->SetEndpointVisibility(device_id.c_str(), TRUE); + if (FAILED(hr)) { + BOOST_LOG(warning) << "[mic] try_enable_endpoint: SetEndpointVisibility(TRUE) failed 0x"sv + << util::hex(hr).to_string_view(); + return false; + } + return true; + } + + void mic_write_wasapi_t::ensure_recommended_steam_mic_format(const std::wstring &device_id) { + IPolicyConfig *raw_policy = nullptr; + auto hr = CoCreateInstance(CLSID_CPolicyConfigClient, nullptr, CLSCTX_ALL, + IID_IPolicyConfig, (void **)&raw_policy); + if (FAILED(hr) || !raw_policy) { + BOOST_LOG(warning) << "[mic] ensure_format: CoCreateInstance(IPolicyConfig) failed 0x"sv + << util::hex(hr).to_string_view(); + return; + } + util::safe_ptr> policy(raw_policy); + + // Set render endpoint to PCM (the critical fix — not IEEE_FLOAT) + auto device_fmt = make_recommended_steam_mic_device_waveformat(); + auto device_id_copy = device_id; + WAVEFORMATEXTENSIBLE prev_fmt {}; + auto render_hr = policy->SetDeviceFormat(device_id_copy.c_str(), + reinterpret_cast(&device_fmt), + reinterpret_cast(&prev_fmt)); + + if (SUCCEEDED(render_hr)) { + BOOST_LOG(info) << "[mic] Changed Steam microphone render device format for ["sv + << target_device_name << "] to [pcm, 32-bit, 48000 Hz, 2ch]"sv; + } else { + BOOST_LOG(warning) << "[mic] SetDeviceFormat render: 0x"sv + << util::hex(render_hr).to_string_view() + << " (format may already be correct)"sv; + } + + // Also normalize the paired capture endpoint + IMMDeviceCollection *raw_captures = nullptr; + if (FAILED(device_enum->EnumAudioEndpoints(eCapture, DEVICE_STATE_ACTIVE, &raw_captures))) return; + util::safe_ptr> captures(raw_captures); + + UINT count = 0; + captures->GetCount(&count); + for (UINT i = 0; i < count; ++i) { + IMMDevice *raw_dev = nullptr; + if (FAILED(captures->Item(i, &raw_dev)) || !raw_dev) continue; + util::safe_ptr> cap_dev(raw_dev); + + IPropertyStore *raw_prop = nullptr; + if (FAILED(cap_dev->OpenPropertyStore(STGM_READ, &raw_prop)) || !raw_prop) continue; + util::safe_ptr> prop(raw_prop); + + PROPVARIANT pv; + PropVariantInit(&pv); + if (FAILED(prop->GetValue(local_PKEY_Device_FriendlyName, &pv)) || + pv.vt != VT_LPWSTR || !pv.pwszVal) { + PropVariantClear(&pv); + continue; + } + std::wstring cap_name(pv.pwszVal); + PropVariantClear(&pv); + + // Match capture endpoint against steam patterns or generic "Steam Streaming" substring + std::wstring cap_lower = cap_name; + std::transform(cap_lower.begin(), cap_lower.end(), cap_lower.begin(), ::towlower); + + bool is_steam_cap = false; + for (const auto &pattern : autodetect_patterns) { + std::wstring pat_lower = pattern; + std::transform(pat_lower.begin(), pat_lower.end(), pat_lower.begin(), ::towlower); + if (cap_lower.find(pat_lower) != std::wstring::npos) { + is_steam_cap = true; + break; + } + } + if (!is_steam_cap && cap_lower.find(L"steam streaming") != std::wstring::npos) { + is_steam_cap = true; + } + + if (is_steam_cap) { + WCHAR *raw_cap_id = nullptr; + if (FAILED(cap_dev->GetId(&raw_cap_id)) || !raw_cap_id) continue; + std::wstring cap_id(raw_cap_id); + CoTaskMemFree(raw_cap_id); + + auto cap_fmt = make_recommended_steam_mic_device_waveformat(); + WAVEFORMATEXTENSIBLE prev {}; + auto cap_hr = policy->SetDeviceFormat(cap_id.c_str(), + reinterpret_cast(&cap_fmt), + reinterpret_cast(&prev)); + + if (SUCCEEDED(cap_hr)) { + BOOST_LOG(info) << "[mic] Changed Steam microphone capture device format for ["sv + << platf::to_utf8(cap_name) << "] to [pcm, 32-bit, 48000 Hz, 2ch]"sv; + } else { + BOOST_LOG(warning) << "[mic] SetDeviceFormat capture: 0x"sv + << util::hex(cap_hr).to_string_view() + << " (format may already be correct)"sv; + } + } + } + } + + bool mic_write_wasapi_t::initialize_device(const std::wstring &device_id) { + IMMDevice *raw_device = nullptr; + auto hr = device_enum->GetDevice(device_id.c_str(), &raw_device); + if (FAILED(hr) || !raw_device) { + BOOST_LOG(warning) << "[mic] initialize_device: GetDevice failed 0x"sv + << util::hex(hr).to_string_view(); + return false; + } + util::safe_ptr> device(raw_device); + + IAudioClient *raw_client = nullptr; + hr = device->Activate(IID_IAudioClient, CLSCTX_ALL, nullptr, (void **)&raw_client); + if (FAILED(hr) || !raw_client) { + BOOST_LOG(warning) << "[mic] initialize_device: Activate failed 0x"sv + << util::hex(hr).to_string_view(); + return false; + } + audio_client.reset(raw_client); + + // Use IEEE_FLOAT for WASAPI render stream (after PCM SetDeviceFormat above). + // Buffer duration in REFERENCE_TIME (100-nanosecond units): ms * 10000. + // Uses mic_buffer_ms from config (default 50ms). Previously hardcoded to 100ms. + const REFERENCE_TIME buffer_duration = + static_cast(config::audio.mic_buffer_ms) * 10000LL; + auto render_fmt = make_required_steam_mic_render_waveformat(); + hr = audio_client->Initialize(AUDCLNT_SHAREMODE_SHARED, + AUDCLNT_STREAMFLAGS_EVENTCALLBACK, buffer_duration, 0, + reinterpret_cast(&render_fmt), nullptr); + if (FAILED(hr)) { + BOOST_LOG(warning) << "[mic] initialize_device: Initialize failed 0x"sv + << util::hex(hr).to_string_view(); + return false; + } + + hr = audio_client->GetBufferSize(&buffer_frame_count); + if (FAILED(hr)) { + BOOST_LOG(warning) << "[mic] initialize_device: GetBufferSize failed"sv; + return false; + } + + IAudioRenderClient *raw_render = nullptr; + hr = audio_client->GetService(IID_IAudioRenderClient, (void **)&raw_render); + if (FAILED(hr) || !raw_render) { + BOOST_LOG(warning) << "[mic] initialize_device: GetService(IAudioRenderClient) failed 0x"sv + << util::hex(hr).to_string_view(); + return false; + } + audio_render = raw_render; + + HANDLE evt = CreateEvent(nullptr, FALSE, FALSE, nullptr); + if (!evt) { + BOOST_LOG(warning) << "[mic] initialize_device: CreateEvent failed"sv; + return false; + } + render_event.reset(evt); + audio_client->SetEventHandle(render_event.get()); + + hr = audio_client->Start(); + if (FAILED(hr)) { + BOOST_LOG(warning) << "[mic] initialize_device: Start failed 0x"sv + << util::hex(hr).to_string_view(); + return false; + } + + BOOST_LOG(info) << "[mic] WASAPI init / render target: "sv << target_device_name + << " buf="sv << buffer_frame_count << " frames"sv; + + // Only spawn the render thread on first init. When called from render_loop's + // WASAPI re-init path, render_thread is joinable (we ARE the thread). + // Assigning over a joinable std::thread calls std::terminate() — guard against that. + if (!render_thread.joinable()) { + stop_render_thread.store(false); + render_dead.store(false); + render_thread = std::thread(&mic_write_wasapi_t::render_loop, this); + } + return true; + } + + int mic_write_wasapi_t::write_pcm(const float *samples, std::uint32_t frame_count) { + if (!render_event || stop_render_thread.load(std::memory_order_acquire) || + render_dead.load(std::memory_order_acquire)) { + if (render_dead.load(std::memory_order_acquire) && !render_dead_logged) { + render_dead_logged = true; + BOOST_LOG(warning) << "[mic] write_pcm: render thread dead — mic audio being dropped"sv; + } + return -1; + } + { + std::lock_guard lk(queue_mutex); + for (std::uint32_t i = 0; i < frame_count; ++i) { + pending_frames.push_back(std::clamp(samples[i], -1.0f, 1.0f)); + } + // Cap at 1 second to bound latency + if (pending_frames.size() > 48000) { + auto trim = pending_frames.size() - 48000; + pending_frames.erase(pending_frames.begin(), pending_frames.begin() + (std::ptrdiff_t) trim); + } + } + if (!first_packet_written_logged) { + first_packet_written_logged = true; + BOOST_LOG(info) << "[mic] first PCM write to Steam mic backend"sv; + } + SetEvent(render_event.get()); + return 0; + } + + void mic_write_wasapi_t::render_loop() { + auto coinit_hr = CoInitializeEx(nullptr, COINIT_MULTITHREADED | COINIT_SPEED_OVER_MEMORY); + if (FAILED(coinit_hr) && coinit_hr != RPC_E_CHANGED_MODE) { + BOOST_LOG(error) << "[mic] render_loop: CoInitializeEx failed 0x"sv + << util::hex(coinit_hr).to_string_view(); + return; + } + platf::adjust_thread_priority(platf::thread_priority_e::high); + + while (!stop_render_thread.load(std::memory_order_acquire)) { + WaitForSingleObject(render_event.get(), 20); + if (stop_render_thread.load(std::memory_order_acquire)) break; + + UINT32 padding = 0; + HRESULT pad_hr = audio_client->GetCurrentPadding(&padding); + if (FAILED(pad_hr)) { + BOOST_LOG(warning) << "[mic] GetCurrentPadding failed 0x"sv + << util::hex(pad_hr).to_string_view() + << " — attempting WASAPI re-init"sv; + // Stop the existing client before re-init + if (audio_client) audio_client->Stop(); + if (audio_render) { audio_render->Release(); audio_render = nullptr; } + audio_client.reset(); + // Re-find device and re-initialize + std::wstring device_id; + if (!find_target_device(device_id) || !initialize_device(device_id)) { + BOOST_LOG(error) << "[mic] WASAPI re-init failed — render dead"sv; + render_dead.store(true, std::memory_order_release); + break; + } + BOOST_LOG(info) << "[mic] WASAPI re-init succeeded — resuming render"sv; + // Clear stale audio queued before the re-init to prevent replaying it + // into the freshly-opened WASAPI session. + { std::lock_guard lk(queue_mutex); pending_frames.clear(); } + continue; + } + + auto avail = buffer_frame_count - padding; + if (avail == 0) continue; + + UINT32 to_write; + { std::lock_guard lk(queue_mutex); to_write = std::min(avail, (UINT32) pending_frames.size()); } + if (to_write == 0) continue; + + BYTE *buf = nullptr; + if (FAILED(audio_render->GetBuffer(to_write, &buf)) || !buf) continue; + + auto *dst = reinterpret_cast(buf); + { + std::lock_guard lk(queue_mutex); + for (UINT32 f = 0; f < to_write; ++f) { + const float s = pending_frames.front(); + pending_frames.pop_front(); + dst[f * 2] = s; + dst[f * 2 + 1] = s; + } + } + audio_render->ReleaseBuffer(to_write, 0); + } + + if (SUCCEEDED(coinit_hr)) CoUninitialize(); + } + + mic_write_wasapi_t::~mic_write_wasapi_t() { + stop_render_thread.store(true); + if (render_event) SetEvent(render_event.get()); + if (render_thread.joinable()) render_thread.join(); + if (audio_client) audio_client->Stop(); + if (audio_render) { + audio_render->Release(); + audio_render = nullptr; + } + // render_event auto-closes via CloseHandle + // audio_client auto-releases + // device_enum auto-releases + } + +} // namespace platf::audio diff --git a/src/platform/windows/mic_write.h b/src/platform/windows/mic_write.h new file mode 100644 index 000000000..cfc47bc7d --- /dev/null +++ b/src/platform/windows/mic_write.h @@ -0,0 +1,91 @@ +/** + * @file src/platform/windows/mic_write.h + * @brief WASAPI render backend for Steam Streaming Microphone. + * + * Provides mic_redirect_backend_t (abstract base) and mic_write_wasapi_t + * (concrete WASAPI implementation). The critical fix over previous attempts: + * SetDeviceFormat uses KSDATAFORMAT_SUBTYPE_PCM; WASAPI Initialize uses + * KSDATAFORMAT_SUBTYPE_IEEE_FLOAT. Using IEEE_FLOAT for both failed. + */ +#pragma once + +#include +#include +#include +#include +#include +#include +#include +#include + +// WinSock2.h must precede any windows.h pulled in by the WASAPI headers below, +// otherwise the toolchain emits "#warning Please include winsock2.h before +// windows.h", which is fatal under -Werror. +#include + +#include +#include + +#include "src/utility.h" + +namespace platf::audio { + + template + void release_com(T *p) { + p->Release(); + } + + /** + * @brief Abstract base for mic redirect backends. + */ + class mic_redirect_backend_t { + public: + virtual int init() = 0; + virtual int write_pcm(const float *samples, std::uint32_t frame_count) = 0; + virtual const char *backend_id() const = 0; + virtual ~mic_redirect_backend_t() = default; + }; + + /** + * @brief WASAPI render backend that writes decoded PCM to Steam Streaming Microphone. + * + * init() discovers the Steam Streaming Microphone render endpoint, calls + * IPolicyConfig::SetDeviceFormat with a PCM (not float!) waveformat to normalize + * the endpoint, then initializes WASAPI with IEEE_FLOAT for the render stream. + * + * write_pcm() queues mono float32 samples and signals the render thread. + * The render thread expands mono to stereo float32 and writes to WASAPI. + */ + class mic_write_wasapi_t: public mic_redirect_backend_t { + public: + int init() override; + int write_pcm(const float *samples, std::uint32_t frame_count) override; + const char *backend_id() const override { return "mic_write_wasapi"; } + ~mic_write_wasapi_t() override; + + // Set before calling init() + std::vector autodetect_patterns; + + private: + bool find_target_device(std::wstring &out_device_id); + bool try_enable_endpoint(const std::wstring &device_id); + void ensure_recommended_steam_mic_format(const std::wstring &device_id); + bool initialize_device(const std::wstring &device_id); + void render_loop(); + + util::safe_ptr> device_enum; + util::safe_ptr> audio_client; + IAudioRenderClient *audio_render = nullptr; + UINT32 buffer_frame_count = 0; + std::string target_device_name; + bool first_packet_written_logged = false; + bool render_dead_logged = false; + util::safe_ptr_v2 render_event; + std::mutex queue_mutex; + std::deque pending_frames; + std::thread render_thread; + std::atomic stop_render_thread {false}; + std::atomic render_dead {false}; + }; + +} // namespace platf::audio diff --git a/src/platform/windows/vibepollo_vmic.cpp b/src/platform/windows/vibepollo_vmic.cpp new file mode 100644 index 000000000..077688c77 --- /dev/null +++ b/src/platform/windows/vibepollo_vmic.cpp @@ -0,0 +1,55 @@ +/** + * @file src/platform/windows/vibepollo_vmic.cpp + * @brief vibepollo_vmic_t implementation. + */ + +#include "vibepollo_vmic.h" +#include "misc.h" +#include "src/config.h" +#include "src/logging.h" +#include "src/platform/common.h" + +using namespace std::literals; + +namespace platf::audio { + + int vibepollo_vmic_t::init() { + auto backend = std::make_unique(); + + // If mic_sink is configured, use it as the primary search pattern so + // find_target_device honours the user's explicit render endpoint choice. + // Fall back to well-known Steam Streaming Microphone names if mic_sink + // is empty or doesn't match any device. + std::vector patterns; + if (!config::audio.mic_sink.empty()) { + patterns.push_back(platf::from_utf8(config::audio.mic_sink)); + } + // Always include the canonical fallback names so the backend can + // auto-discover the device even when mic_sink uses a shorter alias. + patterns.push_back(L"Steam Streaming Microphone"); + patterns.push_back(L"Speakers (Steam Streaming Microphone)"); + backend->autodetect_patterns = std::move(patterns); + + if (backend->init() != 0) { + log_missing_driver_once(); + return -1; + } + + speaker_backend = std::move(backend); + return 0; + } + + int vibepollo_vmic_t::write_pcm(const float *samples, std::uint32_t frame_count) { + if (!speaker_backend) return -1; + return speaker_backend->write_pcm(samples, frame_count); + } + + void vibepollo_vmic_t::log_missing_driver_once() { + if (!missing_driver_logged) { + missing_driver_logged = true; + BOOST_LOG(warning) << "[mic] Steam Streaming Microphone render endpoint not found — " + "is Steam running on this machine? Mic passthrough disabled."sv; + } + } + +} // namespace platf::audio diff --git a/src/platform/windows/vibepollo_vmic.h b/src/platform/windows/vibepollo_vmic.h new file mode 100644 index 000000000..66f892610 --- /dev/null +++ b/src/platform/windows/vibepollo_vmic.h @@ -0,0 +1,29 @@ +/** + * @file src/platform/windows/vibepollo_vmic.h + * @brief vibepollo_vmic_t — high-level mic redirect backend for Steam Streaming Microphone. + * + * Wraps mic_write_wasapi_t and provides the autodetect patterns for + * Steam Streaming Microphone. This is the entry point called by audio_control_t. + */ +#pragma once + +#include "mic_write.h" +#include + +namespace platf::audio { + + class vibepollo_vmic_t: public mic_redirect_backend_t { + public: + int init() override; + int write_pcm(const float *samples, std::uint32_t frame_count) override; + const char *backend_id() const override { return "steam_streaming_microphone"; } + ~vibepollo_vmic_t() override = default; + + private: + void log_missing_driver_once(); + + std::unique_ptr speaker_backend; + bool missing_driver_logged = false; + }; + +} // namespace platf::audio diff --git a/src/rtsp.cpp b/src/rtsp.cpp index f4806bed1..348d6dd70 100644 --- a/src/rtsp.cpp +++ b/src/rtsp.cpp @@ -23,6 +23,7 @@ extern "C" { #include // lib includes +#include #include #include #ifdef _WIN32 @@ -50,6 +51,11 @@ using asio::ip::udp; using namespace std::literals; +// Sunshine/Foundation extension: microphone stream encryption (not in upstream Limelight.h yet). +#ifndef SS_ENC_MIC + #define SS_ENC_MIC 0x08 +#endif + #ifdef _WIN32 namespace { constexpr wchar_t kVulkanHdrLayerGlobalActiveEventName[] = L"Global\\SunshineVirtualHdrActive"; @@ -214,6 +220,8 @@ namespace rtsp_stream { snapshot->frame_generation_provider = frame_generation_provider; snapshot->lossless_scaling_target_fps = lossless_scaling_target_fps; snapshot->lossless_scaling_rtss_limit = lossless_scaling_rtss_limit; + snapshot->enable_mic = enable_mic; + snapshot->mic_protocol_version = mic_protocol_version; #ifdef _WIN32 snapshot->display_helper_gate = display_helper_gate; #endif @@ -1287,6 +1295,14 @@ namespace rtsp_stream { } } + // Advertise mic encryption when mic uplink can be offered to this client. + const auto mic_status = stream::get_mic_status(); + const bool advertise_mic = mic_status.capable && + (!config::audio.mic_require_steam || mic_status.ready); + if (mic_status.capable) { + encryption_flags_supported |= SS_ENC_MIC; + } + // Report supported and required encryption flags ss << "a=x-ss-general.encryptionSupported:" << encryption_flags_supported << std::endl; ss << "a=x-ss-general.encryptionRequested:" << encryption_flags_requested << std::endl; @@ -1336,17 +1352,29 @@ namespace rtsp_stream { ss << std::endl; } + // Versioned client microphone uplink (UDP/RTP on MIC_STREAM_PORT). + if (advertise_mic) { + const auto mic_port = net::map_port(stream::MIC_STREAM_PORT); + ss << "m=audio " << mic_port << " RTP/AVP 97 127" << std::endl; + ss << "a=rtpmap:97 opus/48000/2" << std::endl; + ss << "a=fmtp:97 minptime=20;useinbandfec=1;stereo=0;sprop-stereo=0" << std::endl; + ss << "a=rtpmap:127 moonlight-rs-fec/48000" << std::endl; + ss << "a=x-ss-mic-protocol:moonlight-mic" << std::endl; + ss << "a=x-ss-mic-versions:1" << std::endl; + } + respond(socket->sock, *session, &option, 200, "OK", req->sequenceNumber, ss.str()); return false; } bool cmd_setup(rtsp_server_t *server, std::shared_ptr socket, std::shared_ptr session, msg_t &&req) { - OPTION_ITEM options[4] {}; + OPTION_ITEM options[5] {}; auto &seqn = options[0]; auto &session_option = options[1]; auto &port_option = options[2]; auto &payload_option = options[3]; + auto &mic_protocol_option = options[4]; seqn.option = const_cast("CSeq"); @@ -1365,6 +1393,36 @@ namespace rtsp_stream { port = net::map_port(stream::VIDEO_STREAM_PORT); } else if (type == "control"sv) { port = net::map_port(stream::CONTROL_PORT); + } else if (type == "mic"sv) { + // Accept SETUP even when the Steam backend is temporarily missing so Foundation + // clients can complete handshake; the stream path no-ops without a sink. + const auto mic_status = stream::get_mic_status(); + if (!mic_status.capable) { + BOOST_LOG(warning) << "Rejecting mic SETUP: microphone uplink is unavailable"sv; + cmd_not_found(server, socket, session, std::move(req)); + return false; + } + if (config::audio.mic_require_steam && !mic_status.ready) { + BOOST_LOG(warning) << "Mic SETUP accepted but Steam Streaming Microphone is not ready"sv; + } + + const char *requested_protocol = nullptr; + for (auto option = req->options; option != nullptr; option = option->next) { + if (boost::iequals(std::string_view {option->option}, "X-SS-Mic-Protocol"sv)) { + requested_protocol = option->content; + break; + } + } + + if (requested_protocol && std::string_view {requested_protocol} != "moonlight-mic/1"sv) { + BOOST_LOG(warning) << "Rejecting unsupported microphone protocol: "sv << requested_protocol; + respond(socket->sock, *session, &seqn, 461, "Unsupported Transport", req->sequenceNumber, {}); + return false; + } + + session->enable_mic = true; + session->mic_protocol_version = requested_protocol ? MIC_PROTOCOL_MOONLIGHT_V1 : MIC_PROTOCOL_FOUNDATION_LEGACY; + port = net::map_port(stream::MIC_STREAM_PORT); } else { cmd_not_found(server, socket, session, std::move(req)); return false; @@ -1395,6 +1453,12 @@ namespace rtsp_stream { port_option.next = &payload_option; + if (type == "mic"sv && session->mic_protocol_version == MIC_PROTOCOL_MOONLIGHT_V1) { + payload_option.next = &mic_protocol_option; + mic_protocol_option.option = const_cast("X-SS-Mic-Protocol"); + mic_protocol_option.content = const_cast("moonlight-mic/1"); + } + respond(socket->sock, *session, &seqn, 200, "OK", req->sequenceNumber, {}); return false; } diff --git a/src/rtsp.h b/src/rtsp.h index 0952fac00..5397bcb89 100644 --- a/src/rtsp.h +++ b/src/rtsp.h @@ -37,6 +37,8 @@ namespace stream { namespace rtsp_stream { constexpr auto RTSP_SETUP_PORT = 21; + constexpr std::uint8_t MIC_PROTOCOL_FOUNDATION_LEGACY = 0; + constexpr std::uint8_t MIC_PROTOCOL_MOONLIGHT_V1 = 1; struct launch_session_t { uint32_t id; @@ -75,6 +77,10 @@ namespace rtsp_stream { int surround_info; std::string surround_params; bool continuous_audio; + /// Client requested microphone uplink via RTSP SETUP type "mic". + bool enable_mic = false; + /// Locked microphone wire protocol selected by RTSP SETUP (0=Foundation legacy, 1=moonlight-mic/1). + std::uint8_t mic_protocol_version = MIC_PROTOCOL_FOUNDATION_LEGACY; bool enable_hdr; // Resolved global/per-client preference for Main10 SDR when the client requests SDR. bool prefer_sdr_10bit = false; diff --git a/src/stream.cpp b/src/stream.cpp index 50d68d920..940ef376a 100644 --- a/src/stream.cpp +++ b/src/stream.cpp @@ -7,13 +7,16 @@ #include #include #include +#include #include #include #include #include #include +#include #include #include +#include #include #include #include @@ -21,20 +24,26 @@ #include #include #include +#include // lib includes #include #include #include +#include + +#include extern "C" { // clang-format off #include +#include #include "rswrapper.h" // clang-format on } // local includes +#include "audio.h" #include "config.h" #include "crypto.h" #include "display_device.h" @@ -113,6 +122,19 @@ using asio::ip::udp; using namespace std::literals; +// Sunshine/Foundation extension: microphone stream encryption (not in upstream Limelight.h yet). +#ifndef SS_ENC_MIC + #define SS_ENC_MIC 0x08 +#endif + +// Foundation mic packet type when using 16-bit extended RTP header. +#ifndef IDX_MIC_DATA_TYPE + #define IDX_MIC_DATA_TYPE 0x5504 +#endif +#ifndef MIC_PACKET_TYPE_OPUS + #define MIC_PACKET_TYPE_OPUS 97 +#endif + namespace stream { namespace { std::string current_server_version() { @@ -225,7 +247,7 @@ namespace stream { enum class socket_e : int { video, ///< Video - audio ///< Audio + audio ///< Audio (host → client) }; namespace session { @@ -532,16 +554,64 @@ namespace stream { std::thread recv_thread; std::thread video_thread; std::thread audio_thread; + std::thread mic_thread; std::thread control_thread; asio::io_context io_context; + // Separate io_context so mic recv can be stopped independently. + asio::io_context mic_io_context; udp::socket video_sock {io_context}; udp::socket audio_sock {io_context}; + udp::socket mic_sock {mic_io_context}; + + std::atomic mic_socket_enabled {false}; + std::atomic mic_sessions_count {0}; + std::mutex mic_start_mutex; control_server_t control_server; }; + struct mic_rs_block_key_t { + std::uint32_t ssrc = 0; + std::uint16_t base_seq = 0; + + bool operator<(const mic_rs_block_key_t &other) const { + return ssrc < other.ssrc || (ssrc == other.ssrc && base_seq < other.base_seq); + } + }; + + struct mic_rs_block_t { + std::array, RTPA_DATA_SHARDS> data; + std::array, RTPA_FEC_SHARDS> fec; + std::array fec_sequences {}; + std::array marks {}; + std::array payload_lengths {}; + std::uint32_t data_ssrc = 0; + std::uint32_t base_timestamp = 0; + std::uint16_t base_seq = 0; + std::size_t block_size = 0; + int data_count = 0; + int fec_count = 0; + bool has_fec_metadata = false; + std::chrono::steady_clock::time_point updated_at; + }; + + struct mic_recent_data_key_t { + std::uint32_t ssrc = 0; + std::uint16_t sequence = 0; + + bool operator<(const mic_recent_data_key_t &other) const { + return ssrc < other.ssrc || (ssrc == other.ssrc && sequence < other.sequence); + } + }; + + struct mic_recent_data_t { + std::vector wire_payload; + std::uint32_t timestamp = 0; + std::chrono::steady_clock::time_point updated_at; + }; + struct session_t { config_t config; int stream_fps = 0; @@ -636,6 +706,39 @@ namespace stream { safe::mail_raw_t::event_t shutdown_event; safe::signal_t controlEnd; + // Microphone uplink (client → host via MIC_STREAM_PORT) + struct { + bool enabled = false; ///< Session requested mic (RTSP SETUP mic) + std::uint8_t protocol_version = rtsp_stream::MIC_PROTOCOL_FOUNDATION_LEGACY; + bool capture_switched = false; + platf::capture_snapshot_t capture_snap {}; + // Keep the shared audio context alive so the Steam mic backend outlives this session. + audio::audio_ctx_ref_t audio_ctx; + std::unique_ptr decoder {nullptr, opus_decoder_destroy}; + std::mutex lock; + + // Baseline jitter buffer (Foundation-compatible, no RS) + struct queued_t { + std::vector opus; + std::uint16_t seq = 0; + }; + std::map pending; + bool has_playout_cursor = false; + std::uint16_t expected_seq = 0; + std::uint32_t data_ssrc = 0; + std::uint32_t parity_ssrc = 0; + + // moonlight-mic/1 keeps a small block window for data/parity reordering. + std::map rs_blocks; + std::map recent_data; + std::unique_ptr rs {nullptr, reed_solomon_release}; + + std::uint64_t packets_received = 0; + std::uint64_t frames_written = 0; + std::uint64_t decode_errors = 0; + std::uint64_t rs_recovered = 0; + } mic; + std::atomic state; // Real-time performance counters (updated by broadcast/control threads) @@ -721,6 +824,15 @@ namespace stream { static auto broadcast = safe::make_shared(start_broadcast, end_broadcast); + bool ensure_mic_sock_open(broadcast_ctx_t &ctx); + bool start_mic_receiver(broadcast_ctx_t &ctx); + void mic_session_acquire(broadcast_ctx_t &ctx); + void mic_session_release(broadcast_ctx_t &ctx); + int mic_session_start(session_t &session); + void mic_session_stop(session_t &session); + void micRecvThread(broadcast_ctx_t &ctx); + + std::optional decode_control_packet(std::string_view packet_bytes) { if (packet_bytes.size() < sizeof(std::uint16_t)) { return std::nullopt; @@ -2417,6 +2529,1036 @@ namespace stream { shutdown_event->raise(true); } + namespace { + // Mono Opus @ 48 kHz with 10 ms packets (matches DESCRIBE minptime=10). + constexpr int kMicSampleRate = 48000; + constexpr int kMicFrameSamples = 480; // 10 ms + constexpr int kMicMaxFrameSamples = 5760; // 120 ms Opus limit + constexpr std::size_t kMicMaxQueued = 32; + + platf::audio_control_t *mic_audio_control(session_t &session) { + if (session.mic.audio_ctx && session.mic.audio_ctx->control) { + return session.mic.audio_ctx->control.get(); + } + return nullptr; + } + + int decrypt_mic_payload(session_t &session, std::uint16_t seq, std::vector &payload) { + if (!(session.config.encryptionFlagsEnabled & SS_ENC_MIC)) { + return 0; + } + if (payload.empty()) { + return -1; + } + + // Same IV construction as audio downlink CBC, but keyed by mic seq. + crypto::aes_t iv(16, 0); + *(std::uint32_t *) iv.data() = util::endian::big(session.audio.avRiKeyId + (seq & 0xFFFF)); + + auto &key = session.audio.cipher.key; + if (key.size() < 16) { + return -1; + } + + crypto::cipher_ctx_t ctx {EVP_CIPHER_CTX_new()}; + if (!ctx) { + return -1; + } + + if (EVP_DecryptInit_ex(ctx.get(), EVP_aes_128_cbc(), nullptr, key.data(), iv.data()) != 1) { + return -1; + } + EVP_CIPHER_CTX_set_padding(ctx.get(), 1); + + std::vector plain(payload.size() + 16); + int update_len = 0; + int final_len = 0; + if (EVP_DecryptUpdate(ctx.get(), plain.data(), &update_len, payload.data(), (int) payload.size()) != 1) { + return -1; + } + if (EVP_DecryptFinal_ex(ctx.get(), plain.data() + update_len, &final_len) != 1) { + return -1; + } + + plain.resize((std::size_t) update_len + (std::size_t) final_len); + payload = std::move(plain); + return 0; + } + + void mic_decode_and_write(session_t &session, const std::uint8_t *opus, std::size_t opus_len) { + auto *control = mic_audio_control(session); + if (!control || !session.mic.decoder) { + return; + } + + std::array pcm {}; + const int samples = opus_decode_float( + session.mic.decoder.get(), + opus_len ? opus : nullptr, + (int) opus_len, + pcm.data(), + kMicMaxFrameSamples, + 0 + ); + if (samples < 0) { + ++session.mic.decode_errors; + BOOST_LOG(verbose) << "[mic] opus_decode_float failed: "sv << samples; + return; + } + if (samples == 0) { + return; + } + + if (control->write_mic_pcm(pcm.data(), (std::uint32_t) samples) == 0) { + ++session.mic.frames_written; + } + } + + std::optional mic_v1_block_index(std::uint16_t base_sequence, std::uint16_t sequence) { + const auto offset = (std::uint16_t) (sequence - base_sequence); + if (offset >= RTPA_DATA_SHARDS) { + return std::nullopt; + } + return (int) offset; + } + + bool mic_waiting_for_v1_fec(const session_t &session, std::uint16_t seq) { + if (session.mic.protocol_version != rtsp_stream::MIC_PROTOCOL_MOONLIGHT_V1) { + return false; + } + + const auto now = std::chrono::steady_clock::now(); + for (const auto &[key, block] : session.mic.rs_blocks) { + const auto index = mic_v1_block_index(key.base_seq, seq); + if (block.has_fec_metadata && index && block.marks[*index] != 0 && + now - block.updated_at <= std::chrono::milliseconds(RTPQ_OOS_WAIT_TIME_MS)) { + return true; + } + } + + // Before parity arrives there is no authoritative block base. A nearby + // later data packet is enough to defer PLC briefly without inventing one. + for (const auto &[key, recent] : session.mic.recent_data) { + const auto distance = (std::uint16_t) (key.sequence - seq); + if (distance > 0 && distance < RTPA_DATA_SHARDS && + now - recent.updated_at <= std::chrono::milliseconds(RTPQ_OOS_WAIT_TIME_MS)) { + return true; + } + } + return false; + } + + void mic_queue_opus(session_t &session, std::uint16_t seq, std::vector &&opus) { + auto &mic = session.mic; + const int configured_prebuffer = std::clamp(config::audio.mic_buffer_packets, 1, 16); + int prebuffer = configured_prebuffer; + if (mic.protocol_version == rtsp_stream::MIC_PROTOCOL_MOONLIGHT_V1) { + prebuffer = std::max(prebuffer, RTPA_DATA_SHARDS); + } + + if (!mic.has_playout_cursor) { + if (mic.pending.empty()) { + mic.expected_seq = seq; + } else { + const auto delta = (std::int16_t) (seq - mic.expected_seq); + if (delta < 0 && delta >= -16) { + mic.expected_seq = seq; + } + } + mic.pending[seq] = decltype(mic.pending)::mapped_type {std::move(opus), seq}; + if ((int) mic.pending.size() >= prebuffer) { + mic.has_playout_cursor = true; + } else if (mic.pending.size() >= kMicMaxQueued) { + mic.has_playout_cursor = true; + } else { + return; + } + } else { + // Drop hopelessly late packets (more than half the seq space behind). + const std::int16_t delta = (std::int16_t) (seq - mic.expected_seq); + if (delta < -16) { + return; + } + mic.pending[seq] = decltype(mic.pending)::mapped_type {std::move(opus), seq}; + } + + while (mic.pending.size() > kMicMaxQueued) { + auto farthest = std::max_element(mic.pending.begin(), mic.pending.end(), [&](const auto &left, const auto &right) { + return (std::uint16_t) (left.first - mic.expected_seq) < + (std::uint16_t) (right.first - mic.expected_seq); + }); + if (farthest == mic.pending.end() || farthest->first == mic.expected_seq) { + break; + } + mic.pending.erase(farthest); + } + + // Contiguous playout; limited PLC via Opus null-packet on small gaps. + int missing_skipped = 0; + while (true) { + auto it = mic.pending.find(mic.expected_seq); + if (it != mic.pending.end()) { + mic_decode_and_write(session, it->second.opus.data(), it->second.opus.size()); + mic.pending.erase(it); + ++mic.expected_seq; + missing_skipped = 0; + continue; + } + + if (mic.pending.empty()) { + break; + } + + // If we already have later packets, fill a short gap with PLC then continue. + auto next_it = std::min_element(mic.pending.begin(), mic.pending.end(), [&](const auto &left, const auto &right) { + return (std::uint16_t) (left.first - mic.expected_seq) < + (std::uint16_t) (right.first - mic.expected_seq); + }); + const auto next = next_it->first; + const auto gap = (std::uint16_t) (next - mic.expected_seq); + if (gap >= 0x8000) { + mic.pending.erase(next_it); + continue; + } + if (mic_waiting_for_v1_fec(session, mic.expected_seq)) { + break; + } + if (gap > 8 || missing_skipped >= 4) { + // Large loss: jump forward rather than synthesizing a long stretch. + mic.expected_seq = next; + missing_skipped = 0; + continue; + } + + mic_decode_and_write(session, nullptr, 0); + ++mic.expected_seq; + ++missing_skipped; + } + } + + bool mic_ensure_rs(session_t &session) { + if (!session.mic.rs) { + session.mic.rs.reset(reed_solomon_new(RTPA_DATA_SHARDS, RTPA_FEC_SHARDS)); + if (session.mic.rs) { + // Fixed OpenFEC/NVIDIA matrix used by moonlight-mic/1. + const unsigned char parity[] = {0x77, 0x40, 0x38, 0x0e, 0xc7, 0xa7, 0x0d, 0x6c}; + std::memcpy(session.mic.rs->p, parity, sizeof(parity)); + } + } + return session.mic.rs != nullptr; + } + + mic_rs_block_t &mic_get_rs_block(session_t &session, std::uint32_t data_ssrc, std::uint16_t base_seq) { + auto &blocks = session.mic.rs_blocks; + const mic_rs_block_key_t key {data_ssrc, base_seq}; + auto [it, inserted] = blocks.try_emplace(key); + if (inserted) { + it->second.data_ssrc = data_ssrc; + it->second.base_seq = base_seq; + it->second.marks.fill(1); + } + it->second.updated_at = std::chrono::steady_clock::now(); + + if (inserted && blocks.size() > RTPA_CACHED_FEC_BLOCK_LIMIT) { + auto oldest = blocks.end(); + for (auto candidate = blocks.begin(); candidate != blocks.end(); ++candidate) { + if (candidate != it && (oldest == blocks.end() || candidate->second.updated_at < oldest->second.updated_at)) { + oldest = candidate; + } + } + if (oldest != blocks.end()) { + blocks.erase(oldest); + } + } + + return it->second; + } + + bool mic_store_recent_v1_data( + session_t &session, + std::uint32_t data_ssrc, + std::uint16_t sequence, + std::uint32_t timestamp, + const std::vector &wire_payload + ) { + constexpr std::size_t kRecentDataLimit = RTPA_CACHED_FEC_BLOCK_LIMIT * RTPA_DATA_SHARDS; + auto &recent_data = session.mic.recent_data; + const mic_recent_data_key_t key {data_ssrc, sequence}; + auto [it, inserted] = recent_data.try_emplace(key); + if (!inserted) { + return false; + } + + it->second.wire_payload = wire_payload; + it->second.timestamp = timestamp; + it->second.updated_at = std::chrono::steady_clock::now(); + if (recent_data.size() > kRecentDataLimit) { + auto oldest = recent_data.end(); + for (auto candidate = recent_data.begin(); candidate != recent_data.end(); ++candidate) { + if (candidate != it && (oldest == recent_data.end() || candidate->second.updated_at < oldest->second.updated_at)) { + oldest = candidate; + } + } + if (oldest != recent_data.end()) { + recent_data.erase(oldest); + } + } + return true; + } + + bool mic_v1_data_matches_block( + const mic_rs_block_t &block, + int index, + std::uint32_t timestamp, + const std::vector &wire_payload + ) { + return block.has_fec_metadata && + wire_payload.size() == block.payload_lengths[index] && + (index != 0 || timestamp == block.base_timestamp); + } + + void mic_attach_v1_data_to_block( + mic_rs_block_t &block, + int index, + const std::vector &wire_payload + ) { + block.data[index] = wire_payload; + block.marks[index] = 0; + ++block.data_count; + block.updated_at = std::chrono::steady_clock::now(); + } + + void mic_rs_try_recover(session_t &session, mic_rs_block_t &block) { + auto &mic = session.mic; + if (!block.has_fec_metadata || block.block_size == 0 || !mic_ensure_rs(session)) { + return; + } + if (block.data_count + block.fec_count < RTPA_DATA_SHARDS) { + return; + } + if (block.data_count == RTPA_DATA_SHARDS) { + // Already complete — nothing to recover. + return; + } + + std::array shards_p {}; + std::array, RTPA_TOTAL_SHARDS> storage {}; + for (int i = 0; i < RTPA_DATA_SHARDS; ++i) { + storage[i].assign(block.block_size, 0); + if (!block.data[i].empty()) { + if (block.data[i].size() > block.block_size) { + return; + } + std::memcpy(storage[i].data(), block.data[i].data(), block.data[i].size()); + } + shards_p[i] = storage[i].data(); + } + for (int i = 0; i < RTPA_FEC_SHARDS; ++i) { + storage[RTPA_DATA_SHARDS + i].assign(block.block_size, 0); + if (!block.fec[i].empty()) { + if (block.fec[i].size() != block.block_size) { + return; + } + std::memcpy(storage[RTPA_DATA_SHARDS + i].data(), block.fec[i].data(), block.fec[i].size()); + } + shards_p[RTPA_DATA_SHARDS + i] = storage[RTPA_DATA_SHARDS + i].data(); + } + + auto marks = block.marks; + if (reed_solomon_decode(mic.rs.get(), shards_p.data(), marks.data(), RTPA_TOTAL_SHARDS, (int) block.block_size) != 0) { + return; + } + + for (int i = 0; i < RTPA_DATA_SHARDS; ++i) { + if (block.marks[i] == 0) { + continue; // already present + } + const auto payload_length = block.payload_lengths[i]; + if (payload_length == 0 || payload_length > block.block_size) { + return; + } + + std::vector recovered_wire_payload(shards_p[i], shards_p[i] + payload_length); + auto opus = recovered_wire_payload; + if (decrypt_mic_payload(session, (std::uint16_t) (block.base_seq + i), opus) != 0) { + BOOST_LOG(verbose) << "[mic] decrypt failed for recovered seq="sv << (std::uint16_t) (block.base_seq + i); + continue; + } + + block.data[i] = std::move(recovered_wire_payload); + block.marks[i] = 0; + ++block.data_count; + ++mic.rs_recovered; + mic_queue_opus(session, (std::uint16_t) (block.base_seq + i), std::move(opus)); + } + } + + void mic_handle_v1_fec_shard( + session_t &session, + std::uint32_t parity_ssrc, + std::uint16_t parity_sequence, + std::uint32_t data_ssrc, + std::uint16_t base_seq, + std::uint32_t base_timestamp, + std::uint8_t shard_index, + const std::array &payload_lengths, + std::vector &&parity + ) { + if (parity_ssrc == data_ssrc || + (session.mic.data_ssrc != 0 && session.mic.data_ssrc != data_ssrc) || + (session.mic.parity_ssrc != 0 && session.mic.parity_ssrc != parity_ssrc)) { + return; + } + if (session.mic.parity_ssrc == 0) { + session.mic.parity_ssrc = parity_ssrc; + } + if (session.mic.data_ssrc == 0) { + session.mic.data_ssrc = data_ssrc; + } + + auto &block = mic_get_rs_block(session, data_ssrc, base_seq); + + if (block.has_fec_metadata) { + if (block.base_timestamp != base_timestamp || block.block_size != parity.size() || + block.payload_lengths != payload_lengths) { + return; + } + } else { + for (int i = 0; i < RTPA_DATA_SHARDS; ++i) { + const mic_recent_data_key_t key {data_ssrc, (std::uint16_t) (base_seq + i)}; + const auto recent = session.mic.recent_data.find(key); + if (recent != session.mic.recent_data.end() && + (recent->second.wire_payload.size() != payload_lengths[i] || + (i == 0 && recent->second.timestamp != base_timestamp))) { + return; + } + } + block.base_timestamp = base_timestamp; + block.block_size = parity.size(); + block.payload_lengths = payload_lengths; + block.has_fec_metadata = true; + + for (int i = 0; i < RTPA_DATA_SHARDS; ++i) { + const mic_recent_data_key_t key {data_ssrc, (std::uint16_t) (base_seq + i)}; + const auto recent = session.mic.recent_data.find(key); + if (recent != session.mic.recent_data.end() && block.marks[i] != 0) { + mic_attach_v1_data_to_block(block, i, recent->second.wire_payload); + } + } + } + + if (block.marks[RTPA_DATA_SHARDS + shard_index] == 0) { + return; + } + for (int i = 0; i < RTPA_FEC_SHARDS; ++i) { + if (i != shard_index && block.marks[RTPA_DATA_SHARDS + i] == 0 && + block.fec_sequences[i] == parity_sequence) { + return; + } + } + + block.fec[shard_index] = std::move(parity); + block.fec_sequences[shard_index] = parity_sequence; + block.marks[RTPA_DATA_SHARDS + shard_index] = 0; + ++block.fec_count; + mic_rs_try_recover(session, block); + } + + void mic_handle_v1_data_shard( + session_t &session, + std::uint16_t seq, + std::uint32_t timestamp, + std::uint32_t data_ssrc, + std::vector &&wire_payload + ) { + if ((session.mic.data_ssrc != 0 && session.mic.data_ssrc != data_ssrc) || + session.mic.parity_ssrc == data_ssrc) { + return; + } + if (session.mic.data_ssrc == 0) { + session.mic.data_ssrc = data_ssrc; + } + + mic_rs_block_t *matching_block = nullptr; + int matching_index = 0; + for (auto &[key, block] : session.mic.rs_blocks) { + if (key.ssrc != data_ssrc || !block.has_fec_metadata) { + continue; + } + const auto index = mic_v1_block_index(key.base_seq, seq); + if (!index) { + continue; + } + if (matching_block) { + return; // Reject ambiguous overlapping FEC metadata. + } + matching_block = █ + matching_index = *index; + } + + if (matching_block && + (matching_block->marks[matching_index] == 0 || + !mic_v1_data_matches_block(*matching_block, matching_index, timestamp, wire_payload))) { + return; + } + + auto opus = wire_payload; + if (decrypt_mic_payload(session, seq, opus) != 0) { + BOOST_LOG(verbose) << "[mic] decrypt failed seq="sv << seq; + return; + } + + if (!mic_store_recent_v1_data(session, data_ssrc, seq, timestamp, wire_payload)) { + return; + } + + // FEC protects the exact wire payload. Keep this copy encrypted while a + // second copy follows the normal decrypt/jitter path immediately. + if (matching_block) { + mic_attach_v1_data_to_block(*matching_block, matching_index, wire_payload); + mic_rs_try_recover(session, *matching_block); + } + + mic_queue_opus(session, seq, std::move(opus)); + } + + void mic_handle_legacy_data(session_t &session, std::uint16_t seq, std::vector &&payload) { + if (decrypt_mic_payload(session, seq, payload) != 0) { + BOOST_LOG(verbose) << "[mic] legacy decrypt failed seq="sv << seq; + return; + } + mic_queue_opus(session, seq, std::move(payload)); + } + + enum class mic_packet_kind_e : std::uint8_t { + data, + fec + }; + + struct normalized_mic_packet_t { + mic_packet_kind_e kind = mic_packet_kind_e::data; + std::uint16_t sequence = 0; + std::uint32_t timestamp = 0; + std::uint32_t ssrc = 0; + const std::uint8_t *payload = nullptr; + std::size_t payload_length = 0; + + std::uint8_t fec_shard_index = 0; + std::uint16_t fec_base_sequence = 0; + std::uint32_t fec_base_timestamp = 0; + std::uint32_t fec_data_ssrc = 0; + std::array fec_payload_lengths {}; + }; + + constexpr std::size_t kRtpFixedHeaderSize = 12; + constexpr std::size_t kAudioFecHeaderSize = 12; + constexpr std::size_t kMicFecV1ExtensionSize = 12; + constexpr std::uint8_t kMicFecEncryptedFlag = 0x01; + + std::uint16_t mic_read_be16(const std::uint8_t *data) { + return (std::uint16_t) ((data[0] << 8) | data[1]); + } + + std::uint32_t mic_read_be32(const std::uint8_t *data) { + return ((std::uint32_t) data[0] << 24) | + ((std::uint32_t) data[1] << 16) | + ((std::uint32_t) data[2] << 8) | + data[3]; + } + + bool mic_parse_v1_rtp_header( + const std::uint8_t *data, + std::size_t bytes, + normalized_mic_packet_t &packet, + std::uint8_t &payload_type, + std::size_t &header_length + ) { + if (!data || bytes < kRtpFixedHeaderSize || (data[0] & 0xc0) != 0x80 || (data[0] & 0x20) != 0 || + (data[1] & 0x80) != 0) { + return false; + } + + header_length = kRtpFixedHeaderSize + (std::size_t) (data[0] & 0x0f) * 4; + if (header_length > bytes) { + return false; + } + + if ((data[0] & 0x10) != 0) { + if (bytes - header_length < 4) { + return false; + } + const auto extension_words = mic_read_be16(data + header_length + 2); + const auto extension_bytes = (std::size_t) extension_words * 4; + if (extension_bytes > bytes - header_length - 4) { + return false; + } + header_length += 4 + extension_bytes; + } + + payload_type = data[1] & 0x7f; + packet.sequence = mic_read_be16(data + 2); + packet.timestamp = mic_read_be32(data + 4); + packet.ssrc = mic_read_be32(data + 8); + return packet.ssrc != 0; + } + + bool mic_parse_moonlight_v1( + const session_t &session, + const std::uint8_t *data, + std::size_t bytes, + normalized_mic_packet_t &packet + ) { + std::uint8_t payload_type = 0; + std::size_t header_length = 0; + if (!mic_parse_v1_rtp_header(data, bytes, packet, payload_type, header_length)) { + return false; + } + + const auto *rtp_payload = data + header_length; + const auto rtp_payload_length = bytes - header_length; + if (payload_type == MIC_PACKET_TYPE_OPUS) { + if (rtp_payload_length == 0 || rtp_payload_length > MAX_AUDIO_PACKET_SIZE) { + return false; + } + packet.kind = mic_packet_kind_e::data; + packet.payload = rtp_payload; + packet.payload_length = rtp_payload_length; + return true; + } + + if (payload_type != 127 || rtp_payload_length < kAudioFecHeaderSize + kMicFecV1ExtensionSize + 1) { + return false; + } + + const auto *fec = rtp_payload; + const auto *extension = fec + kAudioFecHeaderSize; + const auto extension_length = mic_read_be16(extension + 2); + if (fec[0] >= RTPA_FEC_SHARDS || fec[1] != MIC_PACKET_TYPE_OPUS || extension[0] != 1 || + extension_length != kMicFecV1ExtensionSize || (extension[1] & ~kMicFecEncryptedFlag) != 0) { + return false; + } + + const bool protected_payload_encrypted = (extension[1] & kMicFecEncryptedFlag) != 0; + const bool session_encrypted = (session.config.encryptionFlagsEnabled & SS_ENC_MIC) != 0; + if (protected_payload_encrypted != session_encrypted) { + return false; + } + + const auto fec_headers_length = kAudioFecHeaderSize + (std::size_t) extension_length; + if (fec_headers_length >= rtp_payload_length) { + return false; + } + const auto parity_length = rtp_payload_length - fec_headers_length; + if (parity_length == 0 || parity_length > MAX_AUDIO_PACKET_SIZE) { + return false; + } + + packet.kind = mic_packet_kind_e::fec; + packet.fec_shard_index = fec[0]; + packet.fec_base_sequence = mic_read_be16(fec + 2); + packet.fec_base_timestamp = mic_read_be32(fec + 4); + packet.fec_data_ssrc = mic_read_be32(fec + 8); + if (packet.fec_data_ssrc == 0) { + return false; + } + + for (int i = 0; i < RTPA_DATA_SHARDS; ++i) { + const auto payload_length = mic_read_be16(extension + 4 + i * 2); + if (payload_length == 0 || payload_length > parity_length || payload_length > MAX_AUDIO_PACKET_SIZE) { + return false; + } + packet.fec_payload_lengths[i] = payload_length; + } + + packet.payload = rtp_payload + fec_headers_length; + packet.payload_length = parity_length; + return true; + } + + bool mic_parse_foundation_legacy( + const std::uint8_t *data, + std::size_t bytes, + normalized_mic_packet_t &packet + ) { + if (!data || bytes <= kRtpFixedHeaderSize || data[0] != 0 || + (data[1] != 96 && data[1] != MIC_PACKET_TYPE_OPUS)) { + return false; + } + + packet.kind = mic_packet_kind_e::data; + packet.sequence = (std::uint16_t) (data[2] | (data[3] << 8)); + packet.payload = data + kRtpFixedHeaderSize; + packet.payload_length = bytes - kRtpFixedHeaderSize; + return packet.payload_length <= MAX_AUDIO_PACKET_SIZE; + } + + bool mic_parse_foundation_ext_5504( + const std::uint8_t *data, + std::size_t bytes, + normalized_mic_packet_t &packet + ) { + if (!data || bytes <= 4 || mic_read_be16(data) != IDX_MIC_DATA_TYPE) { + return false; + } + + const auto outer_sequence = (std::uint16_t) (data[2] | (data[3] << 8)); + std::size_t offset = 4; + if (bytes >= 6) { + const auto declared_length = mic_read_be16(data + 4); + if (declared_length == bytes - 6 || declared_length == bytes - 4) { + offset = 6; + } + } + if (offset >= bytes) { + return false; + } + + if (bytes - offset > kRtpFixedHeaderSize && mic_parse_foundation_legacy(data + offset, bytes - offset, packet)) { + return true; + } + + packet.kind = mic_packet_kind_e::data; + packet.sequence = outer_sequence; + packet.payload = data + offset; + packet.payload_length = bytes - offset; + return packet.payload_length > 0 && packet.payload_length <= MAX_AUDIO_PACKET_SIZE; + } + + bool mic_parse_packet( + const session_t &session, + const std::uint8_t *data, + std::size_t bytes, + normalized_mic_packet_t &packet + ) { + if (session.mic.protocol_version == rtsp_stream::MIC_PROTOCOL_MOONLIGHT_V1) { + return mic_parse_moonlight_v1(session, data, bytes, packet); + } + return mic_parse_foundation_ext_5504(data, bytes, packet) || + mic_parse_foundation_legacy(data, bytes, packet); + } + + session_t *mic_find_session(broadcast_ctx_t &ctx, const udp::endpoint &peer) { + for (auto *session : *ctx.control_server._sessions) { + if (!session || !session->mic.enabled) { + continue; + } + if (session->state.load(std::memory_order_relaxed) != session::state_e::RUNNING) { + continue; + } + + // Prefer the session whose control/audio/video peer IP matches the sender. + if (session->audio.peer.address() == peer.address() || + session->video.peer.address() == peer.address()) { + return session; + } + if (session->control.peer) { + TUPLE_2D(port, addr, platf::from_sockaddr_ex((sockaddr *) &session->control.peer->address.address)); + (void) port; + boost::system::error_code ec; + auto control_addr = boost::asio::ip::make_address(addr, ec); + if (!ec && control_addr == peer.address()) { + return session; + } + } + } + return nullptr; + } + + void mic_on_packet(broadcast_ctx_t &ctx, const udp::endpoint &peer, const std::uint8_t *data, std::size_t bytes) { + // Locate and retain the session before selecting a parser. The negotiated + // wire protocol is locked per session and never inferred from UDP bytes. + auto sessions = ctx.control_server._sessions.lock(); + auto *session = mic_find_session(ctx, peer); + if (!session) { + return; + } + + std::lock_guard lg {session->mic.lock}; + if (!session->mic.enabled || !session->mic.decoder) { + return; + } + + normalized_mic_packet_t packet; + if (!mic_parse_packet(*session, data, bytes, packet)) { + return; + } + + ++session->mic.packets_received; + std::vector payload(packet.payload, packet.payload + packet.payload_length); + + if (packet.kind == mic_packet_kind_e::fec) { + mic_handle_v1_fec_shard( + *session, + packet.ssrc, + packet.sequence, + packet.fec_data_ssrc, + packet.fec_base_sequence, + packet.fec_base_timestamp, + packet.fec_shard_index, + packet.fec_payload_lengths, + std::move(payload) + ); + return; + } + + if (session->mic.protocol_version == rtsp_stream::MIC_PROTOCOL_MOONLIGHT_V1) { + mic_handle_v1_data_shard(*session, packet.sequence, packet.timestamp, packet.ssrc, std::move(payload)); + } else { + mic_handle_legacy_data(*session, packet.sequence, std::move(payload)); + } + } + } // namespace + + bool ensure_mic_sock_open(broadcast_ctx_t &ctx) { + if (ctx.mic_socket_enabled.load(std::memory_order_acquire)) { + return true; + } + if (!config::audio.stream_mic) { + return false; + } + + auto address_family = net::af_from_enum_string(config::sunshine.address_family); + auto bind_addr_str = net::get_bind_address(address_family); + boost::system::error_code ec; + const auto bind_addr = boost::asio::ip::make_address(bind_addr_str, ec); + if (ec) { + BOOST_LOG(error) << "[mic] Invalid bind address: "sv << bind_addr_str << " - " << ec.message(); + return false; + } + + auto protocol = net::udp_protocol_for_address(bind_addr); + auto mic_port = net::map_port(MIC_STREAM_PORT); + + if (ctx.mic_sock.is_open()) { + ctx.mic_sock.close(); + } + + ctx.mic_io_context.restart(); + ctx.mic_sock.open(protocol, ec); + if (ec) { + BOOST_LOG(error) << "[mic] Couldn't open mic UDP socket: "sv << ec.message(); + return false; + } + + ctx.mic_sock.bind(udp::endpoint(bind_addr, mic_port), ec); + if (ec) { + BOOST_LOG(error) << "[mic] Couldn't bind mic UDP socket to port ["sv << mic_port << "]: "sv << ec.message(); + ctx.mic_sock.close(); + return false; + } + + ctx.mic_socket_enabled.store(true, std::memory_order_release); + BOOST_LOG(info) << "[mic] UDP socket listening on port "sv << mic_port; + return true; + } + + bool start_mic_receiver(broadcast_ctx_t &ctx) { + std::lock_guard lg {ctx.mic_start_mutex}; + if (ctx.mic_thread.joinable()) { + return true; + } + if (!ensure_mic_sock_open(ctx)) { + return false; + } + if (ctx.mic_io_context.stopped()) { + ctx.mic_io_context.restart(); + } + ctx.mic_thread = std::thread {micRecvThread, std::ref(ctx)}; + return true; + } + + void mic_session_acquire(broadcast_ctx_t &ctx) { + ctx.mic_sessions_count.fetch_add(1, std::memory_order_acq_rel); + } + + void mic_session_release(broadcast_ctx_t &ctx) { + const auto remaining = ctx.mic_sessions_count.fetch_sub(1, std::memory_order_acq_rel) - 1; + if (remaining < 0) { + ctx.mic_sessions_count.store(0, std::memory_order_release); + } + } + + int mic_session_start(session_t &session) { + if (!session.mic.enabled) { + return 0; + } + if (!config::audio.stream_mic) { + BOOST_LOG(warning) << "[mic] session requested mic but stream_mic is disabled"sv; + session.mic.enabled = false; + return -1; + } + + if (!session.broadcast_ref || !start_mic_receiver(*session.broadcast_ref.get())) { + BOOST_LOG(warning) << "[mic] failed to start UDP receiver"sv; + session.mic.enabled = false; + return -1; + } + +#ifdef _WIN32 + // Reuse the process-wide audio context so the Steam backend stays warm. + session.mic.audio_ctx = audio::get_audio_ctx_ref(); + auto *control = mic_audio_control(session); + if (!control) { + BOOST_LOG(error) << "[mic] no audio control available"sv; + session.mic.enabled = false; + return -1; + } + + if (config::audio.mic_require_steam && !control->mic_redirect_available()) { + BOOST_LOG(warning) << "[mic] Steam Streaming Microphone not ready; uplink will no-op"sv; + } + + if (control->init_mic_redirect_device() != 0) { + BOOST_LOG(warning) << "[mic] init_mic_redirect_device failed; packets will be dropped"sv; + } + + if (!config::audio.mic_capture_device.empty()) { + session.mic.capture_snap = control->snapshot_capture_defaults(); + control->switch_capture_to(config::audio.mic_capture_device); + session.mic.capture_switched = true; + } +#else + BOOST_LOG(warning) << "[mic] uplink is only implemented on Windows"sv; + session.mic.enabled = false; + return -1; +#endif + + int err = 0; + OpusDecoder *dec = opus_decoder_create(kMicSampleRate, 1, &err); + if (!dec || err != OPUS_OK) { + BOOST_LOG(error) << "[mic] opus_decoder_create failed: "sv << err; + session.mic.enabled = false; + return -1; + } + session.mic.decoder.reset(dec); + + session.mic.pending.clear(); + session.mic.has_playout_cursor = false; + session.mic.expected_seq = 0; + session.mic.data_ssrc = 0; + session.mic.parity_ssrc = 0; + session.mic.rs_blocks.clear(); + session.mic.recent_data.clear(); + session.mic.rs.reset(); + if (session.mic.protocol_version == rtsp_stream::MIC_PROTOCOL_MOONLIGHT_V1 && !mic_ensure_rs(session)) { + BOOST_LOG(warning) << "[mic] failed to initialize moonlight-mic/1 RS decoder"sv; + } + session.mic.packets_received = 0; + session.mic.frames_written = 0; + session.mic.decode_errors = 0; + session.mic.rs_recovered = 0; + + if (session.broadcast_ref) { + mic_session_acquire(*session.broadcast_ref.get()); + } + + BOOST_LOG(info) << "[mic] session uplink started (protocol="sv + << (session.mic.protocol_version == rtsp_stream::MIC_PROTOCOL_MOONLIGHT_V1 ? "moonlight-mic/1"sv : "foundation-legacy"sv) + << ')'; + return 0; + } + + void mic_session_stop(session_t &session) { + if (!session.mic.enabled && !session.mic.decoder && !session.mic.audio_ctx) { + return; + } + + std::lock_guard lg {session.mic.lock}; + + if (session.broadcast_ref && session.mic.enabled) { + mic_session_release(*session.broadcast_ref.get()); + } + + auto *control = mic_audio_control(session); + if (control) { + if (session.mic.capture_switched) { + control->restore_capture_from(session.mic.capture_snap); + session.mic.capture_switched = false; + } + // Keep the shared backend warm for other sessions; only release when no mic sessions remain. + if (session.broadcast_ref && session.broadcast_ref->mic_sessions_count.load(std::memory_order_acquire) <= 0) { + control->release_mic_redirect_device(); + } + } + + session.mic.decoder.reset(); + session.mic.rs_blocks.clear(); + session.mic.recent_data.clear(); + session.mic.rs.reset(); + session.mic.pending.clear(); + session.mic.has_playout_cursor = false; + session.mic.data_ssrc = 0; + session.mic.parity_ssrc = 0; + session.mic.audio_ctx = {}; + session.mic.enabled = false; + + BOOST_LOG(info) << "[mic] session uplink stopped (pkts="sv << session.mic.packets_received + << ", frames="sv << session.mic.frames_written + << ", dec_err="sv << session.mic.decode_errors + << ", rs_rec="sv << session.mic.rs_recovered << ')'; + } + + void micRecvThread(broadcast_ctx_t &ctx) { + platf::set_thread_name("stream::micRecv"); + + auto broadcast_shutdown_event = mail::man->event(mail::broadcast_shutdown); + auto &io = ctx.mic_io_context; + auto &sock = ctx.mic_sock; + + if (!ctx.mic_socket_enabled.load(std::memory_order_acquire) || !sock.is_open()) { + BOOST_LOG(debug) << "[mic] recv thread exiting: socket not enabled"sv; + return; + } + + udp::endpoint peer; + std::array buf {}; + std::function recv_handler; + + recv_handler = [&](const boost::system::error_code &ec, std::size_t bytes) { + if (broadcast_shutdown_event->peek()) { + return; + } + + auto fg = util::fail_guard([&]() { + if (!broadcast_shutdown_event->peek() && sock.is_open()) { + sock.async_receive_from(asio::buffer(buf), peer, 0, recv_handler); + } + }); + + if (ec == boost::system::errc::connection_refused || ec == boost::system::errc::connection_reset) { + return; + } + if (ec == boost::asio::error::operation_aborted) { + return; + } + if (ec || !bytes) { + if (ec) { + BOOST_LOG(verbose) << "[mic] recv error: "sv << ec.message(); + } + return; + } + + mic_on_packet(ctx, peer, reinterpret_cast(buf.data()), bytes); + }; + + sock.async_receive_from(asio::buffer(buf), peer, 0, recv_handler); + + while (!broadcast_shutdown_event->peek()) { + io.run(); + if (broadcast_shutdown_event->peek()) { + break; + } + // run() returns when out of work; restart if the socket is still live. + if (!sock.is_open()) { + break; + } + io.restart(); + sock.async_receive_from(asio::buffer(buf), peer, 0, recv_handler); + } + + BOOST_LOG(debug) << "[mic] recv thread ended"sv; + } + int start_broadcast(broadcast_ctx_t &ctx) { // Reset the shutdown event to ensure it's cleared even if something // raised it between the last end_broadcast and now. @@ -2521,6 +3663,13 @@ namespace stream { ctx.video_sock.close(); ctx.audio_sock.close(); + ctx.mic_socket_enabled.store(false, std::memory_order_release); + { + boost::system::error_code ec; + ctx.mic_sock.close(ec); + } + ctx.mic_io_context.stop(); + video_packets.reset(); audio_packets.reset(); @@ -2532,6 +3681,10 @@ namespace stream { ctx.audio_thread.join(); BOOST_LOG(debug) << "Waiting for main control thread to end..."sv; ctx.control_thread.join(); + if (ctx.mic_thread.joinable()) { + BOOST_LOG(debug) << "Waiting for mic thread to end..."sv; + ctx.mic_thread.join(); + } BOOST_LOG(debug) << "All broadcasting threads ended"sv; broadcast_shutdown_event->reset(); @@ -2775,6 +3928,10 @@ namespace stream { join_deadline_t join_deadline {hung_stage}; BOOST_LOG(debug) << "Waiting for video to end..."sv; + if (session.mic.enabled || session.mic.audio_ctx || session.mic.decoder) { + mic_session_stop(session); + } + session.videoThread.join(); hung_stage->store("audio thread"); BOOST_LOG(debug) << "Waiting for audio to end..."sv; @@ -2941,6 +4098,13 @@ namespace stream { session.audioThread = std::thread {audioThread, &session}; session.videoThread = std::thread {videoThread, &session}; + if (session.mic.enabled) { + if (mic_session_start(session) != 0) { + BOOST_LOG(warning) << "[mic] failed to arm uplink for session"sv; + session.mic.enabled = false; + } + } + session.state.store(state_e::RUNNING, std::memory_order_relaxed); // Record session in persistent history @@ -3162,6 +4326,9 @@ namespace stream { session->audio.sequenceNumber = 0; session->audio.timestamp = 0; + session->mic.enabled = launch_session.enable_mic; + session->mic.protocol_version = launch_session.mic_protocol_version; + session->control.peer = nullptr; session->state.store(state_e::STOPPED, std::memory_order_relaxed); @@ -3170,4 +4337,31 @@ namespace stream { return session; } } // namespace session + + bool mic_backend_ready() { +#ifdef _WIN32 + if (!config::audio.stream_mic) { + return false; + } + return platf::mic_redirect_available(); +#else + return false; +#endif + } + + mic_status_t get_mic_status() { + mic_status_t status {}; +#ifdef _WIN32 + status.capable = config::audio.stream_mic; + if (!status.capable) { + return status; + } + status.ready = mic_backend_ready(); + status.port = net::map_port(MIC_STREAM_PORT); +#else + status.capable = false; +#endif + return status; + } + } // namespace stream diff --git a/src/stream.h b/src/stream.h index edc075333..daea55940 100644 --- a/src/stream.h +++ b/src/stream.h @@ -31,6 +31,7 @@ namespace stream { constexpr auto VIDEO_STREAM_PORT = 9; constexpr auto CONTROL_PORT = 10; constexpr auto AUDIO_STREAM_PORT = 11; + constexpr auto MIC_STREAM_PORT = 12; ///< Client microphone uplink (UDP/RTP) constexpr std::string_view video_format_name(int video_format) { switch (video_format) { @@ -75,6 +76,25 @@ namespace stream { struct session_t; + /** + * @brief Snapshot of microphone uplink state for /serverinfo and the web UI. + */ + struct mic_status_t { + bool capable = false; ///< Feature built/enabled in config + bool ready = false; ///< Steam Streaming Microphone endpoint present + std::uint16_t port = 0; ///< Mapped MIC_STREAM_PORT (0 if disabled) + }; + + /** + * @brief Check whether the Steam mic endpoint is currently available. + */ + bool mic_backend_ready(); + + /** + * @brief Current mic uplink advertisement/status snapshot. + */ + mic_status_t get_mic_status(); + struct config_t { audio::config_t audio; video::config_t monitor; diff --git a/tests/unit/test_rtsp_startup_snapshot.cpp b/tests/unit/test_rtsp_startup_snapshot.cpp index d937bdd45..f3c4aae53 100644 --- a/tests/unit/test_rtsp_startup_snapshot.cpp +++ b/tests/unit/test_rtsp_startup_snapshot.cpp @@ -42,6 +42,8 @@ namespace { ls.frame_generation_provider = "provider"; ls.lossless_scaling_target_fps = 144.0; ls.lossless_scaling_rtss_limit = 90; + ls.enable_mic = true; + ls.mic_protocol_version = rtsp_stream::MIC_PROTOCOL_MOONLIGHT_V1; return ls; } @@ -94,6 +96,8 @@ TEST(RtspStartupSnapshot, CopiesAllConsumedFields) { EXPECT_EQ(clone->frame_generation_provider, source.frame_generation_provider); EXPECT_EQ(clone->lossless_scaling_target_fps, source.lossless_scaling_target_fps); EXPECT_EQ(clone->lossless_scaling_rtss_limit, source.lossless_scaling_rtss_limit); + EXPECT_EQ(clone->enable_mic, source.enable_mic); + EXPECT_EQ(clone->mic_protocol_version, source.mic_protocol_version); // The clone intentionally does NOT copy rtsp_cipher: it is move-only and only the // io_context respond() path uses it, never the startup worker.